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

如何在Kafka Streams中处理并忽略UNKNOWN_TOPIC_OR_PARTITION错误

解决Kafka Streams发往已删除主题的无限错误循环问题

我正在开发基于Kafka Streams的应用,该应用根据消息头动态确定目标主题。场景中应用运行期间主题被删除是正常情况,偶尔会收到发往已删除主题的消息,希望直接忽略这些消息。但仅收到一条不存在主题的消息后,就会触发无限错误循环:

[kafka-producer-network-thread | stream-example-producer] WARN org.apache.kafka.clients.NetworkClient -- [Producer clientId=stream-example-producer] Error while fetching metadata with correlation id 74 : {test1=UNKNOWN_TOPIC_OR_PARTITION}

org.apache.kafka.common.errors.TimeoutException: Topic test1 not present in metadata after 60000 ms.

[kafka-producer-network-thread | stream-example-producer] WARN org.apache.kafka.clients.NetworkClient -- [Producer clientId=stream-example-producer] Error while fetching metadata with correlation id 79 : {test1=UNKNOWN_TOPIC_OR_PARTITION}

这种无限错误循环会导致应用无法正常工作,需要配置Kafka Streams应用,使其忽略发往已删除主题的消息且不触发无限错误循环。

应用简化代码示例

StreamsBuilder builder = new StreamsBuilder();
List<String> dynamicTopics = List.of("good_topic", "deleted_topic");
builder.stream("source_topic").to((k, v, c) -> dynamicTopics.get(new Random().nextInt(dynamicTopics.size()))); //实际场景从消息头获取
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

自动主题创建已禁用。

已尝试的无效方法

  • 使用KafkaAdmin:检查主题存在性的间隙,主题可能被删除,无法解决问题。
  • 设置UncaughtExceptionHandler:代码未进入该处理器
    streams.setUncaughtExceptionHandler(new StreamsUncaughtExceptionHandler() {
        @Override
        public StreamThreadExceptionResponse handle(Throwable throwable) {
            return StreamThreadExceptionResponse.SHUTDOWN_APPLICATION;
        }
    });
    
  • 设置ProductionExceptionHandler:代码未进入该处理器
    props.put(StreamsConfig.DEFAULT_PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG,
              CustomProductionExceptionHandler.class.getName());
    
  • 设置Producer Interceptor:代码能进入拦截器,但无法解决问题
    props.put(StreamsConfig.producerPrefix(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG), ErrorInterceptor.class.getName());
    
  • 调整Producer属性:Kafka Streams仍会无限尝试处理错误
    props.put(StreamsConfig.RETRY_BACKOFF_MS_CONFIG, "5000"); 
    props.put(StreamsConfig.producerPrefix(ProducerConfig.MAX_BLOCK_MS_CONFIG), "8000");
    props.put(StreamsConfig.producerPrefix(ProducerConfig.LINGER_MS_CONFIG), "0");
    props.put(StreamsConfig.producerPrefix(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG), "10000");
    props.put(StreamsConfig.producerPrefix(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG), "10000");
    props.put(StreamsConfig.producerPrefix(ProducerConfig.RETRIES_CONFIG), 0);
    

可行解决方案

1. 结合前置过滤+自定义生产异常处理器

步骤1:前置过滤减少无效发送

在流处理中先过滤掉已知不存在的主题(可结合缓存的主题列表,定期刷新):

builder.stream("source_topic")
       .flatMap((k, v, ctx) -> {
           String targetTopic = ctx.headers().lastHeader("target_topic").value().toString();
           // 这里用缓存的有效主题列表做初步过滤
           if (validTopicsCache.contains(targetTopic)) {
               return Collections.singletonList(new KeyValue<>(k, v));
           } else {
               log.warn("Ignoring message for non-existent topic: {}", targetTopic);
               return Collections.emptyList();
           }
       })
       .to((k, v, ctx) -> ctx.headers().lastHeader("target_topic").value().toString());

步骤2:实现自定义ProductionExceptionHandler

捕获主题不存在的异常并返回CONTINUE,阻止无限重试:

public class IgnoreMissingTopicHandler implements ProductionExceptionHandler {
    private static final Logger log = LoggerFactory.getLogger(IgnoreMissingTopicHandler.class);

    @Override
    public ProductionExceptionHandlerResponse handle(ProducerRecord<Object, Object> record, Exception exception) {
        // 识别主题不存在相关的异常
        boolean isMissingTopic = exception instanceof UnknownTopicOrPartitionException 
                || (exception instanceof TimeoutException && exception.getCause() instanceof UnknownTopicOrPartitionException);
        
        if (isMissingTopic) {
            log.warn("Skip sending message to non-existent topic: {}", record.topic());
            return ProductionExceptionHandlerResponse.CONTINUE;
        }
        // 其他异常正常抛出
        return ProductionExceptionHandlerResponse.FAIL;
    }

    @Override
    public void configure(Map<String, ?> configs) {}
}

步骤3:配置到Kafka Streams

// 注册自定义异常处理器
props.put(StreamsConfig.DEFAULT_PRODUCTION_EXCEPTION_HANDLER_CLASS_CONFIG,
          IgnoreMissingTopicHandler.class.getName());
// 缩短元数据刷新间隔,加快感知主题删除
props.put(StreamsConfig.producerPrefix(ProducerConfig.METADATA_MAX_AGE_MS_CONFIG), "5000");
// 关闭重试,避免重复尝试
props.put(StreamsConfig.producerPrefix(ProducerConfig.RETRIES_CONFIG), 0);
props.put(StreamsConfig.producerPrefix(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG), "10000");

2. 手动控制Producer发送逻辑

通过transform操作手动处理发送,直接捕获异常并忽略:

builder.stream("source_topic")
       .transform(() -> new Transformer<Object, Object, Void>() {
           private Producer<Object, Object> producer;
           private Logger log;

           @Override
           public void init(ProcessorContext context) {
               log = LoggerFactory.getLogger(getClass());
               producer = context.getProducer();
           }

           @Override
           public Void transform(Object key, Object value) {
               String targetTopic = ...; // 从消息头获取目标主题
               try {
                   producer.send(new ProducerRecord<>(targetTopic, key, value), (metadata, exception) -> {
                       if (exception != null) {
                           if (exception instanceof UnknownTopicOrPartitionException || 
                               (exception instanceof TimeoutException && exception.getCause() instanceof UnknownTopicOrPartitionException)) {
                               log.warn("Ignoring message for non-existent topic: {}", targetTopic);
                           } else {
                               log.error("Failed to send message to topic {}", targetTopic, exception);
                           }
                       }
                   });
               } catch (Exception e) {
                   if (e instanceof UnknownTopicOrPartitionException) {
                       log.warn("Ignoring message for non-existent topic: {}", targetTopic);
                   } else {
                       throw e;
                   }
               }
               return null;
           }

           @Override
           public void close() {
               producer.close();
           }
       });

关键说明

  • 无限循环的根源是Kafka Producer会持续重试获取不存在主题的元数据,即使关闭RETRIES,元数据刷新机制仍会重复尝试。
  • 自定义异常处理器必须准确识别UnknownTopicOrPartitionException及其包装的TimeoutException,返回CONTINUE才能让应用继续运行。
  • 前置过滤能减少无效发送次数,但无法完全避免(主题可能在过滤后立即被删除),因此异常处理是核心解决手段。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 16:17:22