From cbaa66026c3b483dbaf4ccbb8ebaa576a82e14f7 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 3 Jul 2025 13:45:09 -0300 Subject: [PATCH] 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. --- ingestor/ingestor.py | 12 ++- ingestor/managers/data_manager.py | 115 ++++++++++++++------------ ingestor/managers/ingestor_manager.py | 14 +++- 3 files changed, 80 insertions(+), 61 deletions(-) diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 8789779..fb039c7 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -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 diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index 466542a..f1c89dc 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -18,6 +18,7 @@ class DataManager: kafka_servers: str, mongo_connection_string: str, mongo_database: str, + export_to_kafka: bool, logger: Logger, notification_handler: NotificationHandler, ) -> None: @@ -35,40 +36,44 @@ class DataManager: self.pod_id = os.getenv("HOSTNAME", "localhost") self.kafka_producer = None - for i in range(0, 3): - logger.info( - f"Trying ({i}) to initializing DataManager with Kafka servers: {kafka_servers}" - ) - try: - self.kafka_producer = KafkaProducer( - bootstrap_servers=kafka_servers, - value_serializer=lambda v: json.dumps(v).encode( - "utf-8" - ), # Serialize JSON messages - key_serializer=lambda k: str( - k).encode("utf-8") if k else None, - ) - # Kafka connected - metrics.KAFKA_CONNECTION_STATUS.labels( - pod_id=self.pod_id).set(1) - break - except NoBrokersAvailable: - logger.error( - f"Kafka servers {kafka_servers} are not available. Retrying..." - ) - sleep(5) - else: - # Kafka not connected - metrics.KAFKA_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(0) - logger.error( - f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts." - ) - raise NoBrokersAvailable( - f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts." - ) + self.export_to_kafka = export_to_kafka - logger.info( - f"DataManager initialized with Kafka servers: {kafka_servers}") + if self.export_to_kafka: + for i in range(0, 3): + logger.info( + f"Trying ({i}) to initializing DataManager with Kafka servers: {kafka_servers}" + ) + try: + self.kafka_producer = KafkaProducer( + bootstrap_servers=kafka_servers, + value_serializer=lambda v: json.dumps(v).encode( + "utf-8" + ), # Serialize JSON messages + key_serializer=lambda k: str( + k).encode("utf-8") if k else None, + ) + # Kafka connected + metrics.KAFKA_CONNECTION_STATUS.labels( + pod_id=self.pod_id).set(1) + break + except NoBrokersAvailable: + logger.error( + f"Kafka servers {kafka_servers} are not available. Retrying..." + ) + sleep(5) + else: + # Kafka not connected + metrics.KAFKA_CONNECTION_STATUS.labels( + pod_id=self.pod_id).set(0) + logger.error( + f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts." + ) + raise NoBrokersAvailable( + f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts." + ) + + logger.info( + f"DataManager initialized with Kafka servers: {kafka_servers}") self.connection_string = mongo_connection_string self.database = mongo_database @@ -147,29 +152,33 @@ class DataManager: Exception: If there is an error during message delivery, it will be handled by the `delivery_error` callback. """ - try: + if self.export_to_kafka: + try: - self.logger.debug(f"Publishing message to topic {topic}: {data}") - self.kafka_producer.send(topic=topic, value=data).add_callback( - self.delivery_report - ).add_errback(self.delivery_error) + self.logger.debug( + f"Publishing message to topic {topic}: {data}") + self.kafka_producer.send(topic=topic, value=data).add_callback( + self.delivery_report + ).add_errback(self.delivery_error) - self.kafka_producer.flush(timeout=10) - metrics.KAFKA_MESSAGES_SENT.labels( - pod_id=self.pod_id, topic=topic).inc() + self.kafka_producer.flush(timeout=10) + metrics.KAFKA_MESSAGES_SENT.labels( + pod_id=self.pod_id, topic=topic).inc() - except Exception as e: - metrics.KAFKA_MESSAGES_ERRORS.labels( - pod_id=self.pod_id, topic=topic).inc() - trace = traceback.format_exc() - self.notification_handler.build_and_send_notification( - notification_id=f"KAFKA_PRODUCER_ERROR_{topic}", - message=f"Error publishing message to topic {topic}: {e}", - block="kafka_producer", - level=NotificationLevel.ERROR, - attachment_content=trace, - ) - self.logger.error(trace) + except Exception as e: + metrics.KAFKA_MESSAGES_ERRORS.labels( + pod_id=self.pod_id, topic=topic).inc() + trace = traceback.format_exc() + self.notification_handler.build_and_send_notification( + notification_id=f"KAFKA_PRODUCER_ERROR_{topic}", + message=f"Error publishing message to topic {topic}: {e}", + block="kafka_producer", + level=NotificationLevel.ERROR, + attachment_content=trace, + ) + self.logger.error(trace) + else: + self.logger.info(f"Skipping message to topic {topic}: {data}") try: collection = self.mongo_db[topic] diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 6fe99c3..dc56e55 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -12,14 +12,20 @@ import ingestor.metrics as metrics class IngestorManager(): def __init__(self, - kafka_servers: str, redis_host: str, redis_port: int, + kafka_servers: str, redis_data: dict, lease_ttl: int, heartbeat_ttl: int, pod_id: str, poll_interval: int, mongo_connection_string: str, mongo_database: str, logger: Logger, notification_handler: NotificationHandler, - redis_username: str = None, redis_password: str = None): + export_to_kafka: bool = False): + + redis_host = redis_data.get('host') + redis_port = redis_data.get('port') + redis_username = redis_data.get('username', None) + redis_password = redis_data.get('password', None) self.data_manager = DataManager( - kafka_servers, mongo_connection_string, mongo_database, logger, notification_handler) + kafka_servers, mongo_connection_string, mongo_database, + export_to_kafka, logger, notification_handler) self.opc_managers = {} self.resource_manager = ResourceManager( redis_host, redis_port, lease_ttl, heartbeat_ttl, pod_id, redis_username, redis_password @@ -156,7 +162,7 @@ class IngestorManager(): _a = [self.managed_tags[slot].pop(server, None) for slot, _value in current_managed_tags.items()] - for server in list(self.opc_managers.keys()): + for server in self.opc_managers.keys(): if server not in registered_servers: self.logger.warning( f"Server {server} not found in managed tags. "