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

集成RetryableTopic与OpenTracing SpringBoot时遇UnsupportedOperationException

问题:Spring Kafka RetryableTopic 与 OpenTracing 集成异常

单独使用RetryableTopic配置时,重试和死信队列(DLT)功能正常,代码如下:

@RetryableTopic(
    attempts = "3",
    backoff = @Backoff(delay = 3000, multiplier = 2.0),
    topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_DELAY_VALUE,
    kafkaTemplate = "dltKafkaTemplate",
    listenerContainerFactory = "retryEventListenerFactory",
    exclude = {
        DeserializationException.class,
        SerializationException.class,
        MessageConversionException.class,
        ConversionException.class,
        MethodArgumentResolutionException.class,
        NoSuchMethodException.class,
        ClassCastException.class
    }
)
@KafkaListener(
    topics = "dlt-msg-test-topic",
    containerFactory = "retryEventListenerFactory")
public void consume(
    String message, @Headers MessageHeaders messageHeaders) {
  LOGGER.info("Received {}", message);
  throw new RuntimeException("Test retry exception");
}

@DltHandler
public void dlt(String in, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) {
  LOGGER.info(in + " from " + topic);
}

引入opentracing-spring-cloud-starter等追踪依赖后,抛出如下异常:

java.lang.UnsupportedOperationException: This implementation doesn't support this method
    at org.springframework.kafka.core.ProducerFactory.getConfigurationProperties(ProducerFactory.java:120)
    at org.springframework.kafka.listener.DeadLetterPublishingRecoverer.determineSendTimeout(DeadLetterPublishingRecoverer.java:661)
    at org.springframework.kafka.listener.DeadLetterPublishingRecoverer.verifySendResult(DeadLetterPublishingRecoverer.java:636)
    at org.springframework.kafka.listener.DeadLetterPublishingRecoverer.publish(DeadLetterPublishingRecoverer.java:628)
    at org.springframework.kafka.listener.DeadLetterPublishingRecoverer.send(DeadLetterPublishingRecoverer.java:524)
    at org.springframework.kafka.listener.DeadLetterPublishingRecoverer.sendOrThrow(DeadLetterPublishingRecoverer.java:489)
    at org.springframework.kafka.listener.DeadLetterPublishingRecoverer.accept(DeadLetterPublishingRecoverer.java:461)
    at org.springframework.kafka.listener.FailedRecordProcessor.getRecoveryStrategy(FailedRecordProcessor.java:181)
    at org.springframework.kafka.listener.DefaultErrorHandler.handleRemaining(DefaultErrorHandler.java:134)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeErrorHandler(KafkaMessageListenerContainer.java:2674)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:2555)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeWithRecords(KafkaMessageListenerContainer.java:2429)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:2307)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:1981)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1365)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1356)
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1251)

问题分析

OpenTracing会创建TracingKafkaConsumer包装原始消费者,重试流程中DefaultErrorHandler.handleRemaining方法传入的是这个包装类,导致SeekUtils.seekOrRecover无法从包装类中获取必要参数,最终触发ProducerFactory.getConfigurationProperties()的未实现异常。


解决方案

1. 自定义ProducerFactory包装类,实现缺失的方法

异常根源是OpenTracing包装的ProducerFactory没有实现getConfigurationProperties方法,我们可以自己写一个包装类,委托原始ProducerFactory的实现:

@Component
public class TracingProducerFactoryWrapper<K, V> implements ProducerFactory<K, V> {

    private final ProducerFactory<K, V> delegate;

    public TracingProducerFactoryWrapper(ProducerFactory<K, V> delegate) {
        this.delegate = delegate;
    }

    @Override
    public Map<String, Object> getConfigurationProperties() {
        // 直接返回原始ProducerFactory的配置属性
        return delegate.getConfigurationProperties();
    }

    // 其他所有方法都委托给原始ProducerFactory
    @Override
    public Producer<K, V> createProducer() {
        return delegate.createProducer();
    }

    @Override
    public Producer<K, V> createProducer(String transactionIdPrefix) {
        return delegate.createProducer(transactionIdPrefix);
    }

    @Override
    public boolean isTransactional() {
        return delegate.isTransactional();
    }

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

    @Override
    public void initTransactions() {
        delegate.initTransactions();
    }

    @Override
    public void closeTransaction() {
        delegate.closeTransaction();
    }
}

2. 替换OpenTracing默认的ProducerFactory

在配置类中,将原始ProducerFactory用自定义包装类包裹后注入,覆盖OpenTracing的默认配置:

@Configuration
public class KafkaTracingConfig {

    @Bean
    public ProducerFactory<Object, Object> producerFactory(ProducerFactory<Object, Object> originalProducerFactory) {
        return new TracingProducerFactoryWrapper<>(originalProducerFactory);
    }
}

3. 确保RetryableTopic使用正确的KafkaTemplate

确认@RetryableTopic中指定的kafkaTemplate依赖的是上述自定义包装后的ProducerFactory,或者直接让KafkaTemplate的配置指向这个Bean。

另外,如果可以升级opentracing-spring-cloud-kafka-starter的版本,部分新版本已经修复了这个未实现方法的问题,升级也是一个可选方案。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 10:15:47