Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

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.

CoordinationClientsKafka ClusterBroker 2Broker 1Broker 3Topic APartition 0FollowerTopic BPartition 0LeaderTopic APartition 0LeaderTopic BPartition 1FollowerTopic APartition 1LeaderTopic BPartition 0FollowerProducersConsumer GroupKRaft ControllerQuorumCoordinationClientsKafka ClusterBroker 2Broker 1Broker 3Topic APartition 0FollowerTopic BPartition 0LeaderTopic APartition 0LeaderTopic BPartition 1FollowerTopic APartition 1LeaderTopic BPartition 0FollowerProducersConsumer GroupKRaft ControllerQuorum

Topics and Partitions

Topics are logical channels for organising messages, while partitions enable parallel processing and horizontal scaling.

Key Concepts

Partition 0 DetailTopic: ordersPartition 0Partition 1Partition 2Offset 0Offset 1Offset 2Offset 3Partition 0 DetailTopic: ordersPartition 0Partition 1Partition 2Offset 0Offset 1Offset 2Offset 3
  • 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

Consumer Group: order-processorTopic: ordersPartition 0Partition 1Partition 2Partition 3Consumer 1Consumer 2Consumer Group: order-processorTopic: ordersPartition 0Partition 1Partition 2Partition 3Consumer 1Consumer 2
  • 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

KRaft ModeController 1Controller 2Controller 3Broker 1Broker 2Broker 3ZooKeeper ModeZooKeeper EnsembleBroker 1Broker 2Broker 3KRaft ModeController 1Controller 2Controller 3Broker 1Broker 2Broker 3ZooKeeper ModeZooKeeper EnsembleBroker 1Broker 2Broker 3
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

Compact PolicyKey: A, V1Key: B, V1Key: A, V2Key: B, V2Key: A, V2Key: B, V2Delete PolicyDelete afterretention.msSegment 1OldestSegment 2Segment 3ActiveDeletedCompact PolicyKey: A, V1Key: B, V1Key: A, V2Key: B, V2Key: A, V2Key: B, V2Delete PolicyDelete afterretention.msSegment 1OldestSegment 2Segment 3ActiveDeleted
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:

  1. Broker failure - Check broker health and restart if needed

    kafka-topics.sh --bootstrap-server localhost:9092 \
      --describe --under-replicated-partitions
    
  2. Network issues - Check connectivity between brokers

    nc -zv broker2 9092
    
  3. Disk full - Check disk space and increase or clean up

    df -h /var/kafka/data
    kafka-log-dirs.sh --bootstrap-server localhost:9092 --describe
    
  4. 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:

  1. Slow processing - Increase consumer instances or optimise processing

    kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
      --describe --group order-processor
    
  2. Too few consumers - Add more consumers (up to partition count)

  3. Inefficient polling - Tune consumer configuration

    max.poll.records=1000
    fetch.max.bytes=52428800
    
  4. Rebalancing storms - Increase session timeout

    session.timeout.ms=60000
    heartbeat.interval.ms=20000
    

Producer Timeouts

Symptoms: TimeoutException, messages not delivered

Causes and solutions:

  1. Broker unavailable - Check broker health and connectivity

    kafka-broker-api-versions.sh --bootstrap-server localhost:9092
    
  2. Buffer full - Increase buffer or reduce send rate

    buffer.memory=67108864
    max.block.ms=120000
    
  3. 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:

  1. Long processing time - Increase max.poll.interval.ms

    max.poll.interval.ms=600000
    
  2. Session timeout - Adjust timeout settings

    session.timeout.ms=45000
    heartbeat.interval.ms=15000
    
  3. Static membership - Use static group membership

    group.instance.id=consumer-1
    

Disk Space Issues

Symptoms: Broker failing, log directory errors

Causes and solutions:

  1. 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
    
  2. 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
    
  3. 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:

  1. Wrong bootstrap servers - Verify advertised.listeners configuration

    kafka-configs.sh --bootstrap-server localhost:9092 \
      --describe --entity-type brokers --entity-name 0 \
      --all | grep advertised
    
  2. SSL/TLS misconfiguration - Check certificates and trust stores

    openssl s_client -connect broker:9093 -CAfile ca.crt
    
  3. 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