Skip to content

Commit 1d2fcb8

Browse files
committed
fix: await kafka delivery futures for actual broker ack (CM-1318)
Signed-off-by: Uroš Marolt <uros@marolt.me>
1 parent 48c6de2 commit 1d2fcb8

1 file changed

Lines changed: 12 additions & 8 deletions

File tree

  • services/apps/mailing_list_integration/src/crowdmail/services/queue

services/apps/mailing_list_integration/src/crowdmail/services/queue/queue_service.py

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -157,15 +157,19 @@ async def _emit_batch(activities_kafka: list[dict[str, str]]):
157157
try:
158158
for i in range(0, len(activities_kafka), _SEND_CHUNK_SIZE):
159159
chunk = activities_kafka[i : i + _SEND_CHUNK_SIZE]
160-
futures = [
161-
_producer.send(
162-
topic=CROWD_KAFKA_TOPIC,
163-
key=activity["message_id"].encode("utf-8", errors="replace"),
164-
value=activity["payload"].encode("utf-8", errors="replace"),
160+
# producer.send() only enqueues and returns a delivery Future; that Future
161+
# must be awaited too, else acks="all" is never actually waited on.
162+
send_futures = await asyncio.gather(
163+
*(
164+
_producer.send(
165+
topic=CROWD_KAFKA_TOPIC,
166+
key=activity["message_id"].encode("utf-8", errors="replace"),
167+
value=activity["payload"].encode("utf-8", errors="replace"),
168+
)
169+
for activity in chunk
165170
)
166-
for activity in chunk
167-
]
168-
await asyncio.gather(*futures, return_exceptions=False)
171+
)
172+
await asyncio.gather(*send_futures, return_exceptions=False)
169173
except KafkaError:
170174
logger.warning("Kafka send failed, reconnecting before retry...")
171175
await disconnect()

0 commit comments

Comments
 (0)