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.

Medium Python 3.9+ Aug 9, 2026 Streaming & messaging 13 views 0 copies

Python code

47 lines
Python 3.9+
import 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

stdout
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

  1. Adding a dead-letter queue after a maximum number of retries.
  2. 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

Run this sample

Open the browser IDE to tweak the example and see results without installing anything.

Open editor

More from Streaming & messaging

Related tutorials and quizzes for this topic.