跨多DC Kafka部署方案:主集群故障时如何实现生产者自动切换
Kafka集群容灾Producer自动化切换落地方案
前置准备
切换逻辑生效的前提是容灾集群数据和主集群对齐,需提前完成以下配置:
- 部署MirrorMaker2服务,实时同步主集群的Topic元数据、消息数据、消费位点、ACL策略到次级DR集群
- 提前完成DR集群的性能压测,确保容量可以承接主集群的全量流量
- 定义明确的切换触发阈值,避免小波动导致误切,推荐阈值参考:主集群 Producer 发送失败率持续10s高于5%、元数据请求连续3次超时、可用Broker数低于副本最小同步数
两种主流实现方案
方案1:客户端侧内置切换逻辑(轻量无额外依赖)
适合中小规模集群,无需额外部署中间件,业务代码零改造:
- 核心逻辑:
- 预先在Producer配置中维护主、DR集群的
bootstrap.servers地址 - 封装原生Kafka Producer,内置失败计数、主集群健康检测逻辑
- 触发切换阈值时,自动销毁原有Producer实例,新建指向DR集群的Producer实例
- 可选配置自动切回逻辑:定时检测主集群状态,恢复正常后自动切回主集群
- 预先在Producer配置中维护主、DR集群的
- 优缺点:
- 优点:部署成本低,切换延迟最低可达毫秒级
- 缺点:客户端独立判断集群状态,极端情况会出现部分客户端切换、部分未切换的不一致问题
方案2:全局控制平面统一切换(适合大规模生产环境)
适合集群节点数过百、业务线多的生产场景:
- 核心逻辑:
- 独立部署容灾控制服务,统一采集主集群所有Broker、Topic的运行状态
- 主集群异常时,通过配置中心(如Nacos、etcd、Consul)向所有Producer客户端推送切换指令
- 客户端监听配置变更,收到指令后执行集群切换动作,同时上报切换状态到控制平面
- 优缺点:
- 优点:切换动作全局一致,支持人工干预、灰度流量切换等复杂策略
- 缺点:需要额外部署控制服务、配置中心,架构复杂度更高
客户端侧切换Java示例代码
以下是基于kafka-clients 3.4.0版本封装的支持自动切换的Producer示例:
import org.apache.kafka.clients.producer.*; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; public class DisasterRecoveryKafkaProducer { private static final String PRIMARY_BOOTSTRAP = "primary-node1:9092,primary-node2:9092,primary-node3:9092"; private static final String DR_BOOTSTRAP = "dr-node1:9092,dr-node2:9092,dr-node3:9092"; private static final int MAX_FAIL_COUNT = 5; private static final long HEALTH_CHECK_INTERVAL = 30000; private volatile KafkaProducer<String, String> producer; private volatile boolean useDrCluster = false; private final AtomicInteger failCount = new AtomicInteger(0); private final Object lock = new Object(); private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); public DisasterRecoveryKafkaProducer() { producer = createProducer(PRIMARY_BOOTSTRAP); // 启动定时健康检查,主集群恢复后自动切回 scheduler.scheduleAtFixedRate(this::checkPrimaryClusterHealth, 10, HEALTH_CHECK_INTERVAL, TimeUnit.MILLISECONDS); } private KafkaProducer<String, String> createProducer(String bootstrapServers) { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); return new KafkaProducer<>(props); } public void send(String topic, String key, String value, Callback callback) { try { producer.send(new ProducerRecord<>(topic, key, value), (metadata, exception) -> { if (exception != null) { int currentFail = failCount.incrementAndGet(); if (currentFail >= MAX_FAIL_COUNT && !useDrCluster) { switchToDrCluster(); } // 切换后重试发送 if (useDrCluster) { producer.send(new ProducerRecord<>(topic, key, value), callback); } else { callback.onCompletion(metadata, exception); } } else { failCount.set(0); callback.onCompletion(metadata, null); } }); } catch (Exception e) { failCount.incrementAndGet(); if (failCount.get() >= MAX_FAIL_COUNT && !useDrCluster) { switchToDrCluster(); producer.send(new ProducerRecord<>(topic, key, value), callback); } } } private void switchToDrCluster() { synchronized (lock) { if (useDrCluster) return; producer.close(); producer = createProducer(DR_BOOTSTRAP); useDrCluster = true; failCount.set(0); } } private void checkPrimaryClusterHealth() { if (!useDrCluster) return; try (KafkaProducer<String, String> testProducer = createProducer(PRIMARY_BOOTSTRAP)) { // 尝试拉取主集群元数据判断是否恢复 testProducer.partitionsFor("health_check_topic"); synchronized (lock) { producer.close(); producer = createProducer(PRIMARY_BOOTSTRAP); useDrCluster = false; failCount.set(0); } } catch (Exception ignored) { // 主集群未恢复,保持DR集群状态 } } public void close() { scheduler.shutdown(); producer.close(); } }
注意:如果使用事务消息,需保证事务的开启、提交/回滚都在同一个集群内完成,不要跨集群执行事务操作;切换过程中发送失败的消息建议落本地盘做兜底,避免消息丢失。
生产环境落地注意事项
- 切换前必须验证DR集群的Topic分区数、副本数、保留时间、ACL策略和主集群完全一致
- 可配置梯度切换规则,首次切换先导10%流量到DR集群,稳定后再全量切换,避免瞬间流量把DR集群打崩
- 每月至少执行1次容灾切换演练,验证切换逻辑和DR集群可用性
- 业务侧需做好消息消费幂等性适配,避免切换过程中消息重复发送导致的业务问题
内容的提问来源于stack exchange,提问作者user16253628
相关产品推荐
相关产品推荐

