New to Kestra?
Use blueprints to kickstart your first workflows.
Automatically triage Kafka DLQ events with Kestra, safely replay transient failures to source topics, and quarantine poisoned records into Amazon S3.
Diagram unavailable
We could not build the topology for this blueprint. The flow itself is valid, use the YAML on the left to run it.
Dead-Letter Queues (DLQs) in Apache Kafka prevent consumers from stalling when they encounter unexpected or corrupted records. However, without automated triage, messages in DLQs accumulate into dark data lakes or risk triggering secondary incident storms when operators attempt manual, unverified replays.
This blueprint implements an automated Dead-Letter Queue management and remediation engine. It intercepts failed messages in real time, classifies their error signatures into transient vs unrecoverable poison-pill categories, applies an exponential retry budget backed by Kestra namespace KV storage, safely republishes recoverable messages to the primary topic, and quarantines toxic records directly into partitioned Amazon S3 buckets while paging data engineers on Slack with forensic diagnostics.
orders-dlq) and target ingestion topic (e.g., orders-raw).Configure the following secrets in your Kestra namespace:
KAFKA_BOOTSTRAP_SERVERS: Kafka broker connection string (e.g. broker1:9092,broker2:9092).AWS_ACCESS_KEY_ID: IAM access key with s3:PutObject permissions for the quarantine bucket.AWS_SECRET_ACCESS_KEY: IAM secret key matching the access key.SLACK_WEBHOOK_URL: Slack Incoming Webhook URL for quarantine alerts and pipeline errors.DLQ_REPLAY_WEBHOOK_KEY: Secret authentication key for manual replay webhook trigger.Adjust flow variables to match your streaming infrastructure:
dlq_topic: The dead-letter topic being monitored (default: orders-dlq).replay_topic: The primary destination topic for replayed messages (default: orders-raw).max_retry_attempts: The maximum number of transient replay cycles allowed before quarantine (default: 3).s3_quarantine_bucket: Target S3 bucket for unrecoverable payloads (default: company-dlq-quarantine-archive).s3_region: AWS region for the S3 bucket (default: us-east-1).┌────────────────────────┐
│ Kafka DLQ / Webhook │
└───────────┬────────────┘
│
▼
┌────────────────────────┐
│ Extract & Classify │
└───────────┬────────────┘
│
▼
┌────────────────────────┐
│ Check KV Retry Budget │
└───────────┬────────────┘
│
[Transient & Attempts < Max?]
/ \
YES NO (Poison / Budget Spent)
/ \
▼ ▼
┌─────────────┐ ┌────────────────────────┐
│ Bump KV TTL │ │ Generate S3 Manifest │
└──────┬──────┘ └───────────┬────────────┘
│ │
▼ ▼
┌─────────────┐ ┌────────────────────────┐
│ Replay to │ │ Upload to S3 Cold │
│ Topic │ │ Storage │
└──────┬──────┘ └───────────┬────────────┘
│ │
▼ ▼
┌─────────────┐ ┌────────────────────────┐
│ Log Success │ │ Alert Slack with Link │
└─────────────┘ └────────────────────────┘
on_dlq_event (io.kestra.plugin.kafka.RealtimeTrigger) subscribes to the DLQ topic and creates an execution immediately upon record receipt.extract_event_payload unifies record contents and inspects error strings to tag the event as TRANSIENT or POISON_PILL.read_replay_counter queries the Kestra KV store for dlq_retry_<record_id> to read historical replay attempts.route_dlq_policy (io.kestra.plugin.core.flow.If) evaluates the classified error type against the maximum allowed retry limit.record_retry_attempt (io.kestra.plugin.core.kv.Set) increments the retry counter with a 1-day TTL.replay_to_kafka_topic (io.kestra.plugin.kafka.Produce) publishes the record back into the raw ingestion topic.log_replayed_record (io.kestra.plugin.core.log.Log) records the successful replay event in the execution log.generate_quarantine_manifest formats the corrupt payload, execution ID, and timestamp into quarantine_record.json.upload_quarantine_to_s3 (io.kestra.plugin.aws.s3.Upload) securely archives the artifact into S3.alert_quarantined_message (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) dispatches a high-priority Slack alert with the S3 URI.s3://<bucket>/quarantine/YYYY/MM/DD/quarantine-<execution_id>.json.errors handler executes immediately, delivering a critical incident alert to Slack with the execution URL.io.kestra.plugin.core.flow.Pause or schedule-based replay to enforce an escalating backoff window between retry attempts.aws.s3.Upload with gcp.gcs.Upload or azure.storage.blob.Upload for multi-cloud parity.