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

Kafka批量监听:仅转发失败消息至DLQ及参数异常解决

Kafka批量监听器:仅转发失败消息至DLQ的解决方案

问题场景

现有Spring Kafka批量监听器,需在消息转换失败(ConversionException)时,仅将触发异常的消息发送至DLQ,其余正常消息由监听器处理。当前实现存在两个问题:

  • 使用DefaultErrorHandler时,批量中只要有一条消息失败,整批消息都会被错误处理器处理;
  • 尝试添加@Header(KafkaHeaders.CONVERSION_FAILURES)参数获取转换失败信息时,触发IllegalStateException,提示List<Message<?>>必须是唯一参数(除特定可选参数外)。

解决方案

1. 修复监听器参数异常

当批量监听器使用List<Message<CancelAuthorizationLinkageResource>>作为参数时,Spring Kafka不允许添加额外的@Header参数。需调整参数结构,改用@Payload标注业务数据列表,同时添加转换失败异常的Header参数:

@KafkaListener(
        id = "${spring.kafka.listener.cancel-auth-linkage.id}",
        topics = "${spring.kafka.listener.cancel-auth-linkage.topic.linkage}",
        autoStartup = "false",
        batch = "true",
        groupId = "cushion")
public void listen(@Payload List<CancelAuthorizationLinkageResource> resources,
                   @Header(value = KafkaHeaders.CONVERSION_FAILURES, required = false) List<ConversionException> conversionExceptions,
                   @Header(KafkaHeaders.RECEIVED_RECORDS) List<ConsumerRecord<?, ?>> rawRecords) {
    // 处理正常消息
    if (resources != null && !resources.isEmpty()) {
        processor.process(resources);
    }

    // 处理转换失败的消息,手动发送至DLQ
    if (conversionExceptions != null && !conversionExceptions.isEmpty()) {
        for (ConversionException ex : conversionExceptions) {
            ConsumerRecord<?, ?> failedRecord = ex.getRecord();
            kafkaTemplate.send(listenerProperties.getDlqTopic(), failedRecord.key(), failedRecord.value());
        }
    }
}

2. 配置错误处理器实现单条失败消息转发

若希望通过DefaultErrorHandler和DeadLetterPublishingRecoverer自动处理失败消息,需确保错误处理器支持批量消息的逐记录处理。Spring Kafka 3.x及以上版本的DefaultErrorHandler默认支持批量场景下的单条错误处理,只需调整恢复器逻辑,并配置容器工厂启用批量错误处理:

调整DeadLetterPublishingRecoverer

@Bean
public DeadLetterPublishingRecoverer recoverer(KafkaTemplate<Object, Object> kafkaTemplate,
                                               ListenerPropertiesServiceInterface listenerPropertiesServiceInterface) {
    return new DeadLetterPublishingRecoverer(kafkaTemplate,
            (record, exception) -> {
                if (exception.getCause() instanceof ConversionException) {
                    // 仅将转换失败的消息转发至DLQ
                    return new TopicPartition(listenerPropertiesServiceInterface.getDlqTopic(), -1);
                } else {
                    throw new RuntimeException(exception);
                }
            });
}

配置DefaultErrorHandler支持批量处理

@Bean
public DefaultErrorHandler errorHandler(DeadLetterPublishingRecoverer recoverer) {
    DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer);
    // 设置仅针对ConversionException进行处理,其他异常抛出
    errorHandler.addNotRetryableExceptions(ConversionException.class);
    // 启用批量模式下的逐记录处理
    errorHandler.setProcessBatchWhileStopping(false);
    return errorHandler;
}

配置KafkaListenerContainerFactory

确保容器工厂关联错误处理器,并配置批量消费相关属性:

@Bean
public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
        ConsumerFactory<Object, Object> consumerFactory,
        DefaultErrorHandler errorHandler) {
    ConcurrentKafkaListenerContainerFactory<Object, Object> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    factory.setBatchListener(true);
    factory.setCommonErrorHandler(errorHandler);
    // 配置批量拉取大小
    factory.getContainerProperties().setPollTimeout(3000);
    return factory;
}

关键说明

  • 使用@Payload List<...>替代List<Message<...>>作为批量监听器参数,即可合法添加@Header参数获取转换失败信息;
  • DefaultErrorHandler在批量模式下默认会逐记录检查异常,仅处理失败的记录,正常记录会被监听器处理后提交偏移量;
  • 若手动处理转换失败消息,需通过KafkaHeaders.RECEIVED_RECORDS获取原始记录,避免丢失消息元数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 03:00:54