Kafka Maintenance
Routine maintenance procedures for Apache Kafka clusters.
Maintenance Overview
Section titled “Maintenance Overview”| Frequency | Tasks |
|---|---|
| Daily | Health checks, log monitoring |
| Weekly | Disk usage review, consumer lag review |
| Monthly | Partition balancing, configuration review |
| Quarterly | Capacity planning, security audit |
Log Cleanup
Section titled “Log Cleanup”Retention-Based Cleanup
Section titled “Retention-Based Cleanup”Kafka automatically deletes log segments based on retention settings.
# Time-based retention (default 7 days)log.retention.hours=168
# Size-based retention (per partition)log.retention.bytes=107374182400
# Check intervallog.retention.check.interval.ms=300000
# Segment settingslog.segment.bytes=1073741824log.segment.ms=604800000Manual Log Cleanup
Section titled “Manual Log Cleanup”# Delete records before offsetkafka-delete-records.sh --bootstrap-server kafka:9092 \ --offset-json-file delete-records.json
# delete-records.json content:# {# "partitions": [# {"topic": "my-topic", "partition": 0, "offset": 1000000},# {"topic": "my-topic", "partition": 1, "offset": 1000000}# ],# "version": 1# }Viewing Log Segment Information
Section titled “Viewing Log Segment Information”# List log directorieskafka-log-dirs.sh --bootstrap-server kafka:9092 \ --describe --broker-list 1,2,3
# Dump log segmentkafka-dump-log.sh --files /var/kafka-logs/my-topic-0/00000000000000000000.log \ --print-data-logLog Compaction
Section titled “Log Compaction”Compaction Configuration
Section titled “Compaction Configuration”# Enable compaction for topiccleanup.policy=compact
# Compaction settingslog.cleaner.enable=truelog.cleaner.threads=2log.cleaner.io.buffer.size=524288log.cleaner.dedupe.buffer.size=134217728log.cleaner.io.max.bytes.per.second=1.7976931348623157E308log.cleaner.min.cleanable.ratio=0.5log.cleaner.min.compaction.lag.ms=0log.cleaner.max.compaction.lag.ms=9223372036854775807Creating Compacted Topic
Section titled “Creating Compacted Topic”kafka-topics.sh --bootstrap-server kafka:9092 \ --create \ --topic compacted-topic \ --partitions 12 \ --replication-factor 3 \ --config cleanup.policy=compact \ --config min.cleanable.dirty.ratio=0.5 \ --config segment.ms=86400000Monitoring Compaction
Section titled “Monitoring Compaction”| Metric | Description |
|---|---|
log-cleaner-recopy-percent | Percentage of log recopy |
max-clean-time-secs | Maximum clean time |
max-buffer-utilization-percent | Buffer utilization |
Disk Management
Section titled “Disk Management”Monitoring Disk Usage
Section titled “Monitoring Disk Usage”#!/bin/bashKAFKA_LOG_DIR="/var/kafka-logs"THRESHOLD=80
usage=$(df -h "$KAFKA_LOG_DIR" | tail -1 | awk '{print $5}' | tr -d '%')
echo "Disk usage for $KAFKA_LOG_DIR: ${usage}%"
if [ "$usage" -gt "$THRESHOLD" ]; then echo "WARNING: Disk usage exceeds ${THRESHOLD}%"
# Show largest topics echo "Largest topics:" du -sh "$KAFKA_LOG_DIR"/*-* | sort -rh | head -10fiAdding Log Directory
Section titled “Adding Log Directory”# Multiple log directorieslog.dirs=/data1/kafka-logs,/data2/kafka-logs,/data3/kafka-logsMoving Partitions Between Disks
Section titled “Moving Partitions Between Disks”# Generate plan to use new diskkafka-reassign-partitions.sh --bootstrap-server kafka:9092 \ --topics-to-move-json-file topics.json \ --broker-list "1,2,3" \ --generate
# Verify disk distributionkafka-log-dirs.sh --bootstrap-server kafka:9092 \ --describe --broker-list 1Consumer Group Maintenance
Section titled “Consumer Group Maintenance”Dead Consumer Groups
Section titled “Dead Consumer Groups”# List groups with statekafka-consumer-groups.sh --bootstrap-server kafka:9092 \ --list --state
# Find empty groupskafka-consumer-groups.sh --bootstrap-server kafka:9092 \ --list --state | grep "Empty"
# Delete empty groupkafka-consumer-groups.sh --bootstrap-server kafka:9092 \ --delete --group dead-groupReset Consumer Offsets
Section titled “Reset Consumer Offsets”# Dry run firstkafka-consumer-groups.sh --bootstrap-server kafka:9092 \ --group my-group \ --reset-offsets \ --to-earliest \ --all-topics \ --dry-run
# Execute resetkafka-consumer-groups.sh --bootstrap-server kafka:9092 \ --group my-group \ --reset-offsets \ --to-earliest \ --all-topics \ --executeBroker Maintenance
Section titled “Broker Maintenance”Graceful Restart
Section titled “Graceful Restart”#!/bin/bashBROKER=$1BOOTSTRAP="kafka:9092"
echo "Restarting broker: $BROKER"
# Trigger controlled shutdownssh $BROKER "sudo systemctl stop kafka"
# Wait for leadership to transfersleep 10
# Check under-replicated partitionsURP=$(kafka-topics.sh --bootstrap-server $BOOTSTRAP \ --describe --under-replicated-partitions 2>/dev/null | wc -l)
echo "Under-replicated partitions: $URP"
# Start brokerssh $BROKER "sudo systemctl start kafka"
# Wait for recoveryecho "Waiting for broker to rejoin..."sleep 30
# Verify healthkafka-broker-api-versions.sh --bootstrap-server $BROKER:9092
# Wait for ISR to recoverwhile true; do URP=$(kafka-topics.sh --bootstrap-server $BOOTSTRAP \ --describe --under-replicated-partitions 2>/dev/null | wc -l) if [ "$URP" -eq 0 ]; then echo "All partitions in sync" break fi echo "Waiting for ISR recovery... ($URP under-replicated)" sleep 10done
echo "Restart complete"Preferred Leader Election
Section titled “Preferred Leader Election”# Rebalance leaders after maintenancekafka-leader-election.sh --bootstrap-server kafka:9092 \ --election-type preferred \ --all-topic-partitionsBroker and Log-Directory Cordoning (Kafka 4.3+)
Section titled “Broker and Log-Directory Cordoning (Kafka 4.3+)”Kafka 4.3 introduces a cordoning mechanism for brokers and individual log directories (KAFKA-19774). A cordoned broker or log directory continues to serve its existing replicas but is excluded from new partition placement decisions. This provides a non-disruptive way to stage decommissioning, evacuate a noisy disk, or prepare a broker for maintenance without immediately reassigning existing partitions.
| Cordon Target | Effect |
|---|---|
| Broker | No new partitions placed; existing replicas retained |
| Log directory | No new partition replicas placed on the directory; existing replicas retained |
Typical workflow:
- Cordon the broker or log directory.
- Plan and execute partition reassignment to move existing replicas off the cordoned target.
- Once the target holds no replicas, take it out of service.
Cordon Versus Decommission
Cordoning is reversible and non-destructive. Existing replicas remain available for reads and writes; only new placement is blocked. Uncordoning restores normal placement eligibility.
Certificate Rotation
Section titled “Certificate Rotation”Rotation Process
Section titled “Rotation Process”- Generate new certificates
- Add new CA to truststores
- Deploy updated truststores
- Generate new keystores with new certs
- Deploy new keystores
- Remove old CA from truststores
#!/bin/bash# 1. Add new CA to existing truststorekeytool -keystore kafka.truststore.jks \ -alias NewCARoot \ -import \ -file new-ca-cert.pem \ -storepass changeit \ -noprompt
# 2. Deploy truststores to all brokersfor broker in broker1 broker2 broker3; do scp kafka.truststore.jks $broker:/etc/kafka/ssl/done
# 3. Rolling restart (truststores only - no downtime)for broker in broker1 broker2 broker3; do ./graceful-restart.sh $brokerdone
# 4. Generate new keystores and deploy# ... (similar process for keystores)Index Maintenance
Section titled “Index Maintenance”Index Verification
Section titled “Index Verification”# Verify index integritykafka-dump-log.sh \ --files /var/kafka-logs/my-topic-0/00000000000000000000.log \ --index-sanity-check
# Deep iteration checkkafka-dump-log.sh \ --files /var/kafka-logs/my-topic-0/00000000000000000000.log \ --deep-iterationRebuilding Indexes
Section titled “Rebuilding Indexes”Kafka automatically rebuilds indexes on broker startup if they are missing or corrupt.
# Force index rebuild by removing index files# (broker must be stopped)rm /var/kafka-logs/my-topic-0/*.indexrm /var/kafka-logs/my-topic-0/*.timeindex
# Start broker - indexes will be rebuiltMaintenance Scripts
Section titled “Maintenance Scripts”Daily Health Check
Section titled “Daily Health Check”#!/bin/bashBOOTSTRAP="kafka:9092"DATE=$(date +%Y-%m-%d)LOG_FILE="/var/log/kafka-health/health-$DATE.log"
{ echo "=== Kafka Daily Health Check ===" echo "Date: $DATE" echo ""
# Broker connectivity echo "--- Broker Status ---" kafka-broker-api-versions.sh --bootstrap-server $BOOTSTRAP 2>/dev/null | head -5
# Offline partitions echo "" echo "--- Offline Partitions ---" kafka-topics.sh --bootstrap-server $BOOTSTRAP \ --describe --unavailable-partitions 2>/dev/null || echo "None"
# Under-replicated echo "" echo "--- Under-replicated Partitions ---" kafka-topics.sh --bootstrap-server $BOOTSTRAP \ --describe --under-replicated-partitions 2>/dev/null || echo "None"
# Consumer lag echo "" echo "--- Consumer Groups with Lag ---" kafka-consumer-groups.sh --bootstrap-server $BOOTSTRAP --list 2>/dev/null | \ while read group; do kafka-consumer-groups.sh --bootstrap-server $BOOTSTRAP \ --describe --group "$group" 2>/dev/null | \ awk -v g="$group" '$6 > 0 {print g, $2, $3, $6}' done
} >> "$LOG_FILE"Related Documentation
Section titled “Related Documentation”- Operations Overview - Operations guide
- Monitoring - Metrics and alerting
- Cluster Management - Cluster operations
- Backup and Restore - DR procedures