Production deployment guide
This guide covers what to set up when the bundle runs in production: workers, retries and failed messages, idempotency, the transactional outbox, health checks, observability and message versioning. Most of it is Symfony Messenger configuration; the bundle-specific parts are pointed out.
Workers
Asynchronous commands and events are consumed by regular Messenger workers.
The bundle registers every handler without an explicit bus on the sync bus of
its type and on the configured async bus (buses.command_async,
buses.event_async). A worker therefore finds the handler on the bus the
message was sent from; you do not need --bus.
Consume the transports your async messages are sent to: the names under
transports.command_async / transports.event_async, the transport of
#[Asynchronous] (default async), and your framework.messenger.routing.
A handler declared with an explicit bus
(#[AsCommandHandler(ShipOrder::class, bus: 'messenger.bus.commands')]) is only
registered on that bus. If the message is also dispatched asynchronously,
repeat the attribute for the async bus, otherwise the worker fails with
No handler for message.
Recommended flags
| Flag | Purpose |
|---|---|
--time-limit=3600 |
Stop the worker after one hour; the process manager starts a fresh one |
--memory-limit=256M |
Stop when memory usage exceeds the limit |
--limit=1000 |
Stop after handling this many messages |
--sleep=1 |
Seconds to wait when no message is available |
Run separate workers per transport when you want to scale commands and events independently.
Supervisor
; /etc/supervisor/conf.d/cqrs-workers.conf
[program:cqrs-command-worker]
command=php /var/www/app/bin/console messenger:consume async_commands --time-limit=3600 --memory-limit=256M
autostart=true
autorestart=true
numprocs=2
process_name=%(program_name)s_%(process_num)02d
stdout_logfile=/var/log/supervisor/cqrs-command-worker.log
stderr_logfile=/var/log/supervisor/cqrs-command-worker-error.log
user=www-data
[program:cqrs-event-worker]
command=php /var/www/app/bin/console messenger:consume async_events --time-limit=3600 --memory-limit=256M
autostart=true
autorestart=true
numprocs=1
process_name=%(program_name)s_%(process_num)02d
stdout_logfile=/var/log/supervisor/cqrs-event-worker.log
stderr_logfile=/var/log/supervisor/cqrs-event-worker-error.log
user=www-data
systemd
; /etc/systemd/system/cqrs-command-worker@.service
[Unit]
Description=CQRS Command Worker %i
After=network.target
[Service]
Type=simple
User=www-data
ExecStart=/usr/bin/php /var/www/app/bin/console messenger:consume async_commands --time-limit=3600 --memory-limit=256M
Restart=always
RestartSec=5
[Install]
WantedBy=multi-user.target
# Start two worker instances
systemctl enable --now cqrs-command-worker@1
systemctl enable --now cqrs-command-worker@2
Deployments and shutdown
- With the
pcntlextension,SIGTERM,SIGINTandSIGQUITmake a worker finish the current message and exit.--time-limit,--memory-limitand--limitstop it the same way. - After deploying new code, run
bin/console messenger:stop-workersso running workers exit after their current message and restart with the new code. - Symfony resets services tagged
kernel.resetbetween messages (unless you pass--no-reset). The bundle'sCausationIdContextis one of them, so a causation id never leaks from one message to the next.
Retries and failed messages
Transport-level retries
Messenger retries a failed message on the transport it was received from. To drive those retries per message class, combine three settings:
# config/services.yaml
services:
app.retry.payment:
class: SomeWork\CqrsBundle\Policy\ExponentialBackoffRetryPolicy
arguments:
$maxRetries: 5
$initialDelay: 1000 # milliseconds
$multiplier: 2.0
# config/packages/messenger.yaml
framework:
messenger:
failure_transport: failed
transports:
async_commands:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
retry_strategy: # used for messages without a RetryConfiguration
max_retries: 3
delay: 1000
multiplier: 2
failed: 'doctrine://default?queue_name=failed'
# config/packages/somework_cqrs.yaml
somework_cqrs:
retry_policies:
command:
map:
App\Application\Command\ProcessPayment: app.retry.payment
retry_strategy:
transports:
async_commands: command
jitter: 0.1
max_delay: 60000
retry_strategy.transportsreplaces the retry strategy ofasync_commandswithCqrsRetryStrategy. For each failed message it resolves theretry_policies.commandentry.- When the policy implements
RetryConfiguration(asExponentialBackoffRetryPolicydoes), its values apply: here 5 retries with delays of about 1, 2, 4, 8 and 16 seconds, varied by up to 10% and capped at 60 seconds. - Other messages use the transport's own
retry_strategy. Without one, Messenger's defaults apply: 3 retries, 1000 ms delay, multiplier 2. ExponentialBackoffRetryPolicyadds no stamps when the message is dispatched. Withoutretry_strategy.transportsits values are not used at all.
See Retry strategy bridge for the details and for custom policies.
Failure transport
After the last retry, Messenger moves the message to the failure_transport
(failed above). The bundle does not change this pipeline.
# List failed messages, or count them by class
bin/console messenger:failed:show
bin/console messenger:failed:show --stats
# Show one message, retry it or remove it
bin/console messenger:failed:show 42
bin/console messenger:failed:retry 42
bin/console messenger:failed:remove 42
# Retry all failed messages interactively
bin/console messenger:failed:retry
bin/console messenger:stats shows how many messages wait in each transport
(for transports that can count them). A growing failure transport means handlers
keep failing after their retries; the custom health check in
Health checks turns that into a warning.
Idempotency
IdempotencyStamp deduplication relies on Messenger's deduplicate middleware.
It only works when:
- symfony/lock is installed;
- the lock component is enabled (
framework.lock), with a store that all application and worker processes share and that keeps locks until their TTL expires, such as Redis or a database.
framework:
lock: '%env(LOCK_DSN)%' # e.g. redis://redis:6379 or the DSN of your database
somework_cqrs:
idempotency:
ttl: 3600 # seconds a key stays locked
Do not rely on the local flock or semaphore stores for deduplication: they
release a lock as soon as Messenger's lock object is gone, so a second dispatch
with the same key goes through.
Missing pieces do not break the build. In debug mode the reason is written to
the container compilation log (var/cache/<env>/*Compiler.log), for example:
Idempotency is enabled but Messenger's deduplicate middleware is not registered, so DeduplicateStamp is not enforced. Enable the lock component ("framework.lock").
For a synchronous dispatch the key stays locked for ttl seconds after the
message was handled; a duplicate makes dispatchSync() / ask() throw
DuplicateMessageException. If the handler fails, the lock is released so the
caller can retry. For an asynchronous dispatch, the lock is released once a
worker handled the message. See Idempotency.
Outbox operations
With outbox.enabled: true, you store messages in the outbox table inside your
database transaction and a relay sends them to Messenger afterwards. See
Transactional outbox for storing messages.
Table
The table is created on first use (auto_setup: true), but never inside an open
transaction: storing the first message inside a transaction throws a
LogicException if the table does not exist yet or lacks the columns of this
version (a table of 0.4). Outside a transaction, storing and the relay add the
missing columns, but leave the indexes to the setup command (the relay and the
health check report it until it has run; missing columns are critical). Create
the table, or upgrade one of an earlier version, before deploying the code:
If Doctrine migrations manage your schema, set outbox.auto_setup: false. With
doctrine/orm installed, doctrine:migrations:diff includes the outbox table of
the configured connection.
Run somework:cqrs:outbox:setup over a direct database connection: it holds a
session lock, which a pooler in transaction mode (PgBouncer) would move to
another client. The automatic setup is safe behind such a pooler: on PostgreSQL
it runs in one transaction.
The relay only decodes rows with a valid signature (outbox.signing), but
whoever writes to the table can still delay, redirect or drop messages: give the
application a role that can only read and write rows, run the setup with a role that may change the schema, and see
Security for the serializer and the Symfony version to use.
Relay
somework:cqrs:outbox:relay sends up to --limit (default 100) due rows (per
transport, new rows in the order of their created_at and id, then retries in the
order of their retry time; see Delivery guarantees
for what that order guarantees) and marks each one published after dispatching it.
- Rows are dispatched on the Messenger bus of their type (the async command or
event bus when configured, otherwise the sync one; the default bus for other
messages) with their stored transport name as
TransportNamesStamp; rows without a transport name followframework.messenger.routing. Workers then hand each message to the bus where its handlers are registered. A row that is not sent to any transport is handled synchronously, and the command prints and logs a warning. - The relay does not run the stamp pipeline. A message stored through the buses
(
DispatchMode::OUTBOX,#[Outbox],dispatch_modes) went through it when it was stored, as an asynchronous dispatch, and its stamps are stored with it. A message stored withOutboxWriter::store()never goes through it: add the stamps you need to the envelope you store. Only a message stored while a handler runs then gets aMessageMetadataStampwithout one being passed (it continues the handled message's correlation); outside a handler, pass one yourself if you need it. - A row that fails is logged, postponed (1 minute, doubling up to 1 hour) and
makes a single run exit with
1(--watchgoes on, and its exit code does not change); the rows behind it are not blocked. Afteroutbox.max_attemptsattempts (default 10) the relay gives up on the row. - The transports take turns, the one whose next row has waited longest first, so one transport's backlog does not hold up the others.
- A transport that fails 3 times in a row with a
TransportException(broker down, or rejecting messages) is paused until the next run (with--watch, for 30 seconds, doubling up to 5 minutes until it accepts a message), while the rows of the other transports are relayed (10 times, or 3 times taking more than 10 seconds, when it accepted a message earlier in the run: it is up and only rejects some messages). Its rows get three timesmax_attempts(about a day) before they are given up. If the database fails, the run stops right away. - Delivery is at least once: if the process stops between dispatching a row and marking it published, the row is sent again. Make handlers idempotent.
- SIGTERM and SIGINT (with the
pcntlextension) let the relay finish the current row, then a single run exits with1, and--watchwith0. A deploy or a container stop therefore does not leave a row half done. A send blocked on the network ends only with the transport's timeout, and a wait for another process upgrading the table (at most 30 seconds) ends first; a second signal stops the relay at once. - When symfony/lock is installed, only one relay runs at a time; a second one
prints "Another outbox relay is already running." and exits with
0(--wait-for-lock=<seconds>waits for the lock first, then exits with3;--watchwaits as long as it runs). The lock useslock.factorywhenframework.lockis enabled. Otherwise it is a local lock, which only protects relays on the same host. The lock name includesframework.cache.prefix_seed; set it to a stable value when every release is deployed to a new directory, so old and new relays share the lock.
Run the relay from cron:
# crontab: relay every minute, purge published rows every night
* * * * * cd /var/www/app && php bin/console somework:cqrs:outbox:relay --limit=500
0 3 * * * cd /var/www/app && php bin/console somework:cqrs:outbox:purge --older-than="7 days"
or keep it running with --watch when a minute of latency is too much:
[program:cqrs-outbox-relay]
command=php /var/www/app/bin/console somework:cqrs:outbox:relay --watch --time-limit=3600
autostart=true
autorestart=true
; exits with 1 when the database or the lock store fails: restarted at once, about every
; second while it is down (startsecs=0: never FATAL); it resumes when the database is back,
; or when the relay lock expires (up to 60 s) if the lock store is that database
startsecs=0
stopsignal=TERM
stopwaitsecs=30
user=www-data
What the watching relay does:
- It looks for due rows every
--sleepseconds (default 1) and relays them in runs of--limit.--time-limitrestarts it now and then, likemessenger:consume, and lets a deploy's new code take over. - It holds the relay lock while it runs. A second watcher (another server, or the new process of a deploy while the old one finishes its row) prints "Waiting for the relay lock held by another relay." and waits: it takes over when the first one stops. Only the relay that holds the lock works, so more watchers add failover, not throughput.
- After each row that its own process handled (no transport,
sync://), it resets the services, as a worker does between messages: Doctrine's entity managers are cleared, and a closed one is reset (--no-resetturns that off). - A transport that keeps failing is left alone for 30 seconds, doubling up to 5 minutes, instead of being tried again every second.
- A handler run in the relay's process that leaves a transaction open on the outbox connection fails its row: the relay rolls that transaction back, with what the handler wrote in it, instead of writing its next rows into it.
- It exits with
0when a signal or--time-limitstops it: restart it whatever its exit code (autorestart=true, systemdRestart=always), not only on failure. - When the database or the lock store fails, it exits with
1: the process manager restarts it (with the settings above, about every second until the database is back), and the rows wait in the table meanwhile. When the lock store is that same database (framework.lockon its DSN), the stopped relay cannot release its lock: the new one waits for it to expire, up to 60 seconds after the database is back. A lock store of its own (Redis) avoids that wait. - With Doctrine's
auto_commit: false, the relay commits after each fetch, so an idle watcher holds no snapshot and no locks. DBAL starts the next transaction right after each commit, though: PostgreSQL shows the connection asidle in transaction, andidle_in_transaction_session_timeoutends it (the relay exits with1and is restarted). Give the relay a connection with auto-commit, or exempt its database user from that timeout.
Purge
somework:cqrs:outbox:purge --older-than="7 days" deletes rows published before
the given age (a relative date such as "12 hours"; default 7 days).
Unpublished rows are never deleted.
Given-up rows
Monitor the health check and the given-up rows:
bin/console somework:cqrs:outbox:failed # what the relay gave up on, and why
bin/console somework:cqrs:outbox:failed --requeue # after fixing the cause
bin/console somework:cqrs:outbox:failed --delete <id> # a row that must not be sent
Rows that failed and wait for another attempt are not listed: the relay logs each failed
attempt with the row's transport and message type (transport, type), and the health check
warns while such rows keep failing.
A broker outage uses up attempts slowly: rows whose transport fails get three
times max_attempts (30 attempts by default, about a day of retries), and each
run tries 3 rows of a failing transport (3 of the rows stored for it by name, and
3 of the rows without a transport name routed to it), new ones first. Once the
outage is over, requeue the rows it gave up on. The health check warns while
rows keep failing, long before they are given up.
Health checks
somework:cqrs:health checks that the CQRS infrastructure can start:
- handler: every CQRS handler service is instantiated. A handler that cannot
be built (missing environment variable, failing constructor) is
CRITICAL; no handlers at all is aWARNING. - transport: every Messenger transport is instantiated, which validates its
DSN and options. Most built-in transports connect on first use, so this does
not need the broker. The Redis transport connects when it is created unless
its
lazyoption istrue: an unreachable Redis server is thenCRITICAL, after the transport'stimeout, which is unlimited by default (set it, e.g.?timeout=2, orlazy=true, when a probe runs the check). - outbox (when
outbox.enabled): aWARNINGwhen the relay gave up on rows, when failed rows wait for another attempt and the oldest was stored more than 10 minutes ago, when due rows have waited more than 10 minutes, or when the table needssomework:cqrs:outbox:setup(e.g. its index is missing);CRITICALwhen the table cannot be read, or when it lacks the columns of this version (storing a message inside a transaction fails until the setup command has run).
The command prints a table of results and exits with the highest severity:
0 OK, 1 warnings, 2 critical. A checker that throws is reported as
CRITICAL.
Use it in a deployment pipeline or as a probe. Container orchestrators treat any non-zero exit code as a failure; to fail only on critical issues:
Custom checks
Implement SomeWork\CqrsBundle\Health\HealthChecker (@api). With
autoconfiguration the service is tagged somework_cqrs.health_checker and its
results appear in the command:
<?php
declare(strict_types=1);
namespace App\Health;
use SomeWork\CqrsBundle\Health\CheckResult;
use SomeWork\CqrsBundle\Health\CheckSeverity;
use SomeWork\CqrsBundle\Health\HealthChecker;
use Symfony\Component\DependencyInjection\Attribute\Autowire;
use Symfony\Component\Messenger\Transport\Receiver\MessageCountAwareInterface;
final class FailedMessagesChecker implements HealthChecker
{
public function __construct(
#[Autowire(service: 'messenger.transport.failed')]
private readonly object $failureTransport,
) {
}
public function check(): array
{
if (!$this->failureTransport instanceof MessageCountAwareInterface) {
return [];
}
$count = $this->failureTransport->getMessageCount();
return [new CheckResult(
$count > 0 ? CheckSeverity::WARNING : CheckSeverity::OK,
'failed_messages',
sprintf('%d message(s) in the failure transport', $count),
)];
}
}
Observability
Logs
The bundle logs through the logger service, on its own cqrs channel when MonologBundle is
installed. Route or silence it like any channel:
# config/packages/monolog.yaml
monolog:
handlers:
cqrs:
type: stream
path: '%kernel.logs_dir%/cqrs.log'
channels: [cqrs]
Warnings to watch for: an asynchronous dispatch without a transport (Messenger handles the
message in the calling process), an event with handlers that a worker received on a bus without them (it is
acknowledged without being handled), and the outbox relay's failures, paused transports and
given-up messages. Every dispatch through the command and event buses logs a debug line with the
bus, the dispatch mode and the classes of the stamps (a dispatch through the outbox logs a
second one once the message is stored, with its transports); ask() logs the number of stamps
before the query is handled, and a second line once it was handled.
Correlation and causation ids
Every message dispatched through the facades gets a MessageMetadataStamp from
the default metadata provider (a provider of your own that returns
null leaves the message without one), with three ids:
- the message id, unique per message (a retry keeps it);
- the correlation id of the flow: the first message uses its own message id, and every message a handler dispatches inherits the correlation id of the message being handled, so one request shares one correlation id;
- the causation id: the message id of the message whose handler dispatched this one (null for the first message).
Group your logs by correlation id to see a whole flow, and follow the causation ids to rebuild its tree of messages.
Read the stamp in a handler through EnvelopeAware:
<?php
declare(strict_types=1);
namespace App\Application\Command;
use Psr\Log\LoggerInterface;
use SomeWork\CqrsBundle\Attribute\AsCommandHandler;
use SomeWork\CqrsBundle\Contract\EnvelopeAware;
use SomeWork\CqrsBundle\Contract\EnvelopeAwareTrait;
use SomeWork\CqrsBundle\Stamp\MessageMetadataStamp;
#[AsCommandHandler(ProcessPayment::class)]
final class ProcessPaymentHandler implements EnvelopeAware
{
use EnvelopeAwareTrait;
public function __construct(
private readonly PaymentGateway $gateway,
private readonly LoggerInterface $logger,
) {
}
public function __invoke(ProcessPayment $command): mixed
{
$metadata = $this->getEnvelope()->last(MessageMetadataStamp::class);
$this->logger->info('Processing payment', [
'message_id' => $metadata?->getMessageId(),
'correlation_id' => $metadata?->getCorrelationId(),
'causation_id' => $metadata?->getCausationId(),
'payment_id' => $command->paymentId,
]);
$this->gateway->charge($command->paymentId);
return null;
}
}
To continue a correlation id that came with a request, pass your own stamp (one
per dispatch: the stamp also carries the message id); a MessageMetadataStamp
from the caller is kept, and the messages its handlers dispatch inherit its
correlation id:
<?php
use SomeWork\CqrsBundle\Bus\DispatchMode;
use SomeWork\CqrsBundle\Stamp\MessageMetadataStamp;
$correlationId = $request->headers->get('X-Correlation-Id');
$stamps = null !== $correlationId && '' !== $correlationId
? [new MessageMetadataStamp($correlationId)]
: [];
$commandBus->dispatch(new ProcessPayment($paymentId), DispatchMode::DEFAULT, ...$stamps);
OpenTelemetry
The bundle traces messages when open-telemetry/api (1.8 or newer) is installed
and the container has an OpenTelemetry\API\Trace\TracerProviderInterface
service. The service must exist when the container is compiled; otherwise the
middleware is not registered. For example, to use the global tracer provider set
up by the OpenTelemetry SDK:
# config/services.yaml
services:
OpenTelemetry\API\Trace\TracerProviderInterface:
factory: ['OpenTelemetry\API\Globals', 'tracerProvider']
What you get:
cqrs.dispatch <ShortClassName>spans (kindPRODUCER) where messages are dispatched; for synchronous dispatches the span also covers the handlers;cqrs.consume <ShortClassName>spans (kindCONSUMER) in the worker;- attributes
cqrs.message.classandcqrs.message.type, and statusERRORwith the recorded exception when handling fails; - trace propagation: the dispatch adds a
TraceContextStampwith the W3C trace headers, and the worker's span continues that trace.
See Middleware: OpenTelemetryMiddleware for the details.
Personal data
Messages often carry personal data, and the bundle keeps or passes on parts of them:
- Outbox rows. The body (the whole serialized message) and
last_errorstay in the table until the row is purged. Published rows are only deleted bysomework:cqrs:outbox:purge: schedule it with an--older-thanthat fits your retention policy. Rows the relay gave up on are never purged; delete them withsomework:cqrs:outbox:failed --delete <id>…once handled (for example to answer an erasure request), after finding them by id or with SQL. - Error texts. Exception messages can contain personal data. They end up in
last_error, in the relay's output and logs, in the output ofoutbox:failed, and, with OpenTelemetry, in the span status and exception events sent to your tracing backend. - Idempotency keys. An
IdempotencyStampkey is stored in the lock store (e.g. Redis) for its TTL, written to the debug log of the bundle, included in the message ofDuplicateMessageException(and so in error trackers) and serialized with the message. Do not put personal data such as e-mail addresses in keys; hash client-supplied values (hash('sha256', $tenantId.':'.$requestId)). - Metadata. Values returned by your
MessageMetadataProviderare serialized with every message: they are stored in the transports, the outbox table and failure transports, and appear wherever messages are dumped or logged whole. The bundle's own spans only carry the message class and type, and its log lines the stamp classes, not their values; keep the extras to ids anyway.
Message versioning
Messages waiting in a transport were serialized with the old version of their class. Plan changes to message classes with that in mind.
Renaming or moving a class breaks the messages already queued: the serialized data refers to the old class name. Drain the queue before deploying the rename, or keep the old class until no message of it is left:
# Consume what is left, then deploy the rename
bin/console messenger:consume async_commands --time-limit=300
Removing or renaming a property loses the data of queued messages or makes them fail to decode.
Stamps are serialized too. A worker running an older version of a library
cannot decode a message carrying a stamp class that version does not have: with
OpenTelemetry enabled, this bundle adds TraceContextStamp since 0.5, so 0.4
workers must be stopped before 0.5 code dispatches. Deploy workers before (or
with) the code that dispatches, and roll back only once the queues hold no
message of the newer version.
Adding a property depends on the serializer:
- Messenger's default PHP serializer restores objects without calling the
constructor. A new promoted property stays uninitialized in queued messages,
and reading it throws an
Error, even when the constructor parameter has a default value. - The Symfony Serializer (
messenger.transport.symfony_serializer) creates the object through its constructor, so a new constructor parameter with a default value is safe.
The outbox relay decodes stored rows with outbox.serializer
(messenger.default_serializer by default), so the same rules apply to rows
waiting in the outbox table.