Implement a retry queue with visibility timeout in Python
This code simulates a message queue with a visibility timeout, allowing messages to be retried if not deleted before the timeout expires.
Python code
47 linesimport time
from collections import deque
class SimpleQueue:
def __init__(self, visibility_timeout=2):
self.queue = deque()
self.in_flight = {}
self.visibility_timeout = visibility_timeout
def send(self, message):
self.queue.append(message)
def receive(self):
if self.queue:
message = self.queue.popleft()
self.in_flight[message] = time.time() + self.visibility_timeout
return message
return None
def delete(self, message):
if message in self.in_flight:
del self.in_flight[message]
def retry_after_timeout(self):
current_time = time.time()
expired = [
msg for msg, deadline in self.in_flight.items() if deadline <= current_time
]
for msg in expired:
del self.in_flight[msg]
self.queue.appendleft(msg)
return expired
if __name__ == "__main__":
queue = SimpleQueue(visibility_timeout=1)
queue.send("task-1")
queue.send("task-2")
print("First receive:", queue.receive())
time.sleep(1.2)
requeued = queue.retry_after_timeout()
print("Requeued after timeout:", requeued)
print("Next receive:", queue.receive())
queue.delete("task-1")
print("Queue size:", len(queue.queue))
Output
First receive: task-1
Requeued after timeout: ['task-1']
Next receive: task-1
Queue size: 0
How it works
The SimpleQueue uses a deque for the main queue and a dictionary to track in-flight messages with their visibility deadlines. When a message is received, it is moved to the in-flight dictionary with a deadline set to time.time() + visibility_timeout. If the message is not deleted before the deadline, retry_after_timeout moves it back to the front of the queue for retry. This pattern mimics the behavior of message brokers like SQS, where a consumer must delete a message within the visibility timeout to prevent it from being redelivered. The delete method removes the message from in-flight, acknowledging successful processing.
Common mistakes
- Forgetting to delete messages after successful processing, causing duplicate deliveries.
- Not handling concurrent access to the queue and in-flight dictionary in a multi-threaded environment.
- Using a list instead of deque, leading to O(n) pops from the front.
- Setting the visibility timeout too short, causing frequent reprocessing of slow tasks.
Variations
- Adding a dead-letter queue after a maximum number of retries.
- Using `queue.PriorityQueue` to prioritize retry messages.
Real-world use cases
- Mocking AWS SQS behavior in local development to test consumer retry logic without cloud costs.
- Simulating a message queue for unit testing consumer applications in CI pipelines.
- Building a simple in-process task queue with retry handling for educational or lightweight prototypes.
Sponsored
More from Streaming & messaging
- At Most Once Fire-and-Forget Mock in Python easy
- Batch Consume Process Commit Pattern in Python medium
- Build a Streaming Messaging Helper in Python easy
- Dead Letter Queue Failed Messages List Mock in Python easy
- Dedupe processed message IDs in Python easy
- Event Envelope with Schema Version Field in Python easy
Keep learning
Related tutorials and quizzes for this topic.