RealtimeTrigger icon
Webhook icon
OutputValues icon
If icon
Set icon
Produce icon
Log icon
Commands icon
Process icon
Upload icon
SlackIncomingWebhook icon

Kafka DLQ Quarantine, Classification, and Safe Replay Gate

Automatically triage Kafka DLQ events with Kestra, safely replay transient failures to source topics, and quarantine poisoned records into Amazon S3.

Categories
CloudDataInfrastructureinfrastructure

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.

Use Cases

  • Automated Stream Self-Healing: Seamlessly recover messages that failed due to temporary network partitions, rate limits, or downstream database locks without human intervention.
  • Poison-Pill Quarantine: Safely isolate malformed, schema-violating, or corrupted records from the streaming pipeline without dropping data.
  • Dead-Letter Observability: Gain end-to-end lineage, failure categorization, and alerting across microservices architectures.
  • Audit Compliance: Maintain an immutable cold-storage audit log in S3 of every quarantined message and associated stack trace.

Prerequisites

  • An Apache Kafka cluster (or Confluent Cloud / Aiven Kafka) with a defined DLQ topic (e.g., orders-dlq) and target ingestion topic (e.g., orders-raw).
  • An Amazon S3 bucket configured for cold-storage quarantine archives.
  • An incoming Slack webhook configured for notifications.

Secrets

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.

Configuration

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).

Workflow

┌────────────────────────┐
│ 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  │
└─────────────┘   └────────────────────────┘

Step-by-Step

  1. Event Interception: on_dlq_event (io.kestra.plugin.kafka.RealtimeTrigger) subscribes to the DLQ topic and creates an execution immediately upon record receipt.
  2. Payload Extraction & Error Taxonomy: extract_event_payload unifies record contents and inspects error strings to tag the event as TRANSIENT or POISON_PILL.
  3. Stateful Retry Tracking: read_replay_counter queries the Kestra KV store for dlq_retry_<record_id> to read historical replay attempts.
  4. Policy Decision Gate: route_dlq_policy (io.kestra.plugin.core.flow.If) evaluates the classified error type against the maximum allowed retry limit.
  5. Transient Remediation Branch:
    • 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.
  6. Poison Quarantine Branch:
    • 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.

Expected Output

  • Recovered Records: Produced to the destination Kafka topic with updated retry counts in Kestra KV store.
  • Quarantined Records: Persisted to S3 under s3://<bucket>/quarantine/YYYY/MM/DD/quarantine-<execution_id>.json.
  • Slack Notifications: Real-time alert messages containing classification rationale, retry count, and S3 URI.

Failure Handling

  • If any internal Kestra task fails (e.g. AWS credential failure or Kafka disconnection), the flow-level errors handler executes immediately, delivering a critical incident alert to Slack with the execution URL.

Customization

  • Exponential Backoff: Introduce a io.kestra.plugin.core.flow.Pause or schedule-based replay to enforce an escalating backoff window between retry attempts.
  • Alternative Cloud Storage: Swap aws.s3.Upload with gcp.gcs.Upload or azure.storage.blob.Upload for multi-cloud parity.
  • Schema Registry Integration: Add an Avro or Protobuf schema validation task using Kestra Python scripts to validate payloads against Confluent Schema Registry before replay.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.