【问题标题】:Brave Tracing for Kafka consumer using @KafkaListener使用 @KafkaListener 对 Kafka 消费者进行勇敢的跟踪
【发布时间】:2020-11-18 15:39:10
【问题描述】:

我正在使用 Brave 库 https://github.com/openzipkin/brave 进行跟踪,现在我也想将它用于 Kafka 消费者。我想避免添加 Spring Sleuth 并仅利用 Brave Kafka 工具https://github.com/openzipkin/brave/tree/master/instrumentation/kafka-clients

对于 Kafka 消费者,我使用 @KafkaListener。代码如下所示:

TestKafkaEndpoint.java

@Service
public class TestKafkaEndpoint {

    @KafkaListener(topics = "myTestTopic", containerFactory = "testKafkaListenerContainerFactory")
    public void procesMyRequest(@Payload final MyRequest request) {
       // do some magic...
    }
}

及配置类TestKafkaConfig.java


@Configuration
@EnableKafka
@ComponentScan
public class TestKafkaConfig {

    @Bean
    public ConsumerFactory<String, MyRequest> testConsumerFactory() {
        final Map<String, Object> consumerProperties = new HashMap<>();
        consumerProperties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka01-localhost:9092");
        consumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        consumerProperties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        consumerProperties.put(ConsumerConfig.GROUP_ID_CONFIG, "TestGROUP");
        return new DefaultKafkaConsumerFactory<>(consumerProperties, new StringDeserializer(), new JsonDeserializer<>(MyRequest.class));
    }

    @Bean
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, MyRequest>> testKafkaListenerContainerFactory() {
        final ConcurrentKafkaListenerContainerFactory<String, MyRequest> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(testConsumerFactory());
        factory.getContainerProperties().setErrorHandler(new LoggingErrorHandler());
        return factory;
    }

但我不知道在使用 Kafka 工厂时如何使用 KafkaConsumer 或利用 KafkaTracing。有没有人有这方面的经验并让它发挥作用?

【问题讨论】:

    标签: java spring apache-kafka spring-kafka zipkin


    【解决方案1】:

    我不熟悉它,但看起来 TracingConsumer 是一个简单的消费者包装器:https://github.com/openzipkin/brave/blob/363ceb4c922305ffb4a68ac47dc152e1d15da0fb/instrumentation/kafka-clients/src/main/java/brave/kafka/clients/TracingConsumer.java#L69-L79

    您应该能够创建DefaultKafkaConsumerFactory 的子类;覆盖 createConsumer 方法 - 侦听器容器使用...

    this.consumer =
            KafkaMessageListenerContainer.this.consumerFactory.createConsumer(
                    this.consumerGroupId,
                    this.containerProperties.getClientId(),
                    KafkaMessageListenerContainer.this.clientIdSuffix,
                    consumerProperties);
    

    ... 调用 super.createConsumer(...) 并将其包装在 TracingConsumer 中。

    如果您使用的是 2.5.3 或更高版本,您可以将ConsumerPostProcessor 添加到 DKCF。

    侦探就是这样做的:

    https://github.com/spring-cloud/spring-cloud-sleuth/blob/6e306e594d20361483fd19739e0f5f8e82354bf5/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/instrument/messaging/TraceMessagingAutoConfiguration.java#L263-L285

    【讨论】:

    • 嗨,加里,感谢您的精彩评论。我已经开始研究对 DefaultKafkaConsumerFactory 进行子类化,这可能是要走的路。我还会看看 ConsumerPostProcessor ,它听起来有点整洁的解决方案,并将在这里回复我的经验,最终解决方案:)
    • DefaultKafkaConsumerFactory 的子类提供了帮助。谢谢你:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-05-07
    • 2019-08-29
    • 2018-12-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多