Skip to content

AxonOps — AI-Native Control Plane for Open Source Data Platforms

Kafka Cluster Management

Internal mechanisms for cluster coordination, metadata management, and administrative operations.


For detailed Raft consensus mechanics, controller failover, and ZooKeeper migration, see KRaft Deep Dive.

KRaft Controller QuorumBrokersController 1(Active Leader)Controller 2(Follower)Controller 3(Follower)Broker 1Broker 2Broker 3RaftRaftRaftmetadata updatesregistration,heartbeatmetadata updatesregistration,heartbeatmetadata updatesregistration,heartbeat

ResponsibilityDescription
Broker registrationTrack broker membership and liveness
Leader electionElect new partition leaders on failure
Partition assignmentAssign partitions to brokers
Topic managementCreate, delete, and modify topics
Configuration managementStore and distribute configurations
Metadata propagationDistribute cluster metadata to brokers
Client quota managementEnforce producer/consumer quotas

In KRaft mode, all metadata is stored in the __cluster_metadata internal log.

__cluster_metadata LogPartition 0BrokerRegistrationTopicCreatedPartitionAssignmentLeaderElectionConfigChangeControllerBrokersSingle partitionReplicated to all controllersBrokers fetch from controllerswritefetch/apply

Brokers do not replicate the metadata log; they fetch metadata updates from the controller quorum.

Common record types include:

Record TypeDescription
RegisterBrokerRecordBroker joins cluster
BrokerRegistrationChangeRecordBroker registration update
FenceBrokerRecordBroker fenced
UnfenceBrokerRecordBroker unfenced
UnregisterBrokerRecordBroker leaves cluster
RegisterControllerRecordController registration
TopicRecordTopic created
RemoveTopicRecordTopic removed
PartitionRecordPartition assignment
PartitionChangeRecordLeader/ISR change
ConfigRecordConfiguration change
ClientQuotaRecordQuota configuration
ProducerIdsRecordProducer ID allocation
FeatureLevelRecordFeature level changes
Terminal window
# Dump metadata log
kafka-metadata-shell.sh --snapshot /var/kafka-logs/__cluster_metadata-0/00000000000000000000.log \
--command "cat"
# Describe cluster
kafka-metadata-shell.sh --snapshot /var/kafka-logs/__cluster_metadata-0/00000000000000000000.log \
--command "describe"
# List brokers
kafka-metadata-shell.sh --snapshot /var/kafka-logs/__cluster_metadata-0/00000000000000000000.log \
--command "brokers"
# Show topic details
kafka-metadata-shell.sh --snapshot /var/kafka-logs/__cluster_metadata-0/00000000000000000000.log \
--command "topic" --topic-name my-topic

New BrokerControllerMetadata LogNew BrokerNew BrokerControllerControllerMetadata LogMetadata LogBrokerRegistrationRequestbroker.id, rack,listeners, featuresRegisterBrokerRecordBrokerRegistrationResponsebroker_epochloop[Heartbeat]BrokerHeartbeatRequestBrokerHeartbeatResponseHeartbeat maintains registrationMissing heartbeats trigger fencing
StateDescription
NOT_RUNNINGBroker not running
STARTINGBroker starting and catching up with metadata
RECOVERYBroker caught up, waiting to be unfenced
RUNNINGBroker registered and accepting requests
PENDING_CONTROLLED_SHUTDOWNBroker initiating controlled shutdown
SHUTTING_DOWNBroker shutting down
# Broker heartbeat interval
broker.heartbeat.interval.ms=2000
# Session timeout (broker considered dead if no heartbeat)
broker.session.timeout.ms=18000

For detailed leader election protocol, ISR mechanics, and leader epochs, see Replication.

