Add EXPORT_TO_KAFKA environment variable and update Ingestor class logic
- Introduced EXPORT_TO_KAFKA variable in values.yaml to control Kafka export behavior. - Updated Ingestor class to parse and set the export_to_kafka attribute based on the environment variable.
This commit is contained in:
@@ -34,7 +34,14 @@ class Ingestor:
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092")
|
kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092")
|
||||||
self.export_to_kafka = getenv("EXPORT_TO_KAFKA", "false")
|
export_to_kafka = getenv("EXPORT_TO_KAFKA", "false")
|
||||||
|
|
||||||
|
if export_to_kafka and export_to_kafka == "true":
|
||||||
|
export_to_kafka = True
|
||||||
|
else:
|
||||||
|
export_to_kafka = False
|
||||||
|
|
||||||
|
self.export_to_kafka = export_to_kafka
|
||||||
self.redis_host = getenv("REDIS_HOST", "localhost")
|
self.redis_host = getenv("REDIS_HOST", "localhost")
|
||||||
self.redis_port = int(getenv("REDIS_PORT", "6379"))
|
self.redis_port = int(getenv("REDIS_PORT", "6379"))
|
||||||
self.redis_username = getenv("REDIS_USERNAME", None)
|
self.redis_username = getenv("REDIS_USERNAME", None)
|
||||||
|
|||||||
@@ -130,6 +130,9 @@ env:
|
|||||||
# Application variables
|
# Application variables
|
||||||
- name: KAFKA_SERVERS
|
- name: KAFKA_SERVERS
|
||||||
value: "kafka.kafka.svc.cluster.local:9092"
|
value: "kafka.kafka.svc.cluster.local:9092"
|
||||||
|
- name: EXPORT_TO_KAFKA
|
||||||
|
value: "false"
|
||||||
|
|
||||||
- name: REDIS_HOST
|
- name: REDIS_HOST
|
||||||
value: "redis-master.redis.svc.cluster.local"
|
value: "redis-master.redis.svc.cluster.local"
|
||||||
- name: REDIS_PORT
|
- name: REDIS_PORT
|
||||||
|
|||||||
Reference in New Issue
Block a user