Apache Kafka
Essential commands and patterns for building distributed event streaming applications with Apache Kafka.
Apache Kafka
Essential commands and patterns for building distributed event streaming applications with Apache Kafka.
Overview
Apache Kafka is a distributed event streaming platform designed for high-throughput, fault-tolerant publish-subscribe messaging. It organises data into topics, which are split into partitions for parallel processing and distributed across brokers for scalability and redundancy.
graph TB
subgraph Kafka Cluster
subgraph Broker1["Broker 1"]
P1_0[Topic A<br/>Partition 0<br/>Leader]
P2_1[Topic B<br/>Partition 1<br/>Follower]
end
subgraph Broker2["Broker 2"]
P1_1[Topic A<br/>Partition 1<br/>Leader]
P2_0[Topic B<br/>Partition 0<br/>Follower]
end
subgraph Broker3["Broker 3"]
P1_0_R[Topic A<br/>Partition 0<br/>Follower]
P2_0_L[Topic B<br/>Partition 0<br/>Leader]
end
end
subgraph Clients
PROD[Producers]
CONS[Consumer Group]
end
subgraph Coordination
ZK[KRaft Controller Quorum]
end
PROD --> P1_0
PROD --> P1_1
CONS --> P1_0
CONS --> P1_1
ZK --> Broker1
ZK --> Broker2
ZK --> Broker3
Topics and Partitions
Topics are logical channels for organising messages, while partitions enable parallel processing and horizontal scaling.
Key Concepts
flowchart LR
subgraph Topic: orders
P0[Partition 0]
P1[Partition 1]
P2[Partition 2]
end
subgraph Partition 0 Detail
direction LR
O0[Offset 0] --> O1[Offset 1] --> O2[Offset 2] --> O3[Offset 3]
end
P0 --> O0
- Topics: Named feeds of messages; producers write to topics, consumers read from them
- Partitions: Ordered, immutable sequences of messages; enable parallelism
- Offsets: Unique sequential IDs for each message within a partition
- Replication Factor: Number of copies of each partition across brokers
- Leader/Follower: Each partition has one leader (handles reads/writes) and multiple followers (replicas)
Managing Topics with kafka-topics
# List all topics
kafka-topics.sh --bootstrap-server localhost:9092 --list
# Create a topic with specific configuration
kafka-topics.sh --bootstrap-server localhost:9092 \
--create \
--topic orders \
--partitions 6 \
--replication-factor 3
# Describe topic details
kafka-topics.sh --bootstrap-server localhost:9092 \
--describe \
--topic orders
# Describe all topics
kafka-topics.sh --bootstrap-server localhost:9092 \
--describe
# Increase partition count (cannot decrease)
kafka-topics.sh --bootstrap-server localhost:9092 \
--alter \
--topic orders \
--partitions 12
# Delete a topic
kafka-topics.sh --bootstrap-server localhost:9092 \
--delete \
--topic orders
# List topics with specific configuration overrides
kafka-topics.sh --bootstrap-server localhost:9092 \
--describe \
--topics-with-overrides
Topic Configuration
# Set topic-level configuration
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type topics \
--entity-name orders \
--add-config retention.ms=604800000,cleanup.policy=delete
# Describe topic configuration
kafka-configs.sh --bootstrap-server localhost:9092 \
--describe \
--entity-type topics \
--entity-name orders
# Remove configuration override (revert to default)
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type topics \
--entity-name orders \
--delete-config retention.ms
Partition Assignment Strategy
| Strategy | Description | Use Case |
|---|---|---|
| Round Robin | Distributes messages evenly across partitions | No ordering requirements |
| Key-based | Messages with same key go to same partition | Ordering by entity (user, order) |
| Custom | Implement custom partitioner | Complex routing logic |
Example: Topic with Compaction
# Create compacted topic for state storage
kafka-topics.sh --bootstrap-server localhost:9092 \
--create \
--topic user-profiles \
--partitions 6 \
--replication-factor 3 \
--config cleanup.policy=compact \
--config min.cleanable.dirty.ratio=0.1 \
--config segment.ms=3600000
Producer Configuration and Patterns
Producers publish messages to Kafka topics with configurable delivery guarantees and performance characteristics.
Key Concepts
- Acknowledgements (acks): Controls durability guarantees
- Batching: Groups messages for efficient network usage
- Compression: Reduces message size on disk and network
- Idempotence: Prevents duplicate messages on retry
- Transactions: Atomic writes across multiple partitions
Producer Configuration
# Essential producer configuration
bootstrap.servers=broker1:9092,broker2:9092,broker3:9092
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer
# Reliability settings
acks=all # Wait for all replicas
retries=2147483647 # Retry indefinitely
retry.backoff.ms=100 # Wait between retries
delivery.timeout.ms=120000 # Total time for send
enable.idempotence=true # Prevent duplicates
# Performance settings
batch.size=16384 # Batch size in bytes
linger.ms=5 # Wait for batch to fill
buffer.memory=33554432 # Total buffer memory
compression.type=lz4 # Compression algorithm
# Transaction settings (exactly-once)
transactional.id=my-transactional-id
transaction.timeout.ms=60000
Acknowledgement Levels
| acks | Durability | Latency | Description |
|---|---|---|---|
| 0 | None | Lowest | Fire and forget; no confirmation |
| 1 | Leader only | Medium | Leader acknowledges; may lose data |
| all/-1 | Full | Highest | All in-sync replicas acknowledge |
Console Producer
# Basic producer
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic orders
# Producer with key-value pairs
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic orders \
--property parse.key=true \
--property key.separator=:
# Producer with custom serialiser
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic orders \
--property value.serializer=org.apache.kafka.common.serialization.IntegerSerializer
# Producer with acknowledgement
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic orders \
--producer-property acks=all
# Send from file
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic orders < messages.txt
Java Producer Example
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
public class OrderProducer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("acks", "all");
props.put("enable.idempotence", "true");
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
// Synchronous send
ProducerRecord<String, String> record =
new ProducerRecord<>("orders", "order-123", "{\"item\": \"widget\"}");
RecordMetadata metadata = producer.send(record).get();
System.out.printf("Sent to partition %d, offset %d%n",
metadata.partition(), metadata.offset());
// Asynchronous send with callback
producer.send(record, (metadata, exception) -> {
if (exception != null) {
exception.printStackTrace();
} else {
System.out.printf("Sent to partition %d%n", metadata.partition());
}
});
}
}
}
Transactional Producer Example
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("transactional.id", "order-processor");
props.put("enable.idempotence", "true");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("orders", "key1", "value1"));
producer.send(new ProducerRecord<>("inventory", "key2", "value2"));
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
}
Consumer Groups and Offsets
Consumers read messages from topics, organised into consumer groups for parallel processing and fault tolerance.
Key Concepts
flowchart LR
subgraph Topic: orders
P0[Partition 0]
P1[Partition 1]
P2[Partition 2]
P3[Partition 3]
end
subgraph Consumer Group: order-processor
C1[Consumer 1]
C2[Consumer 2]
end
P0 --> C1
P1 --> C1
P2 --> C2
P3 --> C2
- Consumer Group: Logical grouping of consumers that share consumption
- Group Coordinator: Broker managing group membership and partition assignment
- Partition Assignment: Each partition assigned to exactly one consumer in a group
- Offset: Position of consumer in a partition; stored in __consumer_offsets topic
- Rebalancing: Redistribution of partitions when consumers join/leave
Consumer Configuration
# Essential consumer configuration
bootstrap.servers=broker1:9092,broker2:9092,broker3:9092
group.id=order-processor
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
# Offset management
enable.auto.commit=false # Manual commit for control
auto.offset.reset=earliest # Start from beginning if no offset
max.poll.records=500 # Records per poll
# Session management
session.timeout.ms=45000 # Consumer failure detection
heartbeat.interval.ms=3000 # Heartbeat frequency
max.poll.interval.ms=300000 # Max processing time
# Performance settings
fetch.min.bytes=1 # Minimum fetch size
fetch.max.wait.ms=500 # Max wait for fetch.min.bytes
max.partition.fetch.bytes=1048576 # Max data per partition per fetch
Auto Offset Reset Options
| Setting | Behaviour | Use Case |
|---|---|---|
| earliest | Read from beginning | Reprocessing, new consumer groups |
| latest | Read only new messages | Real-time processing |
| none | Throw exception if no offset | Strict offset management |
Console Consumer
# Basic consumer
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic orders \
--group order-processor
# Consumer from beginning
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic orders \
--from-beginning
# Consumer with key and value
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic orders \
--property print.key=true \
--property print.timestamp=true \
--property key.separator=,
# Consumer specific partitions
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic orders \
--partition 0 \
--offset 100
# Consumer with max messages
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic orders \
--max-messages 10
Managing Consumer Groups
# List all consumer groups
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
# Describe consumer group (show lag)
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe \
--group order-processor
# Describe all groups
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe \
--all-groups
# Show group state
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe \
--group order-processor \
--state
# Show members and assignments
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe \
--group order-processor \
--members --verbose
Offset Management
# Reset offset to earliest
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-processor \
--topic orders \
--reset-offsets \
--to-earliest \
--execute
# Reset offset to latest
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-processor \
--topic orders \
--reset-offsets \
--to-latest \
--execute
# Reset to specific offset
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-processor \
--topic orders \
--reset-offsets \
--to-offset 1000 \
--execute
# Reset by timestamp
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-processor \
--topic orders \
--reset-offsets \
--to-datetime 2024-01-01T00:00:00.000 \
--execute
# Shift offset by number
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-processor \
--topic orders \
--reset-offsets \
--shift-by -100 \
--execute
# Delete consumer group
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--delete \
--group order-processor
Java Consumer Example
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.*;
public class OrderConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "order-processor");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest");
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
consumer.subscribe(Arrays.asList("orders"));
while (true) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("Partition: %d, Offset: %d, Key: %s, Value: %s%n",
record.partition(), record.offset(),
record.key(), record.value());
// Process message...
}
// Manual commit after processing
consumer.commitSync();
}
}
}
}
Configuration Essentials
Kafka requires proper broker and cluster configuration for production deployments.
ZooKeeper vs KRaft Mode
flowchart TB
subgraph ZooKeeper Mode
ZK[ZooKeeper Ensemble]
B1[Broker 1]
B2[Broker 2]
B3[Broker 3]
ZK --> B1
ZK --> B2
ZK --> B3
end
subgraph KRaft Mode
C1[Controller 1]
C2[Controller 2]
C3[Controller 3]
KB1[Broker 1]
KB2[Broker 2]
KB3[Broker 3]
C1 <--> C2
C2 <--> C3
C1 --> KB1
C2 --> KB2
C3 --> KB3
end
| Feature | ZooKeeper | KRaft |
|---|---|---|
| Architecture | External coordination | Built-in consensus |
| Complexity | Requires separate cluster | Simplified operations |
| Scalability | Limited by ZK | Better partition scaling |
| Status | Removed in 4.0 (deprecated since 3.5) | Recommended (GA in 3.3+) |
Essential Broker Configuration
# server.properties - Core settings
broker.id=0
listeners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://broker1.example.com:9092
log.dirs=/var/kafka/data
num.partitions=6
default.replication.factor=3
# ZooKeeper mode (Kafka 3.x only; removed in 4.0)
# zookeeper.connect=zk1:2181,zk2:2181,zk3:2181/kafka
# KRaft mode (required from Kafka 4.0)
process.roles=broker,controller
node.id=1
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093
controller.listener.names=CONTROLLER
Performance Tuning
# Network and I/O threads
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
# Log settings
num.recovery.threads.per.data.dir=2
log.flush.interval.messages=10000
log.flush.interval.ms=1000
# Replication
num.replica.fetchers=4
replica.fetch.max.bytes=1048576
replica.fetch.wait.max.ms=500
Security Configuration
# SSL/TLS encryption
listeners=SSL://0.0.0.0:9093
ssl.keystore.location=/var/kafka/ssl/kafka.keystore.jks
ssl.keystore.password=keystore-password
ssl.key.password=key-password
ssl.truststore.location=/var/kafka/ssl/kafka.truststore.jks
ssl.truststore.password=truststore-password
ssl.client.auth=required
# SASL authentication
listeners=SASL_SSL://0.0.0.0:9094
sasl.enabled.mechanisms=SCRAM-SHA-512
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-512
security.inter.broker.protocol=SASL_SSL
# Authorisation (KRaft)
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
super.users=User:admin
KRaft Controller Configuration
# controller.properties
process.roles=controller
node.id=1
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093
listeners=CONTROLLER://0.0.0.0:9093
controller.listener.names=CONTROLLER
log.dirs=/var/kafka/controller-logs
Cluster Management Commands
# Get cluster ID
kafka-cluster.sh cluster-id --bootstrap-server localhost:9092
# Describe KRaft quorum status
kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe --status
# Update broker configuration dynamically
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type brokers \
--entity-name 0 \
--add-config log.cleaner.threads=2
# Describe broker configuration
kafka-configs.sh --bootstrap-server localhost:9092 \
--describe \
--entity-type brokers \
--entity-name 0
Common CLI Tools
Kafka provides command-line tools for administration, testing, and debugging.
Essential CLI Tools Reference
| Tool | Purpose |
|---|---|
| kafka-topics.sh | Create, list, describe, alter, delete topics |
| kafka-console-producer.sh | Send messages from command line |
| kafka-console-consumer.sh | Read messages from command line |
| kafka-consumer-groups.sh | Manage consumer groups and offsets |
| kafka-configs.sh | Manage entity configurations |
| kafka-acls.sh | Manage access control lists |
| kafka-reassign-partitions.sh | Reassign partition replicas |
| kafka-log-dirs.sh | Query log directory information |
Performance Testing
# Producer performance test
kafka-producer-perf-test.sh --topic perf-test \
--num-records 1000000 \
--record-size 1024 \
--throughput -1 \
--producer-props bootstrap.servers=localhost:9092
# Consumer performance test
kafka-consumer-perf-test.sh --bootstrap-server localhost:9092 \
--topic perf-test \
--messages 1000000
# End-to-end latency test
kafka-e2e-latency.sh localhost:9092 latency-test 10000 all 1024
Log and Data Inspection
# Dump log segments
kafka-dump-log.sh --files /var/kafka/data/orders-0/00000000000000000000.log \
--print-data-log
# Get log directory information
kafka-log-dirs.sh --bootstrap-server localhost:9092 \
--describe \
--topic-list orders
# Check log end offset
kafka-get-offsets.sh --bootstrap-server localhost:9092 \
--topic orders \
--time latest
Partition Reassignment
# Generate reassignment plan
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--topics-to-move-json-file topics.json \
--broker-list "1,2,3" \
--generate
# Execute reassignment
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file reassignment.json \
--execute
# Verify reassignment
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file reassignment.json \
--verify
ACL Management
# Add producer ACL
kafka-acls.sh --bootstrap-server localhost:9092 \
--add \
--allow-principal User:producer \
--operation Write \
--topic orders
# Add consumer ACL
kafka-acls.sh --bootstrap-server localhost:9092 \
--add \
--allow-principal User:consumer \
--operation Read \
--topic orders \
--group order-processor
# List ACLs
kafka-acls.sh --bootstrap-server localhost:9092 --list
# Remove ACL
kafka-acls.sh --bootstrap-server localhost:9092 \
--remove \
--allow-principal User:producer \
--operation Write \
--topic orders
Retention and Compaction
Kafka supports multiple data retention strategies based on time, size, or key-based compaction.
Retention Strategies
flowchart TB
subgraph Delete Policy
D1[Segment 1<br/>Oldest] --> D2[Segment 2] --> D3[Segment 3<br/>Active]
D1 -.->|Delete after<br/>retention.ms| X1[Deleted]
end
subgraph Compact Policy
C1[Key: A, V1] --> C2[Key: B, V1] --> C3[Key: A, V2] --> C4[Key: B, V2]
C4 --> R[Key: A, V2<br/>Key: B, V2]
end
| Policy | Description | Use Case |
|---|---|---|
| delete | Remove segments after retention period/size | Event logs, metrics |
| compact | Keep latest value per key | State stores, changelogs |
| compact,delete | Compact then delete after retention | Bounded state stores |
Retention Configuration
# Time-based retention (7 days)
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type topics \
--entity-name events \
--add-config retention.ms=604800000
# Size-based retention (10 GB)
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type topics \
--entity-name events \
--add-config retention.bytes=10737418240
# Configure log segment size and roll time
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type topics \
--entity-name events \
--add-config segment.bytes=1073741824,segment.ms=86400000
Compaction Configuration
# Enable compaction
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type topics \
--entity-name user-state \
--add-config cleanup.policy=compact
# Compaction tuning
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter \
--entity-type topics \
--entity-name user-state \
--add-config \
min.cleanable.dirty.ratio=0.5,\
min.compaction.lag.ms=3600000,\
max.compaction.lag.ms=86400000,\
delete.retention.ms=86400000
Compaction Settings Explained
| Setting | Default | Description |
|---|---|---|
| min.cleanable.dirty.ratio | 0.5 | Minimum ratio of dirty to total logs to trigger compaction |
| min.compaction.lag.ms | 0 | Minimum time message stays uncompacted |
| max.compaction.lag.ms | ∞ | Maximum time before compaction guaranteed |
| delete.retention.ms | 86400000 | Time to retain tombstone markers |
| segment.ms | 604800000 | Time before rolling new segment |
Tombstones for Deletion
In compacted topics, send a message with null value to delete a key:
// Java - Delete key from compacted topic
producer.send(new ProducerRecord<>("user-state", "user-123", null));
# Console producer - send tombstone
echo "user-123:" | kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic user-state \
--property parse.key=true \
--property key.separator=: \
--property null.marker=NULL
Monitoring and Troubleshooting
Effective monitoring is critical for maintaining Kafka cluster health and performance.
Key Metrics to Monitor
| Category | Metric | Description | Alert Threshold |
|---|---|---|---|
| Broker | UnderReplicatedPartitions | Partitions without full replication | > 0 |
| Broker | ActiveControllerCount | Number of active controllers | != 1 |
| Broker | OfflinePartitionsCount | Partitions without leader | > 0 |
| Broker | RequestHandlerAvgIdlePercent | Request handler thread utilisation | < 0.3 |
| Producer | record-error-rate | Failed message sends | > 0 |
| Producer | request-latency-avg | Average request latency | > 100ms |
| Consumer | records-lag-max | Maximum consumer lag | Depends on SLA |
| Consumer | fetch-latency-avg | Average fetch latency | > 500ms |
JMX Metrics Access
# Enable JMX on broker (in kafka-server-start.sh or environment)
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote \
-Dcom.sun.management.jmxremote.port=9999 \
-Dcom.sun.management.jmxremote.authenticate=false \
-Dcom.sun.management.jmxremote.ssl=false"
# Query JMX metrics with jmxterm
java -jar jmxterm.jar -l localhost:9999
> bean kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions
> get Value
Monitoring Consumer Lag
# Check consumer lag for all groups
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe \
--all-groups
# Check specific group lag
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe \
--group order-processor
# Output includes: TOPIC, PARTITION, CURRENT-OFFSET, LOG-END-OFFSET, LAG
Log Analysis
# Check broker logs
tail -f /var/log/kafka/server.log
# Search for errors
grep -i error /var/log/kafka/server.log
# Check controller logs
tail -f /var/log/kafka/controller.log
# Check state change logs
tail -f /var/log/kafka/state-change.log
Health Checks
# Check if broker is accepting connections
nc -zv localhost 9092
# Verify cluster metadata (KRaft)
kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe --status
# Check broker registration
kafka-broker-api-versions.sh --bootstrap-server localhost:9092
# Verify topic replication
kafka-topics.sh --bootstrap-server localhost:9092 \
--describe \
--under-replicated-partitions
Prometheus and Grafana Integration
# prometheus.yml - Kafka JMX exporter configuration
scrape_configs:
- job_name: 'kafka'
static_configs:
- targets: ['broker1:7071', 'broker2:7071', 'broker3:7071']
# jmx_exporter_config.yml
lowercaseOutputName: true
rules:
- pattern: kafka.server<type=(.+), name=(.+)><>Value
name: kafka_server_$1_$2
type: GAUGE
- pattern: kafka.server<type=(.+), name=(.+), topic=(.+)><>Value
name: kafka_server_$1_$2
type: GAUGE
labels:
topic: "$3"
Quick Reference
Topic Operations
kafka-topics.sh --bootstrap-server localhost:9092 --list # List topics
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic t1 --partitions 3 --replication-factor 3
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic t1 # Describe topic
kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic t1 # Delete topic
Producer Operations
kafka-console-producer.sh --bootstrap-server localhost:9092 --topic t1 # Basic producer
kafka-console-producer.sh --bootstrap-server localhost:9092 --topic t1 \
--property parse.key=true --property key.separator=: # With keys
Consumer Operations
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic t1 --from-beginning
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic t1 --group g1
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group g1
Offset Management
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group g1 \
--topic t1 --reset-offsets --to-earliest --execute # Reset to start
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group g1 \
--topic t1 --reset-offsets --to-latest --execute # Reset to end
Configuration Management
kafka-configs.sh --bootstrap-server localhost:9092 --describe \
--entity-type topics --entity-name t1 # Show config
kafka-configs.sh --bootstrap-server localhost:9092 --alter \
--entity-type topics --entity-name t1 --add-config retention.ms=86400000 # Set config
Common Issues and Solutions
Under-Replicated Partitions
Symptoms: UnderReplicatedPartitions metric > 0, data loss risk
Causes and solutions:
-
Broker failure - Check broker health and restart if needed
kafka-topics.sh --bootstrap-server localhost:9092 \ --describe --under-replicated-partitions -
Network issues - Check connectivity between brokers
nc -zv broker2 9092 -
Disk full - Check disk space and increase or clean up
df -h /var/kafka/data kafka-log-dirs.sh --bootstrap-server localhost:9092 --describe -
Slow follower - Check replica fetch configuration
kafka-configs.sh --bootstrap-server localhost:9092 \ --describe --entity-type brokers --entity-name 1
Consumer Lag Growing
Symptoms: Consumer lag continuously increasing, processing delays
Causes and solutions:
-
Slow processing - Increase consumer instances or optimise processing
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --describe --group order-processor -
Too few consumers - Add more consumers (up to partition count)
-
Inefficient polling - Tune consumer configuration
max.poll.records=1000 fetch.max.bytes=52428800 -
Rebalancing storms - Increase session timeout
session.timeout.ms=60000 heartbeat.interval.ms=20000
Producer Timeouts
Symptoms: TimeoutException, messages not delivered
Causes and solutions:
-
Broker unavailable - Check broker health and connectivity
kafka-broker-api-versions.sh --bootstrap-server localhost:9092 -
Buffer full - Increase buffer or reduce send rate
buffer.memory=67108864 max.block.ms=120000 -
Acks=all with failed replicas - Check ISR
kafka-topics.sh --bootstrap-server localhost:9092 \ --describe --topic orders
Rebalancing Issues
Symptoms: Frequent rebalances, consumer instability
Causes and solutions:
-
Long processing time - Increase max.poll.interval.ms
max.poll.interval.ms=600000 -
Session timeout - Adjust timeout settings
session.timeout.ms=45000 heartbeat.interval.ms=15000 -
Static membership - Use static group membership
group.instance.id=consumer-1
Disk Space Issues
Symptoms: Broker failing, log directory errors
Causes and solutions:
-
Retention too long - Reduce retention period
kafka-configs.sh --bootstrap-server localhost:9092 \ --alter --entity-type topics --entity-name orders \ --add-config retention.ms=86400000 -
Large segments - Reduce segment size for faster cleanup
kafka-configs.sh --bootstrap-server localhost:9092 \ --alter --entity-type topics --entity-name orders \ --add-config segment.bytes=536870912 -
Delete old data manually - Use retention with immediate effect
# Set very short retention temporarily kafka-configs.sh --bootstrap-server localhost:9092 \ --alter --entity-type topics --entity-name orders \ --add-config retention.ms=1000 # Wait for deletion, then restore kafka-configs.sh --bootstrap-server localhost:9092 \ --alter --entity-type topics --entity-name orders \ --delete-config retention.ms
Connection/Authentication Failures
Symptoms: Connection refused, authentication errors
Causes and solutions:
-
Wrong bootstrap servers - Verify advertised.listeners configuration
kafka-configs.sh --bootstrap-server localhost:9092 \ --describe --entity-type brokers --entity-name 0 \ --all | grep advertised -
SSL/TLS misconfiguration - Check certificates and trust stores
openssl s_client -connect broker:9093 -CAfile ca.crt -
SASL credentials - Verify JAAS configuration
# Client JAAS config sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \ username="user" password="password";
Related Topics
The following topics complement Apache Kafka knowledge and are commonly used together:
- Kafka Connect - Framework for streaming data between Kafka and external systems; connectors for databases, cloud services, and file systems
- Kafka Streams - Client library for building stream processing applications; stateful transformations and windowing operations
- Schema Registry - Central repository for message schemas (Avro, Protobuf, JSON Schema); ensures data compatibility
- Apache Flink - Distributed stream processing engine; complex event processing and real-time analytics with Kafka
- Kubernetes - Container orchestration for deploying Kafka clusters; Strimzi operator for Kafka on K8s
- Prometheus/Grafana - Monitoring stack commonly used with Kafka; JMX metrics export and visualisation