Add support for Kafka export in Ingestor and DataManager classes
- Introduced EXPORT_TO_KAFKA environment variable in Ingestor class. - Refactored DataManager to conditionally initialize Kafka producer based on export setting. - Updated IngestorManager to pass Redis configuration as a dictionary. - Enhanced error handling and logging for Kafka message publishing.
This commit is contained in:
@@ -34,6 +34,7 @@ class Ingestor:
|
||||
"""
|
||||
|
||||
kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092")
|
||||
self.export_to_kafka = getenv("EXPORT_TO_KAFKA", "false")
|
||||
self.redis_host = getenv("REDIS_HOST", "localhost")
|
||||
self.redis_port = int(getenv("REDIS_PORT", "6379"))
|
||||
self.redis_username = getenv("REDIS_USERNAME", None)
|
||||
@@ -137,8 +138,12 @@ class Ingestor:
|
||||
|
||||
self.ingestor_manager = IngestorManager(
|
||||
self.kafka_servers,
|
||||
self.redis_host,
|
||||
self.redis_port,
|
||||
{
|
||||
'host': self.redis_host,
|
||||
'port': self.redis_port,
|
||||
'username': self.redis_username,
|
||||
'password': self.redis_password,
|
||||
},
|
||||
self.lease_ttl,
|
||||
self.heartbeat_ttl,
|
||||
self.pod_id,
|
||||
@@ -147,8 +152,7 @@ class Ingestor:
|
||||
self.mongo_database,
|
||||
self.logger,
|
||||
self.notification_handler,
|
||||
self.redis_username,
|
||||
self.redis_password,
|
||||
self.export_to_kafka,
|
||||
)
|
||||
|
||||
# Declare ingestor ative
|
||||
|
||||
Reference in New Issue
Block a user