Chapter 17: Handling Failures with Dead Letter Queues
A production bug in the ComplianceAlertHandler causes it to crash with a
TypeError on deposits where source_type is None, a pre-upcasting edge
case. The handler exhausts its retry attempts and messages pile up in the
dead-letter queue.
When a handler keeps failing, the message ends up in a dead letter queue. This chapter walks the workflow around that: finding those messages, working out what went wrong, and getting them processed.
How DLQ Works
When a StreamSubscription handler throws an exception:
- The message is retried up to
max_retriestimes (default: 3) - Each retry uses exponential backoff (
retry_delay_seconds * 2^N) - After exhausting retries, the message moves to the DLQ stream
- The subscription continues processing the next message
The DLQ preserves the original payload, the failure reason, the retry count, and the timestamp. Nothing is lost.
Discovering the Problem
$ protean dlq list --domain=fidelis
DLQ ID Subscription Failure Reason Failed At Retries
1719234567890-0 fidelis::account TypeError: 'NoneT... 2025-06-16 14:22:00 3
1719234568123-0 fidelis::account TypeError: 'NoneT... 2025-06-16 14:22:01 3
1719234569456-0 fidelis::account TypeError: 'NoneT... 2025-06-16 14:23:15 3
3 DLQ message(s) found.
You can filter by subscription:
$ protean dlq list --subscription=fidelis::account --domain=fidelis
Inspecting a Failed Message
$ protean dlq inspect 1719234567890-0 --domain=fidelis
DLQ ID: 1719234567890-0
Stream: fidelis::account
Failure Reason: TypeError: 'NoneType' object has no attribute 'startswith'
Failed At: 2025-06-16 14:22:00
Retry Count: 3
Payload:
{
"type": "Fidelis.DepositMade.v1",
"data": {
"account_id": "acc-7742",
"amount": 12000.0,
"reference": "DEP-8834"
}
}
The inspection shows everything needed to diagnose the issue:
- The failure reason (
TypeError) tells you what went wrong. - The payload shows the exact message that caused the failure.
- The type (
Fidelis.DepositMade.v1) reveals it was a v1 event. Thesource_typefield is missing because the upcaster has not run yet at the handler level.
Fixing and Replaying
After fixing the handler code (adding a None check for
source_type) and redeploying:
# Replay a single message
$ protean dlq replay 1719234567890-0 --domain=fidelis
Replayed message 1719234567890-0 to stream 'fidelis::account'.
# Replay all failed messages for a subscription
$ protean dlq replay-all --subscription=fidelis::account --domain=fidelis
Replay all DLQ messages for subscription 'fidelis::account'? [y/N]: y
Replayed 3 message(s) to stream 'fidelis::account'.
Replaying puts the message back on the original stream. The handler (now fixed) processes it normally.
Purging Unrecoverable Messages
If messages cannot be fixed (e.g., they reference deleted data):
$ protean dlq purge --subscription=fidelis::account --domain=fidelis
Purge all DLQ messages for subscription 'fidelis::account'? [y/N]: y
Purged 3 message(s) from DLQ.
Warning
purge permanently removes messages. Use it only when you are
certain the messages are unrecoverable or no longer relevant.
Verifying Recovery
$ protean dlq list --domain=fidelis
No DLQ messages found.
The Fix-and-Replay Cycle
This pattern will become your standard operating procedure:
- Discover:
protean dlq listfinds failed messages - Inspect:
protean dlq inspectreveals the cause - Fix: Update handler code and redeploy
- Replay:
protean dlq replay-allreprocesses the messages - Verify:
protean dlq listconfirms the DLQ is empty
DLQ Configuration
The DLQ is configured in domain.toml:
[server.stream_subscription]
max_retries = 3 # Retry before DLQ
retry_delay_seconds = 1 # Base delay (exponential backoff)
enable_dlq = true # Enable dead-letter queue
Setting enable_dlq = false means failed messages are dropped after
exhausting retries. This is almost never what you want in production.
What We Built
- The fix-and-replay cycle for production incident response.
protean dlq list: Discover failed messages.protean dlq inspect: Diagnose the root cause.protean dlq replayandreplay-allreprocess messages after you fix the cause.protean dlq purge: Discard unrecoverable messages.- DLQ configuration for retries and backoff.
Next, we set up proactive monitoring to catch problems before they fill the DLQ.
Full Source
"""Chapter 17: Dead Letter Queues
Demonstrates the domain setup for dead letter queue management. When an event
handler or projector fails repeatedly, the failing message is moved to a dead
letter queue (DLQ). Protean provides CLI commands to inspect, replay, and
purge these messages:
protean dlq list --domain fidelis.domain
protean dlq inspect <message-id> --domain fidelis.domain
protean dlq replay <message-id> --domain fidelis.domain
protean dlq purge --domain fidelis.domain
This chapter is primarily about CLI workflows; the Python source only needs to
set up the domain so the CLI can operate on it.
"""
from protean import Domain, apply, handle, invariant
from protean.exceptions import ValidationError
from protean.fields import Float, Identifier, String
from protean.utils.globals import current_domain
domain = Domain("fidelis")
@domain.event(part_of="Account")
class AccountOpened:
account_id: Identifier(required=True)
account_number: String(required=True)
holder_name: String(required=True)
opening_deposit: Float(required=True)
@domain.event(part_of="Account")
class DepositMade:
account_id: Identifier(required=True)
amount: Float(required=True)
reference: String()
@domain.event(part_of="Account")
class WithdrawalMade:
account_id: Identifier(required=True)
amount: Float(required=True)
reference: String()
@domain.aggregate(event_sourced=True)
class Account:
account_number: String(max_length=20, required=True)
holder_name: String(max_length=100, required=True)
balance: Float(default=0.0)
status: String(max_length=20, default="ACTIVE")
@invariant.post
def balance_must_not_be_negative(self):
if self.balance is not None and self.balance < 0:
raise ValidationError(
{"balance": ["Insufficient funds: balance cannot be negative"]}
)
@classmethod
def open(cls, account_number: str, holder_name: str, opening_deposit: float):
account = cls._create_new()
account.raise_(
AccountOpened(
account_id=str(account.id),
account_number=account_number,
holder_name=holder_name,
opening_deposit=opening_deposit,
)
)
return account
def deposit(self, amount: float, reference: str | None = None) -> None:
if amount <= 0:
raise ValidationError({"amount": ["Deposit amount must be positive"]})
self.raise_(
DepositMade(
account_id=str(self.id),
amount=amount,
reference=reference,
)
)
def withdraw(self, amount: float, reference: str | None = None) -> None:
if amount <= 0:
raise ValidationError({"amount": ["Withdrawal amount must be positive"]})
self.raise_(
WithdrawalMade(
account_id=str(self.id),
amount=amount,
reference=reference,
)
)
@apply
def on_account_opened(self, event: AccountOpened):
self.id = event.account_id
self.account_number = event.account_number
self.holder_name = event.holder_name
self.balance = event.opening_deposit
self.status = "ACTIVE"
@apply
def on_deposit_made(self, event: DepositMade):
self.balance += event.amount
@apply
def on_withdrawal_made(self, event: WithdrawalMade):
self.balance -= event.amount
@domain.command(part_of=Account)
class OpenAccount:
account_number: String(required=True)
holder_name: String(required=True)
opening_deposit: Float(required=True)
@domain.command(part_of=Account)
class MakeDeposit:
account_id: Identifier(required=True)
amount: Float(required=True)
reference: String()
@domain.command(part_of=Account)
class MakeWithdrawal:
account_id: Identifier(required=True)
amount: Float(required=True)
reference: String()
@domain.command_handler(part_of=Account)
class AccountCommandHandler:
@handle(OpenAccount)
def handle_open_account(self, command: OpenAccount):
account = Account.open(
account_number=command.account_number,
holder_name=command.holder_name,
opening_deposit=command.opening_deposit,
)
current_domain.repository_for(Account).add(account)
return str(account.id)
@handle(MakeDeposit)
def handle_make_deposit(self, command: MakeDeposit):
repo = current_domain.repository_for(Account)
account = repo.get(command.account_id)
account.deposit(command.amount, reference=command.reference)
repo.add(account)
@handle(MakeWithdrawal)
def handle_make_withdrawal(self, command: MakeWithdrawal):
repo = current_domain.repository_for(Account)
account = repo.get(command.account_id)
account.withdraw(command.amount, reference=command.reference)
repo.add(account)