QuotaDescribe icon
OutputValues icon
Log icon
If icon
SlackIncomingWebhook icon
Loop icon
QuotaAlter icon
Fail icon
Schedule icon

Kafka Client Quota Guard

Audit Kafka client-id quotas against a policy and a ceiling, alert Slack, and enforce them with QuotaAlter behind dry-run gates, then verify.

Categories
DataInfrastructure

Kafka quotas are what stop one client from taking the whole broker. They get lost the same way other configs do. A new service ships without its quota. Someone raises a notebook's consumer rate to 100 MB/s for a one-off backfill and leaves it. A quota set by hand differs from the one in the capacity plan. The first sign is usually every other producer timing out.

This blueprint reads every client-id quota on the cluster with the AdminClient and checks it against a declared policy and a ceiling that applies to every client. It alerts Slack, and can apply the policy with QuotaAlter behind two gates, then describes the quotas again to verify.

This blueprint was created by zkasuran.

Findings

Finding Rule Fix applied
NO_QUOTA A client in the policy has no quota on the cluster The policy rates
DRIFT A policy client has a rate that differs from the policy Only the rates that differ
OVER_CEILING Any client, in the policy or not, has a producer or consumer byte rate above ceiling_byte_rate That rate clamped to the ceiling

A client that is right is not touched. A rate not in the policy is never changed or removed.

How it works

  1. describe (io.kestra.plugin.kafka.QuotaDescribe, read-only) lists every quota entity and its rates.
  2. findings compares them with the policy and the ceiling in one jq expression. It keeps, per client, only the rates to set.
  3. summary logs the result every run.
  4. route (io.kestra.plugin.core.flow.If):
    • compliant: logs it.
    • findings: alert posts to Slack.
      • remediate_gate runs only when auto_remediate is true and dry_run is false.
      • remediate (io.kestra.plugin.core.flow.Loop over io.kestra.plugin.kafka.QuotaAlter) sets only the wrong rates. A rate that renders empty is skipped.
      • verify describes again, and verify_result fails if any rate still differs.
  5. errors alerts Slack when the flow fails, for example on an unreachable broker.

What you get

  • One view of every client-id quota on the cluster, checked against the capacity plan.
  • Three findings (missing quota, drifted rate, rate above the ceiling) with the exact rates to set.
  • Fixes limited to the rates that are wrong, behind two gates, verified by a second describe.

Who it's for

  • Platform teams running a shared Kafka cluster for many services.
  • SREs who have seen one client starve the brokers and want a ceiling enforced on every client.
  • Teams evaluating the Kafka plugin who want an example of QuotaDescribe and QuotaAlter.

Why orchestrate this with Kestra

Quotas are set once with kafka-configs.sh and then forgotten. Kestra runs the check on a schedule, keeps the policy next to the flow, turns findings into a Slack alert, gates the change behind auto_remediate and dry_run, and keeps a history of every run, so you can see when a quota changed and who put it back.

Tested end to end

On Kestra 2.0.5 OSS against Apache Kafka 3.9.1 (KRaft), with a local HTTP sink standing in for Slack. Quotas were set beforehand with kafka-configs.sh:

  • etl-loader: below the policy;
  • billing-service: on the policy;
  • adhoc-notebook: consumer rate of 100 MB/s, above the 50 MB/s ceiling;
  • analytics-reader: no quota at all.
Run Inputs Result
1 defaults 3 findings: etl-loader DRIFT, analytics-reader NO_QUOTA, adhoc-notebook OVER_CEILING. billing-service not reported. Alert sent
2 auto_remediate: true, dry_run: true Same report. kafka-configs.sh shows nothing changed
3 auto_remediate: true, dry_run: false 3 QuotaAlter calls, re-describe passes. kafka-configs.sh shows etl-loader at 5 MB/s and 10 MB/s, analytics-reader at 20 MB/s and adhoc-notebook at 50 MB/s
4 defaults again Compliant, no alert
5 unreachable broker AdminClient timeout, the errors alert is sent

Prerequisites

  • A Kafka cluster reachable from the Kestra worker, and a principal allowed to DescribeConfigs and, for remediation, AlterConfigs on the cluster.
  • Local testing: docker run -d --name kafka -p 9092:9092 apache/kafka:3.9.1, then set a few quotas with kafka-configs.sh --alter --entity-type clients.

Secrets

  • SLACK_WEBHOOK_URL: Slack incoming webhook used by alert and the errors block. In Kestra OSS, provide it as the base64-encoded environment variable SECRET_SLACK_WEBHOOK_URL.

Inputs

Input Default Purpose
bootstrap_servers localhost:9092 Brokers
policy 3 example clients Required rates per client-id
ceiling_byte_rate 52428800 (50 MB/s) Highest rate any client may have
auto_remediate false Allow QuotaAlter
dry_run true Must be false for a change to apply

The hourly trigger is shipped disabled and always runs in safe mode.

Quick start

  1. Add the SLACK_WEBHOOK_URL secret.
  2. Import the flow, set bootstrap_servers, and replace policy with your clients and rates.
  3. Run it once with the defaults and read the findings.
  4. Run it with auto_remediate: true and dry_run: false to apply them.
  5. Enable the hourly trigger, which always runs in safe mode.

How to extend

  • Add user and IP quotas with the entityUser and entityIp properties.
  • Load the policy from a Git repository or a namespace file so it is reviewed like code.
  • Compare quotas with actual traffic from your metrics backend and flag clients that never come near theirs.
  • Open a ticket instead of fixing when the client belongs to another team.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.