如何在Kafka Streams中处理并忽略UNKNOWN_TOPIC_OR_PARTITION错误
我正在开发基于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

