New to Kestra?
Use blueprints to kickstart your first workflows.
Audits Apache Kafka topics for batch compression efficiency, alerting engineering teams when producers emit uncompressed high-bandwidth payloads.
In high-scale event streaming pipelines powered by Apache Kafka, network egress fees and broker storage costs are directly tied to payload sizes. High-throughput producers emitting JSON or Avro events without compression saturate cross-availability-zone (AZ) network links and cause premature broker disk exhaustion.
While Kafka natively supports high-performance codecs like lz4, zstd, and snappy, misconfigured producer clients often default to compression.type=none or omit linger.ms, producing uncompressed single-record packets that waste significant infrastructure spend.
This blueprint implements an automated compression sentinel that connects to target Kafka brokers, samples active message batches, evaluates payload compressibility using in-memory compression analysis, and alerts data platform engineers when topics fall below the efficiency threshold.
periodic_compression_audit trigger (io.kestra.plugin.core.trigger.Schedule) runs every 4 hours.sample_kafka_messages task (io.kestra.plugin.kafka.Consume) consumes up to 50 records from the latest offsets without advancing consumer commit offsets.evaluate_compression_efficiency task (io.kestra.plugin.scripts.python.Script) measures uncompressed bytes versus zlib-compressed size and derives the effective compression ratio.evaluate_inefficiency_condition flowable task (io.kestra.plugin.core.flow.If) evaluates if the ratio is below min_compression_ratio_threshold.notify_slack_compression_gap (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) dispatches actionable tuning recommendations.export_compression_manifest task records topic metrics for long-term reporting.compression.type, linger.ms).enable.auto.commit=false.Monitoring topic compression usually requires deploying custom Kafka Connect audit plugins or running external monitoring services. Kestra unifies native Kafka connectivity, isolated container execution, and workflow alerting into a declarative YAML flow that runs reliably on schedule.
| Name | Type | Default | Description |
|---|---|---|---|
bootstrap_servers |
STRING | kafka.internal:9092 |
Kafka broker connection endpoint. |
topic_name |
STRING | telemetry-events-raw |
Topic name to inspect. |
group_id |
STRING | kestra-compression-sentinel |
Consumer group identifier. |
min_compression_ratio_threshold |
FLOAT | 1.5 |
Minimum acceptable compression ratio (e.g. 1.5x). |
slack_channel |
STRING | #data-platform |
Slack channel destination for alerts. |
{{ outputs.evaluate_compression_efficiency.vars.compression_ratio }}: Calculated compression ratio of sampled messages.{{ outputs.evaluate_compression_efficiency.vars.is_inefficient }}: Boolean flag indicating if ratio failed threshold.{{ outputs.evaluate_compression_efficiency.vars.raw_bytes_kb }}: Total uncompressed batch size in kilobytes.{{ outputs.evaluate_compression_efficiency.outputFiles['compression_report.json'] }}: Complete JSON diagnostic report.SLACK_WEBHOOK_URL: Slack Incoming Webhook endpoint URL.SLACK_WEBHOOK_URL in your Kestra namespace secrets.auto.offset.reset is set to latest and no new messages arrive during the sampling window, the sampled message array will be empty; test on topics with steady message ingress.min_compression_ratio_threshold accordingly.io.kestra.plugin.core.flow.EachParallel to audit entire cluster namespaces.io.kestra.plugin.core.http.Request.