集成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
相关产品推荐
相关产品推荐

