如何使用Apache-Camel检查Kafka主题是否可用 缺失时跳过路由
问题说明
需要在Kafka-Camel路由执行前增加校验逻辑,判断目标Kafka主题是否可用,主题缺失时直接跳过路由流程,避免Kafka NetworkClient自动创建不存在的主题,消除如下告警日志:
2022-06-15 14:36:38 WARN [or.ap.ka.cl.NetworkClient] (kafka-producer-network-thread | producer-8) [Producer clientId=producer-8] Error while fetching metadata with correlation id 3 : {my-topic=LEADER_NOT_AVAILABLE}
落地步骤
1. 配置层关闭自动建主题开关
先从Kafka客户端配置层面兜底,禁止生产者自动创建主题,该参数从Kafka Client 2.3.0版本开始支持:
# 原生Kafka生产者配置 spring.kafka.producer.properties.allow.auto.create.topics=false # Camel Kafka组件专属配置 camel.component.kafka.configuration.additional-properties.allow.auto.create.topics=false
注意:不要依赖生产者拉取元数据的方式判断主题是否存在,该方式依然会触发自动建主题逻辑和对应告警,必须使用AdminClient做查询。
2. 实现前置校验逻辑
通过Camel的Processor扩展点,在路由进入Kafka端点前,用Kafka AdminClient查询目标主题是否存在,不存在则直接中断路由流程:
import org.apache.camel.Exchange; import org.apache.camel.Processor; import org.apache.kafka.clients.admin.AdminClient; import java.util.Set; public class KafkaTopicPreCheckProcessor implements Processor { private final AdminClient kafkaAdminClient; private final String targetTopic; // 构造传入AdminClient实例和待校验的目标主题名 public KafkaTopicPreCheckProcessor(AdminClient kafkaAdminClient, String targetTopic) { this.kafkaAdminClient = kafkaAdminClient; this.targetTopic = targetTopic; } @Override public void process(Exchange exchange) throws Exception { // 拉取集群所有已存在的主题 Set<String> existingTopics = kafkaAdminClient.listTopics().names().get(); if (!existingTopics.contains(targetTopic)) { // 标记中断当前路由,不进入后续Kafka发送流程 exchange.setProperty(Exchange.ROUTE_STOP, Boolean.TRUE); exchange.getIn().setHeader("KAFKA_TOPIC_NOT_EXIST", Boolean.TRUE); } } }
3. 挂载校验逻辑到路由
将校验处理器放在Kafka端点之前,保证校验先执行,校验通过才会触发Kafka消息发送:
from("你的业务源端点,比如timer、rest接口、其他MQ端点") // 最前端挂载主题校验逻辑 .process(new KafkaTopicPreCheckProcessor(kafkaAdminClient, "my-topic")) // 校验通过才执行Kafka发送 .to("kafka:my-topic?brokers=你的Kafka集群地址") // 可选:主题不存在时走兜底逻辑,比如记录日志、写入死信队列 .choice() .when(header("KAFKA_TOPIC_NOT_EXIST").isEqualTo(true)) .log("目标主题{}不存在,跳过路由执行", "my-topic") // 其他兜底操作 .end();
可选扩展
如果需要在路由启动阶段就完成校验(而非等第一条消息进来才判断),可以自定义RoutePolicySupport实现类,重写onStart方法,在路由启动时就调用AdminClient检查主题是否存在,不存在直接阻止路由启动,适合主题固定、不允许动态跳过的生产场景。
内容的提问来源于stack exchange,提问作者mhm3091

