New to Kestra?
Use blueprints to kickstart your first workflows.
Audit Kafka client-id quotas against a policy and a ceiling, alert Slack, and enforce them with QuotaAlter behind dry-run gates, then verify.
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.
| 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.
describe (io.kestra.plugin.kafka.QuotaDescribe, read-only) lists every quota entity and its rates.findings compares them with the policy and the ceiling in one jq expression. It keeps, per client, only the rates to set.summary logs the result every run.route (io.kestra.plugin.core.flow.If):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.errors alerts Slack when the flow fails, for example on an unreachable broker.QuotaDescribe and QuotaAlter.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.
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 |
DescribeConfigs and, for remediation, AlterConfigs on the cluster.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.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.| 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.
SLACK_WEBHOOK_URL secret.bootstrap_servers, and replace policy with your clients and rates.auto_remediate: true and dry_run: false to apply them.hourly trigger, which always runs in safe mode.entityUser and entityIp properties.