SIENTIAPDE-1110
Refactor Kafka activity to use instance-level kafka_connector for improved message polling and management, replacing local variable usage with class attribute methods.
This commit is contained in:
@@ -58,20 +58,13 @@ class Kafka(BaseActivity):
|
|||||||
topic = input_data["topic"]
|
topic = input_data["topic"]
|
||||||
|
|
||||||
# Subscribe to the specified topic
|
# Subscribe to the specified topic
|
||||||
kafka_connector = KafkaConsumer(
|
self.kafka_connector.subscribe([topic])
|
||||||
topic,
|
|
||||||
bootstrap_servers=self.bootstrap_servers,
|
|
||||||
auto_offset_reset="earliest",
|
|
||||||
enable_auto_commit=True,
|
|
||||||
group_id=self.group_id,
|
|
||||||
value_deserializer=lambda x: json.loads(x.decode("utf-8"))
|
|
||||||
)
|
|
||||||
|
|
||||||
# List to store message values
|
# List to store message values
|
||||||
message_values = []
|
message_values = []
|
||||||
|
|
||||||
# Poll for messages
|
# Poll for messages
|
||||||
records = kafka_connector.poll(timeout_ms=self.polling_time)
|
records = self.kafka_connector.poll(timeout_ms=self.polling_time)
|
||||||
|
|
||||||
self.debug(
|
self.debug(
|
||||||
f"Polled {len(records)} records from topic: {topic}",
|
f"Polled {len(records)} records from topic: {topic}",
|
||||||
@@ -97,6 +90,6 @@ class Kafka(BaseActivity):
|
|||||||
metadata=metadata
|
metadata=metadata
|
||||||
)
|
)
|
||||||
|
|
||||||
kafka_connector.close()
|
self.kafka_connector.unsubscribe()
|
||||||
|
|
||||||
return DataFrame(message_values).to_dict()
|
return DataFrame(message_values).to_dict()
|
||||||
|
|||||||
Reference in New Issue
Block a user