New to Kestra?
Use blueprints to kickstart your first workflows.
Audit Kafka ACLs for wildcard principals, wildcard resources, cluster admin rights and unknown principals, alert Slack, delete with AclDelete and verify.
Kafka ACLs grow by accretion. A User:* ALLOW added to unblock a migration and never removed. A legacy ETL account still holding ALL on the cluster, which lets it rewrite every ACL. A test principal reading a production topic. Each of these turns one leaked credential, or no credential at all, into access to every topic.
This blueprint lists every ACL with the Kafka AdminClient and flags the grants that are broader than any service needs. It alerts Slack, and can delete the high-severity ACLs one by one with AclDelete. It then lists the ACLs again to prove the bad ones are gone and nothing else was touched.
This blueprint was created by zkasuran.
| Finding | Severity | Rule | Deleted by remediation |
|---|---|---|---|
WILDCARD_PRINCIPAL |
High | ALLOW for User:* |
Yes |
WILDCARD_RESOURCE |
High | ALLOW on every topic, group or transactional ID (*, literal), outside the admin principals |
Yes |
CLUSTER_ADMIN_RIGHTS |
High | ALL or ALTER on the cluster, outside the admin principals |
Yes |
OPERATION_ALL |
Medium | ALL on a topic or group, which includes DELETE and ALTER |
No, report only |
UNKNOWN_PRINCIPAL |
Medium | ACL for a principal that is not in known_principals |
No, report only |
Each ACL is reported once, with every rule it breaks and the most severe one first. DENY ACLs are never flagged.
list_acls (io.kestra.plugin.kafka.AclList, read-only) returns every ACL binding.findings applies the rules in one jq expression. summary logs every finding.route (io.kestra.plugin.core.flow.If):alert posts to Slack.remediate_gate runs only when auto_remediate is true, dry_run is false, and there is a high-severity finding.remediate (io.kestra.plugin.core.flow.Loop over io.kestra.plugin.kafka.AclDelete) deletes each high-severity ACL, matching all seven fields so no other binding matches the filter.relist and verify_result check that every deleted ACL is gone and every other ACL is still there.errors alerts Slack when the audit cannot run.AclList and AclDelete with gates around them.kafka-acls.sh --list prints the ACLs, but nobody reads it every day. Kestra runs the review on a schedule, keeps the list of known and admin principals next to the flow, turns findings into a Slack alert, deletes only behind auto_remediate and dry_run, and keeps a history of every run as evidence that broad grants were found and removed.
On Kestra 2.0.5 OSS against Apache Kafka 3.9.1 (KRaft, StandardAuthorizer), with a local HTTP sink standing in for Slack. The 9 seeded ACLs include:
orders-svc, billing and analytics;User:* with ALL on every topic, and legacy-etl with ALL on the cluster;billing with ALL on a group, and an unknown intern-test reading orders.| Run | Inputs | Result |
|---|---|---|
| 1 | defaults | 4 findings: 2 high (User:* on every topic, legacy-etl on the cluster), 2 medium. Legitimate ACLs not flagged. Alert sent |
| 2 | auto_remediate: true, dry_run: true |
Same report. kafka-acls.sh --list still shows 9 ACLs |
| 3 | auto_remediate: true, dry_run: false |
2 AclDelete calls, verified. kafka-acls.sh shows 7 ACLs, no User:*, no legacy-etl |
| 4 | same again | Only the 2 medium findings left. Nothing deleted |
| 5 | kafka-admin granted ALL on the cluster and on every topic |
Not flagged: it is an admin principal |
| 6 | unreachable broker | errors alert sent |
A Kafka cluster with an authorizer enabled, reachable from the Kestra worker, and a principal allowed to describe and, for remediation, alter ACLs (cluster ALTER).
Local testing: run apache/kafka:3.9.1 with:
KAFKA_AUTHORIZER_CLASS_NAME=org.apache.kafka.metadata.authorizer.StandardAuthorizer;KAFKA_SUPER_USERS=User:ANONYMOUS.Then add ACLs with kafka-acls.sh --add.
SLACK_WEBHOOK_URL: 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 |
known_principals |
3 example services | Expected service accounts |
admin_principals |
[User:kafka-admin] |
Allowed broad grants |
auto_remediate |
false |
Allow AclDelete |
dry_run |
true |
Must be false for a delete to run |
SLACK_WEBHOOK_URL secret.bootstrap_servers, known_principals and admin_principals.User:* grant, check which clients rely on it. Then run with auto_remediate: true and dry_run: false.daily trigger, which always runs in safe mode.