Redis Reliable Queue
Build a reliable queue on Redis Streams with consumer groups. You use Winter Boot for the application runtime and the Winter Redis module from Winter Modules for streams access plus auto-started consumer workers — the equivalent of Redisson PRO Reliable Queue on the JVM side: at-least-once delivery, acknowledgement on success, redelivery of unacked entries, and dead-letter routing for poison messages.
Prerequisites
Section titled “Prerequisites”You need PHP 8.5 or later with the swoole, pcntl, and redis (phpredis) extensions. You also need a Redis 5+ server reachable from the app.
Start Redis before you run the app:
docker run -d -p 30637:6379 --name redis \ redis:7 redis-server --requirepass redis123Project structure
Section titled “Project structure”The sample uses this layout:
redis-queue/├── bin/│ └── application.php # Application entry point├── config/│ ├── application.yml # Winter Boot application config│ └── redis-config.yml # Redis connections and queue consumers config├── src/│ ├── RedisQueueSampleApplication.php # Main application class│ ├── consumer/│ │ ├── OrderEventsConsumer.php # Queue worker (file writer)│ │ └── RetryDemoConsumer.php # Fails first delivery transiently│ └── controller/│ └── RedisQueueDemoController.php # Send/inspect REST endpoints└── composer.json # DependenciesInstall dependencies
Section titled “Install dependencies”Require the framework and the modules package:
composer require suvera/winter-boot suvera/winter-modulesSource files
Section titled “Source files”Switch between the source files. Each tab shows the exact file from the sample.
The single entry point. It points Winter Boot at the config directory and at the sample namespace.
<?php
namespace dev\example;
use dev\winterframework\stereotype\WinterBootApplication;
#[WinterBootApplication( configDirectory: [__DIR__ . "/../config"], scanNamespaces: [ ['dev\\example', __DIR__ . ''] ])]class RedisQueueSampleApplication {
public static function main(): void { $winterApp = new \dev\winterframework\core\app\WinterWebSwooleApplication(); $winterApp->run(self::class); }}The worker. It extends AbstractConsumer and appends one line per message with timestamp, stream, entry ID, delivery count, and payload. A poison: payload throws a non-transient exception, which routes the entry to the dead-letter stream.
<?php
namespace dev\example\consumer;
use dev\winterframework\data\redis\consumer\AbstractConsumer;use dev\winterframework\data\redis\consumer\ConsumerRecord;use dev\winterframework\data\redis\consumer\ConsumerRecords;
class OrderEventsConsumer extends AbstractConsumer {
private const OUTPUT_FILE = '/tmp/redis-queue-messages.txt';
public function consume(ConsumerRecords $records): void { foreach ($records as $record) { /** @var ConsumerRecord $record */ $payload = $record->getPayload();
if (str_starts_with($payload, 'poison:')) { // Permanent failure: routes to the dead-letter stream. throw new \LogicException('poison message refused: ' . $payload); }
$line = sprintf( "[%s] Stream: %s | Id: %s | Deliveries: %d | Payload: %s", date('Y-m-d H:i:s'), $record->getStream(), $record->getId(), $record->getDeliveryCount(), $payload ) . PHP_EOL;
file_put_contents( self::OUTPUT_FILE, $line, FILE_APPEND | LOCK_EX );
self::logInfo('Wrote message to file: ' . trim($line)); } }}AbstractConsumer gives you the logInfo() helper used on the last line. The worker process XACKs each entry only after consume() returns successfully — unlike the SQS worker, which deletes messages even when consumption fails.
A second worker that demonstrates redelivery. It fails the first delivery of fail-once:<token> payloads with a transient error; the redelivered copy is processed normally.
<?php
namespace dev\example\consumer;
use dev\winterframework\data\redis\consumer\AbstractConsumer;use dev\winterframework\data\redis\consumer\ConsumerRecord;use dev\winterframework\data\redis\consumer\ConsumerRecords;
class RetryDemoConsumer extends AbstractConsumer {
private const OUTPUT_FILE = '/tmp/redis-queue-retry.txt';
/** @var array<string, bool> */ private static array $failedOnce = [];
public function consume(ConsumerRecords $records): void { foreach ($records as $record) { /** @var ConsumerRecord $record */ $payload = $record->getPayload();
if (str_starts_with($payload, 'fail-once:') && !isset(self::$failedOnce[$payload])) { self::$failedOnce[$payload] = true; throw new \RuntimeException('transient failure (first delivery): ' . $payload); }
// ... append to self::OUTPUT_FILE as above } }}RuntimeException is listed in that consumer’s transientExceptions, so the failure is retried in-worker and then left pending for XCLAIM redelivery instead of going to the dead letter.
A small REST controller for producing messages and inspecting streams. It uses the autowired RedisQueueServiceImpl for sends and PhpRedisTemplate for pending-list and dead-letter inspection.
#[RestController]class RedisQueueDemoController {
#[Autowired] private RedisQueueServiceImpl $queues;
#[Autowired] private PhpRedisTemplate $redis;
#[PostMapping(path: "/queuedemo/send")] public function send( #[RequestParam] string $consumer, #[RequestParam] string $message ): array { // $consumer is a consumer name or a raw stream name. // Returns the stream entry id. return ['success' => true, 'entryId' => $this->queues->send($consumer, $message)]; }
#[GetMapping(path: "/queuedemo/pending")] public function pending( #[RequestParam] string $stream, #[RequestParam] string $group ): array { $pending = $this->redis->xpending($stream, $group, '-', '+', 100); return ['success' => true, 'pendingCount' => is_array($pending) ? count($pending) : 0]; }
#[GetMapping(path: "/queuedemo/dlq")] public function dlq(#[RequestParam] string $stream): array { $entries = $this->redis->xrange($stream, '-', '+', 100); return ['success' => true, 'entries' => $entries === false ? [] : $entries]; }}The sample declares framework dependencies with PSR-4 autoloading for its own namespace.
{ "name": "suvera/winter-boot-redis-queue-sample", "require": { "ext-pcntl": "*", "ext-swoole": "*", "ext-redis": "*", "suvera/winter-boot": "@dev", "suvera/winter-modules": "@dev" }, "autoload": { "psr-4": { "dev\\example\\": "src/" } }}The runnable sample in winter-boot-samples adds local path repositories for winter-boot and winter-modules so it resolves them from a sibling checkout. You do not need those entries when you install released packages from Packagist.
Configuration
Section titled “Configuration”Switch between the two config files. application.yml enables the module, and redis-config.yml defines the connections plus the queue consumers.
The full sample file registers RedisModule and sets the server and app identity:
server: port: 8080 address: 0.0.0.0 context-path: /winter: application: name: Redis Queue Sample Application id: redis-queue-sample-app version: 1.0.0modules: - module: dev\winterframework\data\redis\RedisModule enabled: true configFile: redis-config.ymlSee Configuration for every application.yml key.
Defines the primary single-node connection and two queue consumers. stream defaults to the consumer name and group to <stream>-group:
phpredis: singles: - name: __default__ host: localhost port: 30637 auth: redis123 timeout: 2.5 readTimeout: 2.5
- name: primary host: localhost port: 30637 auth: redis123 timeout: 2.5 readTimeout: 2.5
redis: consumers: - name: __default__ workerNum: 1 blockMs: 2000 # XREADGROUP BLOCK per iteration batchSize: 10 # max entries per read claimIdleMs: 3000 # reclaim PEL entries idle this long reclaimIntervalMs: 2000 # how often to scan for stale entries maxDeliveries: 5 # then moved to deadLetterStream retries: 2 # in-worker retries on transient errors retryWaitMs: 200 transientExceptions: []
- name: order-events-consumer redis: primary stream: order-events group: order-events-group workerNum: 2 workerClass: dev\example\consumer\OrderEventsConsumer deadLetterStream: order-events-dlq
- name: retry-demo-consumer redis: primary stream: retry-demo group: retry-demo-group workerNum: 1 retries: 1 # leave transient failures pending for redelivery workerClass: dev\example\consumer\RetryDemoConsumer deadLetterStream: retry-demo-dlq transientExceptions: [ RuntimeException ]See the Redis module for the connection properties.
Reliability semantics
Section titled “Reliability semantics”Each consumer spawns workerNum Swoole worker processes at boot. Every worker loops over XREADGROUP with BLOCK on its stream and group:
- Ack on success. Entries are
XACKed only afterconsume()returns without throwing. Unacked entries stay in the group pending list. - Redelivery. A throttled
XPENDINGscan reclaims entries idle pastclaimIdleMs(for example after a transient failure or a crashed worker) viaXCLAIM, bumpingConsumerRecord::getDeliveryCount(). - Retries. Exceptions listed in
transientExceptionsare retriedretriestimes in-worker first; when attempts run out the entry is left pending for redelivery. Any other exception is permanent. - Dead letter. Entries redelivered past
maxDeliveries, or failing permanently, areXADDed todeadLetterStream(withfailedStream,failedId,failedReason,failedAtfields) and acked. Without a dead-letter stream they are dropped with an error log so one poison message never blocks the stream. - Producing.
RedisQueueServiceImpl::send($consumerOrStream, $message, $fields)appends to the consumer’s stream (or a raw stream name) and returns the entry id. Non-string payloads are JSON-encoded into thepayloadfield.
Run the app
Section titled “Run the app”Start the app, send messages, then check the output files and the pending list.
1. Start the application:
composer installphp bin/application.phpThe queue workers start consuming as soon as the app boots.
2. Send a message:
curl -X POST "http://localhost:8080/queuedemo/send?consumer=order-events-consumer&message=hello"3. Verify consumption:
cat /tmp/redis-queue-messages.txtcurl "http://localhost:8080/queuedemo/pending?stream=order-events&group=order-events-group"You see one line per consumed message and an empty pending list ("pendingCount": 0), for example:
[2026-09-26 03:23:05] Stream: order-events | Id: 1790392985060-0 | Deliveries: 1 | Payload: hello4. Verify redelivery and the dead letter:
curl -X POST "http://localhost:8080/queuedemo/send?consumer=retry-demo-consumer&message=fail-once:1"curl -X POST "http://localhost:8080/queuedemo/send?consumer=order-events-consumer&message=poison:1"The fail-once:1 entry fails transiently on first delivery, is reclaimed, and then appears in /tmp/redis-queue-retry.txt with Deliveries: 2. The poison:1 entry appears in the dead-letter stream:
curl "http://localhost:8080/queuedemo/dlq?stream=order-events-dlq"The runnable sample ships a test-app.sh that automates all of the above: happy-path ack, redelivery, and dead-letter assertions with stream cleanup.
Next steps
Section titled “Next steps”- Read the Redis module for the connection templates (
PhpRedisTemplate, cluster, sentinel) behind the queue consumers. - Compare with the SQS consumer and Kafka consumer: same auto-start worker shape, different delivery guarantees.
- Browse all Libraries when you need S3, OpenSearch, or Doctrine in the same app.