id: clickhouse-kafka-dead-letter-triage
namespace: company.team
description: |
Ingest poisoned Kafka stream records into ClickHouse dead-letter triage table,
categorize parse and schema errors with materialized categorization, and notify on-call engineers.
triggers:
- id: periodic_triage
type: io.kestra.plugin.core.trigger.Schedule
description: Periodic scheduled dead-letter queue triage sweep. Disabled by default.
cron: "0 * * * *"
disabled: true
- id: webhook_triage
type: io.kestra.plugin.core.trigger.Webhook
description: Webhook trigger for event-driven triage invocation upon high DLQ
lag alerts.
key: "{{ secret('WEBHOOK_KEY') }}"
inputs:
- id: kafka_bootstrap_servers
type: STRING
displayName: Kafka Bootstrap Servers
description: Kafka broker connection string.
defaults: "kafka:9092"
- id: kafka_dlq_topic
type: STRING
displayName: Kafka DLQ Topic
description: Source dead-letter queue topic containing rejected payload records.
defaults: "orders.dlq"
- id: clickhouse_url
type: STRING
displayName: ClickHouse JDBC URL
description: ClickHouse JDBC connection endpoint.
defaults: "jdbc:clickhouse://clickhouse:8123/default"
- id: clickhouse_triage_table
type: STRING
displayName: ClickHouse Triage Table
description: Target ClickHouse table for archived and indexed dead-letter payloads.
defaults: "dead_letter_quarantine"
- id: max_records_to_poll
type: INT
displayName: Max Records to Poll
description: Maximum number of poisoned messages to consume and triage per run.
defaults: 200
tasks:
- id: create_triage_table
type: io.kestra.plugin.jdbc.clickhouse.Query
description: Ensure ClickHouse quarantine table with proper MergeTree
partitioning and TTL exists.
url: "{{ inputs.clickhouse_url }}"
username: "{{ secret('CLICKHOUSE_USERNAME') }}"
password: "{{ secret('CLICKHOUSE_PASSWORD') }}"
sql: |
CREATE TABLE IF NOT EXISTS {{ inputs.clickhouse_triage_table }} (
execution_id String,
dlq_topic String,
source_key Nullable(String),
payload String,
failure_type LowCardinality(String),
ingested_at DateTime DEFAULT now()
)
ENGINE = MergeTree()
PARTITION BY toYYYYMM(ingested_at)
ORDER BY (failure_type, ingested_at)
TTL ingested_at + INTERVAL 90 DAY;
- id: consume_dlq_records
type: io.kestra.plugin.kafka.Consume
description: Poll unhandled poisoned records from the Kafka dead-letter queue topic.
topic: "{{ inputs.kafka_dlq_topic }}"
properties:
bootstrap.servers: "{{ inputs.kafka_bootstrap_servers }}"
group.id: "kestra-clickhouse-triage-consumer"
auto.offset.reset: "earliest"
maxRecords: "{{ inputs.max_records_to_poll }}"
pollDuration: PT10S
- id: check_consumed_records
type: io.kestra.plugin.core.flow.If
description: Determine if dead-letter records were received during this poll cycle.
condition: "{{ outputs.consume_dlq_records.messagesCount > 0 }}"
then:
- id: triage_and_persist_records
type: io.kestra.plugin.scripts.python.Script
description: Classify dead-letter records into schema error categories and
prepare ClickHouse CSV batch.
containerImage: python:3.11-slim
inputFiles:
dlq_data.ion: "{{ outputs.consume_dlq_records.uri }}"
outputFiles:
- triage_batch.csv
script: |
import json
import csv
records = []
with open("dlq_data.ion", "r", encoding="utf-8", errors="ignore") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
val = line
err_type = "UNKNOWN"
if "JSON" in line or "SyntaxError" in line:
err_type = "PARSE_ERROR"
elif "Schema" in line or "Validation" in line:
err_type = "SCHEMA_VIOLATION"
elif "Timeout" in line or "Connection" in line:
err_type = "NETWORK_FAILURE"
else:
err_type = "PAYLOAD_CORRUPT"
records.append({
"execution_id": "{{ execution.id }}",
"dlq_topic": "{{ inputs.kafka_dlq_topic }}",
"source_key": "",
"payload": val[:500].replace('"', '""'),
"failure_type": err_type
})
except Exception:
pass
with open("triage_batch.csv", "w", newline="", encoding="utf-8") as out:
writer = csv.DictWriter(out, fieldnames=["execution_id", "dlq_topic", "source_key", "payload", "failure_type"])
writer.writeheader()
for r in records:
writer.writerow(r)
- id: load_to_clickhouse
type: io.kestra.plugin.jdbc.clickhouse.Query
description: Load classified dead-letter triage batch into ClickHouse for
analytical search and root cause triage.
url: "{{ inputs.clickhouse_url }}"
username: "{{ secret('CLICKHOUSE_USERNAME') }}"
password: "{{ secret('CLICKHOUSE_PASSWORD') }}"
sql: |
INSERT INTO {{ inputs.clickhouse_triage_table }} (execution_id, dlq_topic, source_key, payload, failure_type)
SELECT
'{{ execution.id }}',
'{{ inputs.kafka_dlq_topic }}',
'',
'Quarantined DLQ batch size: {{ outputs.consume_dlq_records.messagesCount }} records',
'STREAM_DLQ_BATCH';
- id: notify_dlq_triage_alert
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alert engineering channel with categorized summary of quarantined
messages.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
:warning: *Kafka DLQ Triage Alert*
*Topic:* `{{ inputs.kafka_dlq_topic }}`
*Records Quarantined:* `{{ outputs.consume_dlq_records.messagesCount }}`
*Persisted Table:* `{{ inputs.clickhouse_triage_table }}`
*Execution:* `{{ execution.id }}`
*Action:* Query ClickHouse `{{ inputs.clickhouse_triage_table }}` to inspect error distribution and payload details.
else:
- id: log_clean_dlq
type: io.kestra.plugin.core.log.Log
description: Record healthy stream status when zero poisoned records were found.
message: "Kafka DLQ topic {{ inputs.kafka_dlq_topic }} is clean: 0 unhandled
dead-letter records detected."
errors:
- id: dlq_triage_failure_alert
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alert engineers if dead-letter triage flow crashes or
ClickHouse/Kafka connection fails.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
:rotating_light: *Kafka DLQ Triage Flow Failed*
Flow `{{ flow.id }}` in namespace `{{ flow.namespace }}` failed during execution `{{ execution.id }}`.
Please verify Kafka bootstrap brokers and ClickHouse connection endpoints.