Schedule icon
Consume icon
Script icon
Docker icon
If icon
SlackIncomingWebhook icon
Log icon
Return icon

Kafka Producer Compression Ratio Sentinel

Audits Apache Kafka topics for batch compression efficiency, alerting engineering teams when producers emit uncompressed high-bandwidth payloads.

Categories
DataInfrastructure

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.

How it works

  1. Scheduled Inspection: The periodic_compression_audit trigger (io.kestra.plugin.core.trigger.Schedule) runs every 4 hours.
  2. Message Batch Sampling: The sample_kafka_messages task (io.kestra.plugin.kafka.Consume) consumes up to 50 records from the latest offsets without advancing consumer commit offsets.
  3. Python Compressibility Scoring: The evaluate_compression_efficiency task (io.kestra.plugin.scripts.python.Script) measures uncompressed bytes versus zlib-compressed size and derives the effective compression ratio.
  4. Condition Branching: The evaluate_inefficiency_condition flowable task (io.kestra.plugin.core.flow.If) evaluates if the ratio is below min_compression_ratio_threshold.
  5. Slack Advisory: When inefficient compression is detected, notify_slack_compression_gap (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) dispatches actionable tuning recommendations.
  6. Audit Manifest: The export_compression_manifest task records topic metrics for long-term reporting.

What you get

  • Automated audit of Kafka producer batch compression efficiency without broker downtime.
  • Early warning of uncompressed high-throughput topics causing cloud egress cost spikes.
  • Actionable Slack notifications with exact producer tuning recommendations (compression.type, linger.ms).
  • Non-intrusive consumer implementation with enable.auto.commit=false.

Who it is for

  • Data Platform Engineers managing enterprise Apache Kafka or AWS MSK clusters.
  • Cloud FinOps practitioners optimizing network egress and storage spend.
  • Software Architects designing resilient event streaming schemas.

Why orchestrate this with Kestra

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.

Inputs

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.

Expected outputs

  • {{ 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.

Prerequisites

  • Network access to target Apache Kafka brokers (standard port 9092 or 9094).
  • Read permissions on the specified topic.
  • Slack Incoming Webhook configured in Kestra secrets.

Secrets

  • SLACK_WEBHOOK_URL: Slack Incoming Webhook endpoint URL.

Quick start

  1. Configure SLACK_WEBHOOK_URL in your Kestra namespace secrets.
  2. Import this flow YAML into your Kestra instance.
  3. Click Execute in the UI to run an initial topic compression test.
  4. Review execution outputs to check the calculated compression ratio.

Common pitfalls and troubleshooting

  • Offset Reset: If 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.
  • SASL/SSL Security: For managed clusters like AWS MSK or Confluent Cloud, add required SASL/SSL credentials to the task properties block.
  • Small Payloads: Very small payloads (< 100 bytes) may have low compression ratios due to fixed compression headers; adjust min_compression_ratio_threshold accordingly.

How to extend

  • Scan multiple topics in parallel using io.kestra.plugin.core.flow.EachParallel to audit entire cluster namespaces.
  • Correlate compression ratios with Prometheus broker network metrics via io.kestra.plugin.core.http.Request.
  • Automatically track average message sizes over time to spot payload schema bloat.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.