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

如何使用Apache-Camel检查Kafka主题是否可用 缺失时跳过路由

Kafka-Camel路由前置主题校验实现方案

问题说明

需要在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 01:42:17