33import ssl
44
55from aiokafka import AIOKafkaProducer
6+ from aiokafka .errors import KafkaError
67from tenacity import (
78 retry ,
89 retry_if_exception_type ,
@@ -25,7 +26,7 @@ class QueueService(BaseService):
2526 def __init__ (self ):
2627 super ().__init__ ()
2728 self .kafka_topic = CROWD_KAFKA_TOPIC
28- self .kafka_producer = AIOKafkaProducer ( ** self . _build_kafka_config ())
29+ self .kafka_producer : AIOKafkaProducer | None = None
2930 self ._connected = False
3031
3132 def _build_kafka_config (self ):
@@ -64,7 +65,11 @@ async def ensure_connected(self):
6465 async def _is_connection_healthy (self ) -> bool :
6566 """Check if the current connection is healthy"""
6667 try :
67- return self .kafka_producer ._sender is not None and not self .kafka_producer ._closed
68+ producer = self .kafka_producer
69+ if producer is None or producer ._closed :
70+ return False
71+ sender_task = producer ._sender .sender_task
72+ return sender_task is not None and not sender_task .done ()
6873 except Exception :
6974 return False
7075
@@ -85,19 +90,23 @@ async def connect(self):
8590 return
8691 self .logger .info ("Connecting to kafka..." )
8792 try :
93+ self .kafka_producer = AIOKafkaProducer (** self ._build_kafka_config ())
8894 await self .kafka_producer .start ()
8995 self .logger .info ("Connected to kafka" )
9096 self ._connected = True
9197 except Exception as e :
9298 self .logger .error (f"Queue connection failed: { e } " )
99+ await self .disconnect ()
93100 raise QueueConnectionError () from e
94101
95102 async def disconnect (self ):
96103 """
97104 Disconnect from Kafka and close producer
98105 """
99- if self .kafka_producer ._closed :
106+ if self .kafka_producer is None or self . kafka_producer ._closed :
100107 self .logger .debug ("Producer already closed, skipping" )
108+ self .kafka_producer = None
109+ self ._connected = False
101110 return
102111
103112 try :
@@ -106,6 +115,7 @@ async def disconnect(self):
106115 except Exception as e :
107116 self .logger .error (f"Error during disconnect: { e } " )
108117 finally :
118+ self .kafka_producer = None
109119 self ._connected = False
110120
111121 async def shutdown (self ):
@@ -135,20 +145,15 @@ async def shutdown(self):
135145 self .logger .error (f"Failed to shutdown queue service: { repr (e )} " )
136146 # Don't raise - allow application to continue shutdown
137147
138- async def send_batch_activities (self , activities_kafka : list [dict [str , str ]]):
139- """
140- Send multiple pre-prepared activities to Kafka in a non-blocking way.
141- Args:
142- activities_kafka: List of dicts with 'message_id' and 'payload' keys
143- (prepared by CommitService.prepare_activity_for_db_and_queue)
144- """
145- if not activities_kafka :
146- return
147-
148+ @retry (
149+ stop = stop_after_attempt (3 ),
150+ wait = wait_exponential (multiplier = 1 , min = 2 , max = 10 ),
151+ retry = retry_if_exception_type (KafkaError ),
152+ reraise = True ,
153+ )
154+ async def _emit_batch (self , activities_kafka : list [dict [str , str ]]):
155+ """Emit a single batch, rebuilding the producer and retrying on Kafka errors."""
148156 await self .ensure_connected ()
149-
150- self .logger .info (f"Emitting { len (activities_kafka )} activities to kafka queue..." )
151-
152157 try :
153158 futures = [
154159 self .kafka_producer .send (
@@ -160,7 +165,25 @@ async def send_batch_activities(self, activities_kafka: list[dict[str, str]]):
160165 ]
161166 # Wait for all messages to be sent
162167 await asyncio .gather (* futures , return_exceptions = False )
168+ except KafkaError :
169+ self .logger .warning ("Kafka send failed, reconnecting before retry..." )
170+ await self .disconnect ()
171+ raise
163172
173+ async def send_batch_activities (self , activities_kafka : list [dict [str , str ]]):
174+ """
175+ Send multiple pre-prepared activities to Kafka in a non-blocking way.
176+ Args:
177+ activities_kafka: List of dicts with 'message_id' and 'payload' keys
178+ (prepared by CommitService.prepare_activity_for_db_and_queue)
179+ """
180+ if not activities_kafka :
181+ return
182+
183+ self .logger .info (f"Emitting { len (activities_kafka )} activities to kafka queue..." )
184+
185+ try :
186+ await self ._emit_batch (activities_kafka )
164187 self .logger .info (f"Successfully emitted { len (activities_kafka )} activities to queue" )
165188 except Exception as e :
166189 self .logger .error (f"Failed to emit batch to queue with error: { repr (e )} " )
0 commit comments