TriggerDescription
Broker failureLeader broker becomes unavailable
Controlled shutdownBroker initiates graceful shutdown
Manual electionAdministrator triggers election
Preferred leaderAutomatic rebalancing to preferred replica
ControllerOld LeaderNew LeaderFollowerControllerControllerOld LeaderBroker 1Old LeaderBroker 1New LeaderBroker 2New LeaderBroker 2FollowerBroker 3FollowerBroker 3Detect broker 1 failureSelect new leader from ISRBroker 2 chosenLeaderAndIsrRequest(become leader)LeaderAndIsrRequest(new leader info)Accept writesFetch from new leader
Terminal window
# Trigger preferred leader election
kafka-leader-election.sh --bootstrap-server kafka:9092 \
--election-type preferred \
--all-topic-partitions
# For specific topic
kafka-leader-election.sh --bootstrap-server kafka:9092 \
--election-type preferred \
--topic my-topic
# Unclean election (data loss risk)
kafka-leader-election.sh --bootstrap-server kafka:9092 \
--election-type unclean \
--topic my-topic --partition 0

When creating a topic, partitions are assigned to brokers considering:

  • Rack awareness (distribute across racks)
  • Broker load (balance partition count)
  • Existing assignments (minimize movement)
Terminal window
# Create topic with specific assignment
kafka-topics.sh --bootstrap-server kafka:9092 \
--create \
--topic my-topic \
--replica-assignment 1:2:3,2:3:1,3:1:2
# Format: partition0_replicas,partition1_replicas,...
Terminal window
# Generate reassignment plan
cat > topics.json << 'EOF'
{
"topics": [{"topic": "my-topic"}],
"version": 1
}
EOF
kafka-reassign-partitions.sh --bootstrap-server kafka:9092 \
--topics-to-move-json-file topics.json \
--broker-list "1,2,3,4" \
--generate
# Execute plan
kafka-reassign-partitions.sh --bootstrap-server kafka:9092 \
--reassignment-json-file reassignment.json \
--throttle 100000000 \
--execute
# Verify completion
kafka-reassign-partitions.sh --bootstrap-server kafka:9092 \
--reassignment-json-file reassignment.json \
--verify

Configurations can be changed without broker restart.

Terminal window
# Broker configuration
kafka-configs.sh --bootstrap-server kafka:9092 \
--entity-type brokers \
--entity-name 1 \
--alter \
--add-config log.cleaner.threads=4
# Topic configuration
kafka-configs.sh --bootstrap-server kafka:9092 \
--entity-type topics \
--entity-name my-topic \
--alter \
--add-config retention.ms=86400000
# Client quota
kafka-configs.sh --bootstrap-server kafka:9092 \
--entity-type users \
--entity-name producer-user \
--alter \
--add-config producer_byte_rate=10485760
LevelDescription
Per-topicHighest priority, topic-specific
Per-brokerBroker-specific override
Cluster-wide dynamicDynamic default for all brokers
Static (server.properties)File-based configuration
DefaultBuilt-in defaults

ControllerBrokerControllerOther BrokersBrokerBrokerControllerControllerOther BrokersOther BrokersControllerBrokerHeartbeatRequest(want_shutdown=true)Elect new leaders foraffected partitionsLeaderAndIsrRequest(new leaders)AckBrokerHeartbeatResponse(permissions + state)ShutdownEnsures no data lossMaintains availability

ControlledShutdownRequest

ControlledShutdownRequest was removed in Kafka 4.0. In KRaft, brokers signal shutdown intent via BrokerHeartbeatRequest.

# Enable controlled shutdown
controlled.shutdown.enable=true
# Maximum retries
controlled.shutdown.max.retries=3
# Retry backoff
controlled.shutdown.retry.backoff.ms=5000

MetricDescriptionAlert Threshold
ActiveControllerCountActive controllers≠ 1
OfflinePartitionsCountOffline partitions> 0
UnderReplicatedPartitionsUnder-replicated partitions> 0
GlobalPartitionCountTotal partitionsGrowing unexpectedly
Terminal window
# Check controller
kafka-metadata-shell.sh --snapshot /var/kafka-logs/__cluster_metadata-0/*.log \
--command "describe" | grep -i controller
# Check offline partitions
kafka-topics.sh --bootstrap-server kafka:9092 \
--describe --unavailable-partitions
# Check under-replicated
kafka-topics.sh --bootstrap-server kafka:9092 \
--describe --under-replicated-partitions