Apache Kafka How-to Guides¶
Scope
Task recipes for Apache Kafka 4.x in KRaft mode: run a dev broker, plan a production cluster, tune producers, brokers, and consumers, upgrade, secure listeners, troubleshoot, control cost, and everyday CLI, tiered storage, Strimzi, and monitoring recipes. Defaults and version facts are in Reference. The internals behind these steps are in Explanation.
Related Notes
4.x path changes
Since 4.0 the sample configs live in config/ (the old config/kraft/ directory is gone). Commands below assume the extracted kafka_2.13-4.3.1 directory and Java 17+ on the broker host.
Run a Local Dev Broker¶
This follows the official 4.3 quickstart. It starts a single combined broker and controller.
# Option A: tarball (needs Java 17+)
tar -xzf kafka_2.13-4.3.1.tgz && cd kafka_2.13-4.3.1
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format --standalone -t "$KAFKA_CLUSTER_ID" -c config/server.properties
bin/kafka-server-start.sh config/server.properties
# Option B: official Docker images (JVM or GraalVM native)
docker run -p 9092:9092 apache/kafka:4.3.1
docker run -p 9092:9092 apache/kafka-native:4.3.1
# Smoke test
bin/kafka-topics.sh --create --topic quickstart-events --bootstrap-server localhost:9092
echo "hello" | bin/kafka-console-producer.sh --topic quickstart-events --bootstrap-server localhost:9092
bin/kafka-console-consumer.sh --topic quickstart-events --from-beginning --bootstrap-server localhost:9092
Production Deployment Patterns¶
Cluster Topology (KRaft)¶
| Pattern | Controllers | Brokers | When to use |
|---|---|---|---|
Combined (process.roles=broker,controller) |
1 to 3 | Same nodes | Dev, staging, or very small production (fewer than 5 brokers) |
Isolated (process.roles=controller or broker) |
3 dedicated | N dedicated | Production. Recommended for any cluster you would page someone for |
| Cross-AZ | 3 controllers across 3 AZs | Brokers across 3 AZs, broker.rack set, RF=3, min.insync.replicas=2 |
Single-region high availability |
| Cross-region (active-passive) | One cluster per region | MirrorMaker 2 replicating | Disaster recovery with RPO of seconds |
| Cross-region (active-active) | One cluster per region | Bidirectional MM2 with prefixed (renamed) topics | Lowest RPO, at the cost of dual writes and topic prefixing |
A 3-controller quorum tolerates one controller failure. Five controllers tolerate two. Most production clusters run three.
Broker Sizing¶
- CPU: 8 to 16 vCPU per broker is typical. zstd compression and TLS are the heaviest CPU consumers. TLS also disables zero-copy.
- RAM: the page cache does the caching. Plan for 32 GiB or more (often 64 to 128 GiB) so the active working set fits in cache. Keep the JVM heap small, usually 4 to 6 GiB, so it does not crowd out the cache.
- Disk: SSD or NVMe, with XFS preferred over ext4. Keep log directories off the OS disk. Several log directories on several disks (JBOD) spread partition load.
- Network: 10 GbE minimum, 25 GbE or more for large clusters. Replication multiplies inbound bytes by the replication factor.
- JVM: G1GC is the default. Some large deployments use ZGC for shorter pauses. Validate either on your workload.
Topic Design¶
| Decision | Guidance |
|---|---|
| Partitions per topic | Size for peak consumer parallelism plus headroom (for example 12, 24, or 48). You cannot decrease partitions later |
| Replication factor | 3 for production. RF=2 leaves no headroom for a failure during a rolling restart |
min.insync.replicas |
RF minus 1 (2 for RF=3). With acks=all, this gives durable writes |
| Cleanup policy | delete for events, compact for state, compact,delete for hybrids such as CDC |
| Retention | Set per topic with retention.ms and/or retention.bytes |
| Compression | compression.type=zstd on the producer usually gives the best ratio. lz4 and snappy use less CPU |
| Tiered storage | Worth considering once retention exceeds about a week at high ingest. Enable remote.storage.enable=true per topic (not for compacted topics) |
Partition Planning¶
- Total partitions scale with broker count and hardware. A modest 6-broker KRaft cluster handles tens of thousands of partitions. Validate your own ceiling with a load test.
- Per-broker partition count drives metadata size, follower fetch work, and restart recovery time. Watch
kafka.controller:type=KafkaController,name=PreferredReplicaImbalanceCountand restart times. - Adding partitions later is cheap, but it changes the key-to-partition mapping, so per-key ordering breaks across the change.
Performance Tuning¶
Producer¶
| Setting | Default | Tuning advice |
|---|---|---|
acks |
all (since 3.0) |
Keep all for durability. Use 1 only for low-stakes telemetry |
enable.idempotence |
true (since 3.0) |
Leave on. It prevents duplicates from retries at negligible cost |
compression.type |
none |
zstd for the best ratio, lz4/snappy for lower CPU |
batch.size |
16384 (16 KiB) | Raise to 64 KiB to 1 MiB for throughput-oriented workloads |
linger.ms |
5 (was 0 before 4.0) | Higher values make bigger batches: better throughput, higher latency |
max.in.flight.requests.per.connection |
5 | Must stay at 5 or less while idempotence is enabled, which preserves ordering on retries |
buffer.memory |
32 MiB | Raise if bufferpool-wait-time-ns-total grows |
Broker¶
| Setting | Default | Tuning advice |
|---|---|---|
num.io.threads |
8 | Match disk parallelism. Raise for hosts with many disks |
num.network.threads |
3 | Raise for high connection counts or TLS-heavy traffic |
num.replica.fetchers |
1 | Raise to 4 to 8 when each broker follows many partitions |
socket.send.buffer.bytes / socket.receive.buffer.bytes |
102400 | Match the bandwidth-delay product on fat, long-distance links |
log.flush.interval.messages / log.flush.interval.ms |
Unset (the OS flushes) | Leave unset unless you have a specific durability mandate. Rely on RF and acks=all |
log.segment.bytes |
1 GiB | Smaller segments give finer retention and tiering at the cost of more files |
Consumer¶
| Setting | Default | Tuning advice |
|---|---|---|
fetch.min.bytes |
1 | Raise to 1 to 10 KiB for batching efficiency |
fetch.max.wait.ms |
500 | Pair with fetch.min.bytes to amortize fetches |
max.poll.records |
500 | Lower for slow per-record processing to avoid max.poll.interval.ms timeouts |
isolation.level |
read_uncommitted |
Set read_committed for transactional pipelines |
group.protocol |
classic |
Set consumer to use the KIP-848 protocol (see below) |
Page Cache and Zero-Copy¶
- Kafka leaves caching to the Linux page cache. Do not size the JVM heap so large that it crowds out the cache. A 4 to 6 GiB heap on a 64 GiB host is typical.
- Tail reads, from consumers near the end of the log, are served almost entirely from page cache. Lagging consumers that read old data cause disk reads, and on tiered topics remote reads.
- On plaintext listeners the broker uses
sendfile(2)to stream batches without a user-space copy. On TLS listeners it cannot, so budget extra broker CPU for TLS. See Explanation.
Switch Consumers to the KIP-848 Protocol¶
The consumer protocol is GA since 4.0 and is enabled on the server once the cluster is finalized at 4.0 or later. Opt clients in one group at a time:
# consumer.properties
group.protocol=consumer
# Optional: server-side assignor ("uniform" or "range")
group.remote.assignor=uniform
Client-side settings tied to the classic protocol (partition.assignment.strategy, session.timeout.ms, heartbeat.interval.ms) are not used with group.protocol=consumer. The equivalents are group configs on the broker. A running classic group can be converted by rolling its members to the new protocol. Once any group uses the new protocol, the cluster can only be downgraded to 3.4.1 or later.
Upgrade a Cluster¶
Rolling Upgrade Within 4.x¶
- Read the upgrade notes for every version you cross.
- Upgrade the brokers and controllers one at a time: stop the node, install the new binaries, and restart it. Wait for
UnderReplicatedPartitionsto return to 0 before moving on. - After you have verified behavior, finalize the metadata and feature versions:
bin/kafka-features.sh --bootstrap-server localhost:9092 describe
bin/kafka-features.sh --bootstrap-server localhost:9092 upgrade --release-version 4.3
Downgrades
A metadata version with metadata changes (4.3's 4.3-IV0 is one) cannot be downgraded. Check the upgrade notes before finalizing. From 4.4, storage directories formatted by kafka-storage are not forward-compatible: the broker must be at least the version of the tool that formatted it (KIP-1170).
Upgrade From ZooKeeper or Old KRaft¶
- ZooKeeper clusters must migrate to KRaft on a 3.x bridge release (3.9.x recommended) before moving to 4.x. There is no ZooKeeper support or migration path in 4.0 and later.
- KRaft clusters older than 3.3 should go to 3.9.x first.
- Java clients older than 2.1 cannot connect to 4.x brokers (KIP-896). Upgrade clients first, including non-Java clients built on old protocol versions.
Move a Static Controller Quorum to a Dynamic One¶
Clusters formatted with controller.quorum.voters can switch to a dynamic quorum on 4.1 or later:
bin/kafka-features.sh --bootstrap-server localhost:9092 upgrade --feature kraft.version=1
# Then replace controller.quorum.voters with controller.quorum.bootstrap.servers on every node and roll them.
bin/kafka-metadata-quorum.sh --bootstrap-controller localhost:9093 describe --status
Secure a Cluster¶
The threat model and mechanism trade-offs are in Explanation: Security Model. The checklist is in Reference: Hardening Checklist.
Java properties files have no inline comments
In server.properties, value # comment makes the comment part of the value. Put comments on their own lines, as below.
Configure Listeners¶
# server.properties (broker): TLS for clients, mTLS between brokers
listeners=SASL_SSL://:9094,SSL://:9095
advertised.listeners=SASL_SSL://broker1.example.com:9094,SSL://broker1.example.com:9095
controller.listener.names=CONTROLLER
listener.security.protocol.map=SASL_SSL:SASL_SSL,SSL:SSL,CONTROLLER:SSL
# Broker-to-broker traffic uses the mTLS listener
inter.broker.listener.name=SSL
sasl.enabled.mechanisms=SCRAM-SHA-512,OAUTHBEARER
Set Up SCRAM¶
# Create a SCRAM-SHA-512 credential for User:alice (stored in KRaft metadata)
bin/kafka-configs.sh --bootstrap-server kafka:9094 \
--command-config admin.properties \
--alter --add-config 'SCRAM-SHA-512=[iterations=8192,password=ChangeMeNow]' \
--entity-type users --entity-name alice
# Client config
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
username="alice" password="ChangeMeNow";
To create SCRAM credentials for the inter-broker principal before the first start, pass --add-scram to kafka-storage.sh format.
Set Up OAUTHBEARER With OIDC¶
# Client: client_credentials grant against the identity provider
security.protocol=SASL_SSL
sasl.mechanism=OAUTHBEARER
sasl.oauthbearer.token.endpoint.url=https://idp.example.com/realms/prod/protocol/openid-connect/token
sasl.login.callback.handler.class=org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler
sasl.jaas.config=org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required \
clientId="orders-svc" clientSecret="${env:OIDC_CLIENT_SECRET}";
Current docs also document dedicated client properties (sasl.oauthbearer.client.credentials.client.id, ...client.secret, sasl.oauthbearer.scope) and client-assertion auth (KIP-1258, 4.3). Check the SASL docs for your version.
# Broker: validate JWTs against the IdP's JWKS
listener.name.sasl_ssl.oauthbearer.sasl.server.callback.handler.class=\
org.apache.kafka.common.security.oauthbearer.OAuthBearerValidatorCallbackHandler
listener.name.sasl_ssl.oauthbearer.sasl.oauthbearer.jwks.endpoint.url=\
https://idp.example.com/realms/prod/protocol/openid-connect/certs
listener.name.sasl_ssl.oauthbearer.sasl.oauthbearer.expected.audience=kafka
listener.name.sasl_ssl.oauthbearer.sasl.oauthbearer.expected.issuer=https://idp.example.com/realms/prod
# Allow-list the IdP URLs. This is a JVM system property, not a server.properties key.
# Required since 4.0, where the default is an empty list (CVE-2025-27817 hardening).
export KAFKA_OPTS="-Dorg.apache.kafka.sasl.oauthbearer.allowed.urls=https://idp.example.com/realms/prod/protocol/openid-connect/certs,https://idp.example.com/realms/prod/protocol/openid-connect/token"
bin/kafka-server-start.sh config/server.properties
OAUTHBEARER version gotchas
- 4.1.0 and 4.1.1: the default broker JWT validator accepts unsigned tokens (CVE-2026-33557). Upgrade, or set
sasl.oauthbearer.jwt.validator.class=org.apache.kafka.common.security.oauthbearer.BrokerJwtValidator. - 4.4 (upcoming): a broker with a JWKS URL but no
expected.audienceorexpected.issuerfails at startup unless you setsasl.oauthbearer.allow.unverified.audienceorallow.unverified.issuer. - On 3.9.1 the allowed-URLs property defaults to "allow all" for backward compatibility. Set it explicitly.
Set Up Mutual TLS¶
# server.properties: TLS listener that requires client certificates
listeners=SSL://:9095
ssl.keystore.location=/etc/kafka/tls/server.keystore.jks
ssl.keystore.password=${env:KAFKA_KEYSTORE_PASSWORD}
ssl.key.password=${env:KAFKA_KEY_PASSWORD}
ssl.truststore.location=/etc/kafka/tls/server.truststore.jks
ssl.truststore.password=${env:KAFKA_TRUSTSTORE_PASSWORD}
ssl.client.auth=required
ssl.enabled.protocols=TLSv1.3,TLSv1.2
ssl.endpoint.identification.algorithm=https
# Optional: map the certificate DN to a short principal name
ssl.principal.mapping.rules=RULE:^CN=([^,]+).*$/$1/,DEFAULT
The ${env:...} syntax needs config.providers=env and config.providers.env.class=org.apache.kafka.common.config.provider.EnvVarConfigProvider. Rotate certificates with cert-manager or, on Kubernetes, the Strimzi cluster CA.
Enable ACL Authorization¶
# server.properties
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
super.users=User:admin;User:CN=cluster-admin
# Default deny for resources without ACLs (this is the default)
allow.everyone.if.no.acl.found=false
ACL recipes are under ACLs (kafka-acls.sh).
Configure Authorizer Audit Logging¶
Kafka 4.x ships config/log4j2.yaml with an AuthorizerAppender that writes kafka-authorizer.log, and the kafka.authorizer.logger logger at INFO (denials only). To also log allowed requests, raise that logger to DEBUG. This is high volume:
# config/log4j2.yaml (excerpt)
Loggers:
Logger:
- name: kafka.authorizer.logger
level: DEBUG
additivity: false
AppenderRef:
ref: AuthorizerAppender
Ship kafka-authorizer.log to your SIEM and alert on bursts of Denied Operation lines.
Secure MirrorMaker 2¶
# mm2.properties: SASL_SSL on both clusters with separate principals
clusters = src, dst
src.bootstrap.servers = src-broker1:9094,src-broker2:9094,src-broker3:9094
dst.bootstrap.servers = dst-broker1:9094,dst-broker2:9094,dst-broker3:9094
src.security.protocol = SASL_SSL
src.sasl.mechanism = SCRAM-SHA-512
src.sasl.jaas.config = org.apache.kafka.common.security.scram.ScramLoginModule required \
username="mm2-src-reader" password="${env:MM2_SRC_PASSWORD}";
src.ssl.truststore.location = /etc/mm2/src-truststore.jks
dst.security.protocol = SASL_SSL
dst.sasl.mechanism = SCRAM-SHA-512
dst.sasl.jaas.config = org.apache.kafka.common.security.scram.ScramLoginModule required \
username="mm2-dst-writer" password="${env:MM2_DST_PASSWORD}";
dst.ssl.truststore.location = /etc/mm2/dst-truststore.jks
src->dst.enabled = true
src->dst.topics = orders.*, events.*
replication.factor = 3
sync.topic.acls.enabled = true
sync.group.offsets.enabled = true
Checklist for MM2:
- Use a read-only principal on the source and a write principal on the target that can also create MM2 internal topics.
- Restrict MM2 internal topics (
mm2-offsets.*.internal,mm2-configs.*.internal,mm2-status.*.internal,heartbeats,*.checkpoints.internal) with ACLs. - Keep
sync.topic.acls.enabled=trueso a failover does not silently break authorization. - Treat MM2 worker hosts as part of the Kafka security boundary.
Troubleshooting¶
Under-Replicated Partitions¶
Symptom: kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions stays above 0.
Causes and responses:
- Slow follower fetch. Check follower fetch rates and
kafka.server:type=FetcherLagMetrics. Raisenum.replica.fetchers, or look at disk and network on the slow follower. - Broker down. Check
kafka.controller:type=KafkaController,name=ActiveBrokerCountandkafka-metadata-quorum.sh describe --replication. Restart the broker and let it catch up. - GC pauses. A long GC on the leader stalls produce, and on a follower it stalls fetch. Inspect the GC log and heap sizing.
- Disk saturation. Run
iostat -xm 1and look at%utilandawait. Move log directories to faster disks or spread them across more disks.
ISR Shrinkage¶
Symptom: kafka.server:type=ReplicaManager,name=IsrShrinksPerSec is non-zero, or partitions show up under min ISR.
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --under-min-isr-partitions
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders
# Describe output shows Leader, Replicas, Isr and, with ELR enabled, Elr / LastKnownElr
While the ISR is below min.insync.replicas, acks=all produce requests fail with NotEnoughReplicas. ELR (KIP-966) does not change that. What ELR adds is a safe leader candidate if the remaining ISR members also fail, so you do not need an unclean election. Fix the lagging or failed replicas, which are usually disk, network, or GC problems as above.
Broker Bounce Procedure¶
To restart a broker without disrupting clients:
# 1. Stop the broker. SIGTERM triggers a controlled shutdown: the broker asks the
# controller to move its partition leaderships elsewhere before it exits.
bin/kafka-server-stop.sh
# 2. Do the maintenance (kernel patch, JVM upgrade, config change), then restart
bin/kafka-server-start.sh -daemon config/server.properties
# 3. Wait for UnderReplicatedPartitions to return to 0, then restore preferred leaders
# (automatic when auto.leader.rebalance.enable=true, the default)
bin/kafka-leader-election.sh --bootstrap-server localhost:9092 \
--election-type PREFERRED --all-topic-partitions
To decommission a broker or a disk on 4.3+, first cordon it (cordoned.log.dirs, KIP-1066) so no new partitions land there. Then move its partitions off with kafka-reassign-partitions.sh or Cruise Control.
Consumer Lag¶
bin/kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group orders-aggregator
# Check the LAG column. Sustained growth means the consumer is too slow or poorly partitioned.
Mitigations: scale consumers up to the partition count, or use a share group if you need more workers than partitions and do not need ordering. Raise max.poll.records if processing is fast, and lower it if you hit max.poll.interval.ms. Check downstream back-pressure, such as database write latency.
KRaft Controller Issues¶
# Quorum leader, epoch, high watermark, voters and observers
bin/kafka-metadata-quorum.sh --bootstrap-server kafka:9092 describe --status
# Per-replica lag on the metadata log
bin/kafka-metadata-quorum.sh --bootstrap-server kafka:9092 describe --replication
# Talk to controllers directly (for example when brokers are down)
bin/kafka-metadata-quorum.sh --bootstrap-controller controller1:9093 describe --status
If the active controller flaps, check controller-host CPU, network partitions, and disk latency. Controllers fsync the metadata log on commit.
Cost Analysis¶
| Cost driver | Typical contributors | Reduction levers |
|---|---|---|
| Storage | Replication factor, retention, message size | Tiered storage (KIP-405) to object storage, compaction, shorter retention.ms, zstd compression |
| Network egress | Cross-AZ replication, cross-region MM2, consumer fetches | Rack awareness (broker.rack), follower fetching for consumers (KIP-392, client.rack with replica.selector.class), fewer cross-region copies |
| Compute | Broker and controller count, TLS, compression | Right-size brokers, run controllers on small hosts, use CPUs with AES-NI for TLS |
| Managed surcharges | MSK and Confluent Cloud pricing | Self-host on Kubernetes with Strimzi, or use serverless tiers (MSK Serverless, Confluent Cloud Basic) for spiky workloads |
On AWS, cross-AZ data transfer is often the largest variable cost for self-managed Kafka. Rack-aware consumers that fetch from a same-AZ follower (KIP-392) can cut it substantially.
Commands & Recipes¶
KRaft Cluster Bootstrap¶
# One cluster ID for the whole cluster
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
# Controllers: format with the initial voter set (dynamic quorum, KIP-853)
CONTROLLER_0_UUID="$(bin/kafka-storage.sh random-uuid)"
CONTROLLER_1_UUID="$(bin/kafka-storage.sh random-uuid)"
CONTROLLER_2_UUID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format --cluster-id "$KAFKA_CLUSTER_ID" \
--initial-controllers "0@controller-0:9093:${CONTROLLER_0_UUID},1@controller-1:9093:${CONTROLLER_1_UUID},2@controller-2:9093:${CONTROLLER_2_UUID}" \
--config config/controller.properties
# Brokers: format with no extra flags
bin/kafka-storage.sh format --cluster-id "$KAFKA_CLUSTER_ID" --config config/broker.properties
# Start each node
bin/kafka-server-start.sh -daemon config/broker.properties
# Verify
bin/kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe --status
bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092
kafka-broker-api-versions.sh is deprecated in 4.4 in favor of kafka-cluster.sh api-versions (KIP-1220). For a single-node dev cluster, use format --standalone instead of --initial-controllers.
A minimal broker config for an isolated-mode cluster:
# config/broker.properties
process.roles=broker
node.id=10
controller.quorum.bootstrap.servers=controller-0:9093,controller-1:9093,controller-2:9093
controller.listener.names=CONTROLLER
listeners=PLAINTEXT://:9092
advertised.listeners=PLAINTEXT://broker-10.example.com:9092
inter.broker.listener.name=PLAINTEXT
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
log.dirs=/var/lib/kafka/data
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
auto.create.topics.enable=false
Add controllers to a running dynamic quorum with kafka-storage.sh format --no-initial-controllers, start the node, then run kafka-metadata-quorum.sh --bootstrap-controller <host:9093> add-controller.
Topic Lifecycle (kafka-topics.sh)¶
# Create
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic orders \
--partitions 12 --replication-factor 3 \
--config retention.ms=604800000 \
--config min.insync.replicas=2 \
--config compression.type=zstd \
--config cleanup.policy=delete
# List
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
# Describe (leader, replicas, ISR, and ELR when enabled)
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders
# Increase partitions (changes key-to-partition mapping)
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--alter --topic orders --partitions 24
# Find problem partitions
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --under-replicated-partitions
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --under-min-isr-partitions
# Delete
bin/kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic orders
Configuration (kafka-configs.sh)¶
# Inspect topic-level config
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name orders --describe
# Tighten retention for an existing topic
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name orders \
--alter --add-config retention.ms=86400000
# Enable tiered storage on an existing topic (broker must have remote.log.storage.system.enable=true)
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name orders \
--alter --add-config remote.storage.enable=true,local.retention.ms=3600000
# Client quotas per user
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --entity-type users --entity-name billing-svc \
--add-config 'producer_byte_rate=10485760,consumer_byte_rate=20971520'
# Cluster-wide min.insync.replicas (required form when ELR is enabled)
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type brokers --entity-default \
--alter --add-config min.insync.replicas=2
# Dynamic per-broker config (no restart)
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type brokers --entity-name 1 \
--alter --add-config log.cleaner.threads=4
Consumer Groups (kafka-consumer-groups.sh)¶
# List
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
# Describe (lag per partition)
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group orders-aggregator
# Members and assignments
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group orders-aggregator --members --verbose
# Offset reset (dry run first; the group must have no active members)
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group orders-aggregator --topic orders \
--reset-offsets --to-earliest --dry-run
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group orders-aggregator --topic orders \
--reset-offsets --to-earliest --execute
# Reset to a point in time
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group orders-aggregator --topic orders \
--reset-offsets --to-datetime 2026-04-28T00:00:00.000 --execute
# Delete a group
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--delete --group orders-aggregator
Share Groups (Queues)¶
Share groups are production-ready from 4.2. share.version=1 is part of the 4.2+ release version, so new clusters and clusters finalized with --release-version 4.2 or later have it. On 4.1 (preview) you had to enable it explicitly.
# Check or enable the feature
bin/kafka-features.sh --bootstrap-server localhost:9092 describe
bin/kafka-features.sh --bootstrap-server localhost:9092 upgrade --feature share.version=1
# Consume a topic as a share group from the console
bin/kafka-console-share-consumer.sh --bootstrap-server localhost:9092 \
--topic jobs --group job-workers
# Inspect share groups
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --list
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --describe --group job-workers
# Tune the lock duration and delivery limit for one group
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type groups --entity-name job-workers \
--alter --add-config share.record.lock.duration.ms=60000,share.delivery.count.limit=3
On clusters with fewer than 3 brokers, set share.coordinator.state.topic.replication.factor=1 and share.coordinator.state.topic.min.isr=1 before first use.
ACLs (kafka-acls.sh)¶
# Producer permissions on a topic
bin/kafka-acls.sh --bootstrap-server kafka:9094 --command-config admin.properties \
--add --allow-principal User:orders-svc \
--producer --topic orders
# Consumer permissions for a specific group
bin/kafka-acls.sh --bootstrap-server kafka:9094 --command-config admin.properties \
--add --allow-principal User:billing-svc \
--consumer --topic orders --group billing-aggregator
# Prefixed topic ACL (every topic starting with "events.")
bin/kafka-acls.sh --bootstrap-server kafka:9094 --command-config admin.properties \
--add --allow-principal User:platform-events \
--producer --topic events. --resource-pattern-type prefixed
# Transactional producer (transactional.id ACL)
bin/kafka-acls.sh --bootstrap-server kafka:9094 --command-config admin.properties \
--add --allow-principal User:orders-svc \
--operation Write --operation Describe \
--transactional-id orders-svc-txn-1
# Cluster admin (use sparingly)
bin/kafka-acls.sh --bootstrap-server kafka:9094 --command-config admin.properties \
--add --allow-principal User:CN=cluster-admin \
--operation All --cluster
# List existing ACLs
bin/kafka-acls.sh --bootstrap-server kafka:9094 --command-config admin.properties --list
Group describe permission
ConsumerGroupDescribe (the KIP-848 describe API) requires DESCRIBE on the group, not READ. Earlier docs said READ (CVE-2026-41115). Review group ACLs with that in mind.
kcat (formerly kafkacat)¶
# Produce a single message (key=k1, value=v1)
echo "k1:v1" | kcat -b localhost:9092 -t orders -K: -P
# Consume from the beginning, print key, offset, and value, exit at the end
kcat -b localhost:9092 -t orders -C -o beginning -e \
-f '%k\t%o\t%s\n'
# Consume as a member of a balanced consumer group
kcat -b localhost:9092 -G orders-aggregator orders \
-f '%k\t%o\t%s\n'
# Produce newline-delimited JSON from a file
kcat -b localhost:9092 -t events -P -l events.ndjson
# Cluster metadata
kcat -b localhost:9092 -L
Tiered Storage Setup¶
Apache Kafka ships no production RemoteStorageManager. The snippet below uses the test-only LocalTieredStorage from the official quickstart. Replace the class and path with your plugin's (for example Aiven's S3/GCS/Azure plugin).
# server.properties (broker)
remote.log.storage.system.enable=true
remote.log.metadata.manager.listener.name=PLAINTEXT
remote.log.storage.manager.class.name=org.apache.kafka.server.log.remote.storage.LocalTieredStorage
remote.log.storage.manager.class.path=/opt/kafka/plugins/kafka-storage-4.3.1-test.jar
rsm.config.dir=/var/lib/kafka/remote-tier
rlmm.config.remote.log.metadata.topic.replication.factor=3
rlmm.config.remote.log.metadata.topic.min.isr=2
# Per-topic enablement
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic events-archive \
--partitions 12 --replication-factor 3 \
--config remote.storage.enable=true \
--config local.retention.ms=3600000 \
--config retention.ms=2592000000 \
--config segment.bytes=536870912
# Stop uploading but keep remote data readable
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --entity-type topics --entity-name events-archive \
--add-config 'remote.storage.enable=true,remote.log.copy.disable=true,local.retention.ms=-2,local.retention.bytes=-2'
# Disable tiering and delete remote data for a topic
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --entity-type topics --entity-name events-archive \
--add-config 'remote.storage.enable=false,remote.log.delete.on.disable=true'
Strimzi Operator (Kubernetes)¶
Strimzi (CNCF Incubating) is KRaft-only since 0.46 and moved its CRDs to kafka.strimzi.io/v1 in 1.0. Brokers and controllers are defined as KafkaNodePool resources. This follows the upstream kafka-persistent.yaml example.
# kafka-cluster.yaml (Strimzi 1.2+, Kafka 4.3.1)
apiVersion: kafka.strimzi.io/v1
kind: KafkaNodePool
metadata:
name: controller
labels:
strimzi.io/cluster: orders-cluster
spec:
replicas: 3
roles:
- controller
storage:
type: jbod
volumes:
- id: 0
type: persistent-claim
size: 100Gi
kraftMetadata: shared
---
apiVersion: kafka.strimzi.io/v1
kind: KafkaNodePool
metadata:
name: broker
labels:
strimzi.io/cluster: orders-cluster
spec:
replicas: 6
roles:
- broker
storage:
type: jbod
volumes:
- id: 0
type: persistent-claim
size: 500Gi
kraftMetadata: shared
---
apiVersion: kafka.strimzi.io/v1
kind: Kafka
metadata:
name: orders-cluster
spec:
kafka:
version: 4.3.1
metadataVersion: 4.3-IV0
listeners:
- name: tls
port: 9093
type: internal
tls: true
authentication:
type: tls
config:
offsets.topic.replication.factor: 3
transaction.state.log.replication.factor: 3
transaction.state.log.min.isr: 2
default.replication.factor: 3
min.insync.replicas: 2
metricsConfig:
type: jmxPrometheusExporter
valueFrom:
configMapKeyRef:
name: kafka-metrics
key: kafka-metrics-config.yml
entityOperator:
topicOperator: {}
userOperator: {}
# Install the Strimzi operator (OCI Helm chart)
helm install strimzi-kafka-operator oci://quay.io/strimzi-helm/strimzi-kafka-operator \
--namespace kafka --create-namespace
# Apply the node pools and Kafka CR
kubectl apply -f kafka-cluster.yaml -n kafka
# Watch reconciliation
kubectl get kafka -n kafka -w
kubectl get kafkanodepool -n kafka
Prometheus JMX Exporter¶
# kafka-metrics-config.yml: a minimal set of high-value rules
rules:
- pattern: "kafka.server<type=(.+), name=(.+)PerSec\\w*, topic=(.+)><>Count"
name: kafka_server_$1_$2_total
labels:
topic: "$3"
type: COUNTER
- pattern: "kafka.server<type=ReplicaManager, name=UnderReplicatedPartitions><>Value"
name: kafka_server_replicamanager_underreplicated_partitions
type: GAUGE
- pattern: "kafka.server<type=ReplicaManager, name=IsrShrinksPerSec><>Count"
name: kafka_server_replicamanager_isr_shrinks_total
type: COUNTER
- pattern: "kafka.controller<type=KafkaController, name=ActiveControllerCount><>Value"
name: kafka_controller_active_controllers
type: GAUGE
- pattern: "kafka.network<type=RequestMetrics, name=RequestsPerSec, request=(.+)><>Count"
name: kafka_network_requests_total
labels:
request: "$1"
type: COUNTER
Run the exporter as a Java agent on each node. Use the current release of io.prometheus.jmx:jmx_prometheus_javaagent:
export KAFKA_OPTS="-javaagent:/opt/jmx_prometheus_javaagent.jar=7071:/opt/kafka-metrics-config.yml"
bin/kafka-server-start.sh config/server.properties
Strimzi wires this up automatically when metricsConfig.type: jmxPrometheusExporter is set.
Useful PromQL¶
# Total bytes in and out per second across the cluster
sum(rate(kafka_server_brokertopicmetrics_bytesin_total[5m]))
sum(rate(kafka_server_brokertopicmetrics_bytesout_total[5m]))
# Under-replicated partitions per broker
kafka_server_replicamanager_underreplicated_partitions
# ISR shrink rate (alert if > 0 for a sustained period)
sum(rate(kafka_server_replicamanager_isr_shrinks_total[5m])) by (instance)
# Active controller count (must be exactly 1)
sum(kafka_controller_active_controllers)
# Consumer lag (requires kafka_exporter)
max(kafka_consumergroup_lag) by (consumergroup, topic)
Sources¶
- Apache Kafka quickstart
- Apache Kafka: Operations documentation
- Apache Kafka upgrade notes (4.x)
- Apache Kafka: KRaft operations
- Apache Kafka: Tiered storage
- Apache Kafka: Authentication using SASL
- Apache Kafka: Authorization and ACLs
- Apache Kafka CVE list
- Strimzi documentation
- Strimzi CHANGELOG
- Prometheus JMX Exporter
- kcat