You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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)}")

使用步骤

  1. 安装依赖:
    pip install kafka-python==2.0.2
    
    (指定版本适配Kafka 2.8的兼容性)
  2. 修改脚本中BOOTSTRAP_SERVERS、OLD_BROKER_IDS为集群实际配置
  3. 运行脚本生成JSON计划文件
  4. 执行副本重分配:
    bin/kafka-reassign-partitions.sh --bootstrap-server broker0:9092 --reassignment-json-file kafka_reassignment_plan.json --execute
    
  5. 验证重分配进度:
    bin/kafka-reassign-partitions.sh --bootstrap-server broker0:9092 --reassignment-json-file kafka_reassignment_plan.json --verify
    

注意事项

  • 大规模副本重分配会占用集群IO和网络资源,建议在业务低峰期执行
  • 可先选取单个小topic测试脚本逻辑,确认无误后再全量执行
  • 若集群原有副本分布不符合默认逻辑,脚本中通过describe_topics获取实际副本列表的逻辑已保证准确性

内容的提问来源于stack exchange,提问作者maverick

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 17:55:25