Kafka 2.8集群扩容后批量提升副本因子至4的自动化脚本咨询
Kafka集群扩容副本重分配自动化脚本(生成JSON payload)
脚本说明
针对开源Kafka 2.8版本,以下Python脚本可自动生成副本重分配的JSON配置,实现将12节点集群中所有topic的副本因子从3提升至4,同时保证副本在所有broker间均匀分布。
核心逻辑
- 遍历集群内所有业务topic(排除Kafka内部topic)
- 对每个分区,在原有3个副本基础上新增1个副本节点,通过计数均衡策略选择当前副本数最少的可用broker
- 输出符合
kafka-reassign-partitions.sh工具要求的JSON格式重分配计划
Python脚本实现
import json from kafka import KafkaAdminClient from kafka.errors import KafkaError # 集群配置参数,需替换为实际值 BOOTSTRAP_SERVERS = "broker0:9092,broker1:9092,broker2:9092,broker3:9092,broker4:9092,broker5:9092,broker6:9092,broker7:9092,broker8:9092,broker9:9092,broker10:9092,broker11:9092" OLD_BROKER_IDS = [0,1,2,3,4,5] # 原6个broker的ID ALL_BROKER_IDS = [0,1,2,3,4,5,6,7,8,9,10,11] # 扩容后12个broker的ID TARGET_REPLICA_FACTOR = 4 def get_all_business_topics(): """获取集群内所有非内部业务topic""" admin_client = KafkaAdminClient(bootstrap_servers=BOOTSTRAP_SERVERS) all_topics = admin_client.list_topics() # 排除Kafka系统内部topic exclude_topics = {"__consumer_offsets", "__transaction_state", "__schema_registry"} return [topic for topic in all_topics if topic not in exclude_topics] def get_topic_partition_info(topic_name): """获取指定topic的分区数及每个分区的现有副本列表""" admin_client = KafkaAdminClient(bootstrap_servers=BOOTSTRAP_SERVERS) topic_meta = admin_client.describe_topics([topic_name])[0] partition_info = [] for p in topic_meta["partitions"]: partition_info.append({ "partition_id": p["partition"], "current_replicas": p["replicas"] }) return partition_info def generate_reassignment_plan(): """生成最终的副本重分配JSON计划""" reassignment_plan = {"version": 1, "partitions": []} # 初始化每个broker的副本计数,先统计现有副本数量 broker_replica_counts = {bid:0 for bid in ALL_BROKER_IDS} # 先遍历所有topic,统计现有副本分布 topics = get_all_business_topics() for topic in topics: partition_details = get_topic_partition_info(topic) for p in partition_details: for replica in p["current_replicas"]: broker_replica_counts[replica] += 1 # 再次遍历,为每个分区新增副本并生成计划 for topic in topics: partition_details = get_topic_partition_info(topic) for p in partition_details: current_replicas = p["current_replicas"] # 找到当前副本数最少且不在现有副本列表中的broker candidate_brokers = sorted( [bid for bid in ALL_BROKER_IDS if bid not in current_replicas], key=lambda x: broker_replica_counts[x] ) new_replica = candidate_brokers[0] # 更新副本计数 broker_replica_counts[new_replica] += 1 # 生成该分区的重分配配置 reassignment_plan["partitions"].append({ "topic": topic, "partition": p["partition_id"], "replicas": current_replicas + [new_replica] }) return reassignment_plan if __name__ == "__main__": try: plan = generate_reassignment_plan() with open("kafka_reassignment_plan.json", "w", encoding="utf-8") as f: json.dump(plan, f, indent=2) print("✅ 副本重分配计划已生成至 kafka_reassignment_plan.json") except KafkaError as e: print(f"❌ Kafka客户端错误: {str(e)}") except Exception as e: print(f"❌ 脚本执行错误: {str(e)}")
使用步骤
- 安装依赖:
(指定版本适配Kafka 2.8的兼容性)pip install kafka-python==2.0.2 - 修改脚本中
BOOTSTRAP_SERVERS、OLD_BROKER_IDS为集群实际配置 - 运行脚本生成JSON计划文件
- 执行副本重分配:
bin/kafka-reassign-partitions.sh --bootstrap-server broker0:9092 --reassignment-json-file kafka_reassignment_plan.json --execute - 验证重分配进度:
bin/kafka-reassign-partitions.sh --bootstrap-server broker0:9092 --reassignment-json-file kafka_reassignment_plan.json --verify
注意事项
- 大规模副本重分配会占用集群IO和网络资源,建议在业务低峰期执行
- 可先选取单个小topic测试脚本逻辑,确认无误后再全量执行
- 若集群原有副本分布不符合默认逻辑,脚本中通过
describe_topics获取实际副本列表的逻辑已保证准确性
内容的提问来源于stack exchange,提问作者maverick
相关产品推荐
相关产品推荐

