|
1 | 1 | import uuid |
2 | 2 | from collections.abc import AsyncGenerator, Awaitable, Callable |
| 3 | +from datetime import timedelta |
3 | 4 | from logging import getLogger |
4 | 5 | from typing import ( |
5 | 6 | TYPE_CHECKING, |
@@ -295,28 +296,31 @@ async def listen(self) -> AsyncGenerator[AckableMessage, None]: |
295 | 296 | ) |
296 | 297 | logger.debug("Starting fetching unacknowledged messages") |
297 | 298 | for stream in [self.queue_name, *self.additional_streams.keys()]: |
298 | | - lock = redis_conn.lock( |
| 299 | + pipe = redis_conn.pipeline() |
| 300 | + lock = pipe.lock( |
299 | 301 | f"autoclaim:{self.consumer_group_name}:{stream}", |
300 | 302 | timeout=self.unacknowledged_lock_timeout, |
301 | 303 | ) |
302 | | - if await lock.locked(): |
303 | | - continue |
304 | | - async with lock: |
305 | | - pending = await redis_conn.xautoclaim( |
306 | | - name=stream, |
307 | | - groupname=self.consumer_group_name, |
308 | | - consumername=self.consumer_name, |
309 | | - min_idle_time=self.idle_timeout, |
310 | | - count=self.unacknowledged_batch_size, |
311 | | - ) |
312 | | - logger.debug( |
313 | | - "Found %d pending messages in stream %s", |
314 | | - len(pending[1]), |
315 | | - stream, |
| 304 | + await lock.acquire() |
| 305 | + await pipe.xautoclaim( |
| 306 | + name=stream, |
| 307 | + groupname=self.consumer_group_name, |
| 308 | + consumername=self.consumer_name, |
| 309 | + min_idle_time=self.idle_timeout, |
| 310 | + count=self.unacknowledged_batch_size, |
| 311 | + ) |
| 312 | + await lock.release() |
| 313 | + results = await pipe.execute() |
| 314 | + pending = results[1] |
| 315 | + |
| 316 | + logger.debug( |
| 317 | + "Found %d pending messages in stream %s", |
| 318 | + len(pending[1]), |
| 319 | + stream, |
| 320 | + ) |
| 321 | + for msg_id, msg in pending[1]: |
| 322 | + logger.debug("Received message: %s", msg) |
| 323 | + yield AckableMessage( |
| 324 | + data=msg[b"data"], |
| 325 | + ack=self._ack_generator(id=msg_id, queue_name=stream), |
316 | 326 | ) |
317 | | - for msg_id, msg in pending[1]: |
318 | | - logger.debug("Received message: %s", msg) |
319 | | - yield AckableMessage( |
320 | | - data=msg[b"data"], |
321 | | - ack=self._ack_generator(id=msg_id, queue_name=stream), |
322 | | - ) |
|
0 commit comments