Apache Kafka & KRaft SRE Cheat Sheet Hub
Battle-tested CLI commands, offset recovery recipes, Flink SQL streaming patterns, and KRaft consensus diagnostics.
KafkaKraft Enterprise SRE Cheat Sheet
labs.kafkakraft.com — Production KRaft & Streaming Reference Guide
Updated: 2026
Inspect KRaft Quorum & Controller Status
View active metadata leader ID, leader epoch, and replication high watermarks
kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe --status
💡 Requires Kafka 3.0+ in KRaft mode. Returns LeaderId, LeaderEpoch, HighWatermark, and CurrentVoters.
Inspect Quorum Voter Replication Offsets
Check if metadata voter nodes are in-sync or lagging behind the controller leader
kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe --replication
💡 Crucial for diagnosing partition split-brain and controller election delays.
Generate Cluster ID & Format KRaft Storage Directory
Format the storage directory before bootstrapping a new KRaft broker/controller
KAFKA_CLUSTER_ID="$(kafka-storage.sh random-uuid)" && kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /etc/kafka/server.properties
💡 All brokers in the same cluster MUST share the exact same cluster ID.
Dump and Decode Metadata Log Segment (@metadata)
Human-readable decode of the __cluster_metadata partition record batches
kafka-dump-log.sh --cluster-metadata-decoder --files /var/lib/kafka/data/@metadata-0/00000000000000000000.log --print-data-log
💡 Shows PartitionChangeRecord, RegisterBrokerRecord, and FeatureLevelRecord.
Production Throughput & Latency Benchmark
Emit 1,000,000 synthetic records at 100k msg/sec to benchmark disk I/O
kafka-producer-perf-test.sh --topic benchmark-stream --num-records 1000000 --record-size 1024 --throughput 100000 --producer-props bootstrap.servers=localhost:9092 acks=all compression.type=zstd linger.ms=20 batch.size=131072
💡 Outputs 50th, 95th, 99th, and 99.9th percentile produce latency and MB/sec throughput.
Interactive Console Producer with Key & Value
Send messages with explicit key:value delimiters to control partition hashing
kafka-console-producer.sh --bootstrap-server localhost:9092 --topic orders --property "parse.key=true" --property "key.separator=:"
💡 Format input as 'user_101:{"order_id": 994, "amount": 42.50}'.
Strict Exactly-Once Producer CLI Properties
Flags required to guarantee no duplicates and in-order partition delivery
enable.idempotence=true acks=all retries=2147483647 max.in.flight.requests.per.connection=5
💡 Default in Kafka 3.0+. Prevents duplicate records caused by network TCP retries.
List All Registered Consumer Groups
Inspect active, empty, and dead consumer groups across all protocols
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
💡 Returns group IDs including CooperativeSticky, Range, and RoundRobin assignors.
Inspect Consumer Group Lag & Partition Offsets
View current-offset, log-end-offset, consumer lag, and client host IDs
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group payment-service --describe
💡 A growing LAG column indicates consumer starvation, GC stalls, or poison pills.
Rewind Consumer Group Offsets to Earliest (Replay All)
Reset offsets to offset 0 across all partitions for event stream reprocessing
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group payment-service --reset-offsets --to-earliest --topic orders --execute
💡 Group must have 0 active members (STOP the consumer processes before executing).
Rewind Offsets to Specific Point-in-Time (ISO-8601)
Seek offsets to exactly 2 hours ago or a specific timestamp
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group payment-service --reset-offsets --to-datetime 2026-10-02T14:00:00.000 --topic orders --execute
💡 Ideal for recovering from bad deployment releases by rewinding to pre-bug time.
Inspect Group Coordinator & Rebalance State
Check if consumer group is in Stable, PreparingRebalance, or CompletingRebalance
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group payment-service --describe --state
💡 Frequent PreparingRebalance indicates heartbeats timing out or max.poll.interval.ms exceeded.
Create Topic with Explicit RF and Retention
Create topic with 12 partitions, RF 3, and 7-day segment retention
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic telemetry-events --partitions 12 --replication-factor 3 --config min.insync.replicas=2 --config retention.ms=604800000
💡 Always set min.insync.replicas=2 when replication-factor=3 for durability.
Describe Topic ISR & Leader Distribution
View leader broker ID, replicas list, and in-sync replicas (ISR) per partition
kafka-topics.sh --bootstrap-server localhost:9092 --topic telemetry-events --describe
💡 Partitions where ISR length < Replicas length indicate offline or lagging brokers.
Dynamic Config Alteration (No Restart Required)
Increase maximum message payload size to 10MB dynamically
kafka-configs.sh --bootstrap-server localhost:9092 --entity-type topics --entity-name telemetry-events --alter --add-config max.message.bytes=10485760
💡 Broker does NOT need to be restarted. Applies within 5 seconds cluster-wide.
Generate & Execute Partition Reassignment Plan
Migrate partitions across newly added brokers to rebalance disk utilization
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 --reassignment-json-file /tmp/reassign.json --execute
💡 Use --verify to monitor percentage progress of replica movement across brokers.
Flink SQL: 5-Minute Tumbling Window Aggregation
Group continuous order stream into non-overlapping 5-minute sales buckets
SELECT window_start, window_end, merchant_id, COUNT(*) as tx_count, SUM(amount) as total_revenue FROM TABLE( TUMBLE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '5' MINUTE) ) GROUP BY window_start, window_end, merchant_id;
💡 Requires watermark definition on order_time column (e.g. WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND).
Flink SQL: Register Real-Time Kafka Source Table
DDL mapping a JSON/Avro Kafka topic into a streaming relational table
CREATE TABLE orders ( order_id STRING, customer_id STRING, amount DECIMAL(10, 2), order_time TIMESTAMP(3), WATERMARK FOR order_time AS order_time - INTERVAL '3' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'orders-v1', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'flink-sales-aggregator', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' );
💡 Provides streaming semantics with exactly-once checkpoints.
Create Granular Read/Write Topic ACLs
Grant User:alice write access to topic orders from any client host
kafka-acls.sh --bootstrap-server localhost:9092 --add --allow-principal User:alice --operation Write --operation Describe --topic orders --resource-pattern-type literal
💡 Must have authorizer.class.name enabled in server.properties.
Provision SASL/SCRAM-SHA-512 Credentials
Create salt and cryptographic hash in KRaft metadata for client authentication
kafka-configs.sh --bootstrap-server localhost:9092 --entity-type users --entity-name alice --alter --add-config 'SCRAM-SHA-512=[password=SecretKafkaPassword2026!]'
💡 SCRAM avoids plaintext password storage and supports rolling credential rotation.