Schedule icon
Webhook icon
Query icon
Consume icon
If icon
Script icon
SlackIncomingWebhook icon
Log icon

Kafka DLQ Ingestion and Error Triage with ClickHouse

Automatically capture, categorize, and archive poisoned Kafka streaming records into ClickHouse for rapid incident diagnosis and post-mortem analysis.

Categories
DataInfrastructure

Dead-letter queues (DLQs) prevent malformed, unparseable, or schema-violating records from halting real-time streaming topologies. However, poisoned messages often pile up in DLQ topics without visibility into failure causes, schema regressions, or volume spikes.

This blueprint implements an automated Kafka dead-letter ingestion and triage pipeline using ClickHouse:

  1. Dead-Letter Polling: Consumes poisoned records from a designated Kafka DLQ topic without acknowledging or advancing offsets until successfully triaged.
  2. Automated Error Classification: Categorizes record errors into actionable failure types (PARSE_ERROR, SCHEMA_VIOLATION, NETWORK_FAILURE, PAYLOAD_CORRUPT).
  3. High-Performance Storage: Writes structured quarantine records into a partitioned ClickHouse MergeTree table with built-in 90-day time-to-live (TTL) retention policies.
  4. Proactive Team Alerting: Dispatches real-time Slack notifications summarizing quarantined volume, error distribution, and remediation queries.

Use Cases

  • Stream Processing Resiliency: Preserve poisoned streaming records for post-mortem root-cause analysis without blocking production consumer groups.
  • Schema Drift Detection: Surface unannounced upstream schema modifications before they impact downstream analytics or data lakes.
  • Audit & Compliance: Retain an immutable audit trail of malformed transaction attempts with automatic TTL expiration.

How It Works

  1. create_triage_table (io.kestra.plugin.jdbc.clickhouse.Query) initializes the destination dead_letter_quarantine table in ClickHouse if it does not already exist, applying monthly partitioning and 90-day automatic data lifecycle purging.
  2. consume_dlq_records (io.kestra.plugin.kafka.Consume) polls up to inputs.max_records_to_poll poisoned messages from inputs.kafka_dlq_topic.
  3. check_consumed_records (io.kestra.plugin.core.flow.If) branches execution based on whether any poisoned records were retrieved:
    • If records are found:
      • triage_and_persist_records analyzes record contents and classifies failure root causes.
      • load_to_clickhouse archives the indexed records into ClickHouse.
      • notify_dlq_triage_alert notifies the engineering channel via Slack.
    • If no records are found, log_clean_dlq logs a clean status confirmation.
  4. dlq_triage_failure_alert catches unhandled exceptions to ensure zero silent failures in monitoring pipelines.

Prerequisites

  • Accessible Kafka cluster with a designated DLQ topic.
  • ClickHouse instance accessible over HTTP/JDBC.
  • Slack Incoming Webhook for operational alerting.

Secrets

  • CLICKHOUSE_USERNAME: ClickHouse database user credentials.
  • CLICKHOUSE_PASSWORD: ClickHouse database password.
  • SLACK_WEBHOOK_URL: Slack Incoming Webhook URL.
  • WEBHOOK_KEY: Shared secret key for triggering on-demand triage runs.

Step-by-Step

┌──────────────────────────────┐
│  Schedule / Webhook Trigger  │
└──────────────┬───────────────┘
               │
               ▼
┌──────────────────────────────┐
│ Ensure ClickHouse Table      │
└──────────────┬───────────────┘
               │
               ▼
┌──────────────────────────────┐
│ Poll Kafka DLQ Topic         │
└──────────────┬───────────────┘
               │
     [Any Records Found?]
           /                     YES         NO
         /                         ▼               ▼
┌───────────────┐ ┌───────────────┐
│ Triage Errors │ │ Log Topic     │
│ & Persist     │ │ Clean Status  │
│ ClickHouse    │ └───────────────┘
└───────┬───────┘
        │
        ▼
┌───────────────┐
│ Dispatch      │
│ Slack Alert   │
└───────────────┘

Expected Output

  • Clean Run: Execution completes with zero messages quarantined and confirmation in logs.
  • Triage Run: Records loaded to ClickHouse triage table with categorized error types, followed by a Slack summary notification.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.