Estoy usando la biblioteca Brave https://github.com/openzipkin/brave para rastrear y ahora me gustaría usarla también para el consumidor de Kafka. Me gustaría evitar agregar Spring Sleuth y aprovechar solo la instrumentación Brave Kafka https://github.com/openzipkin/brave/tree/master/instrumentation/kafka-clients .
Para el consumidor de Kafka, uso @KafkaListener . El código se ve así:
TestKafkaEndpoint.java
@Service public class TestKafkaEndpoint { @KafkaListener(topics = "myTestTopic", containerFactory = "testKafkaListenerContainerFactory") public void procesMyRequest(@Payload final MyRequest request) { // do some magic... } }Y clase de configuración 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; }Pero no sé cómo usar KafkaConsumer cuando uso Kafka Factory o aprovecho KafkaTracing . ¿Alguien tiene alguna experiencia con eso y lo hizo funcionar?
No estoy familiarizado con él, pero parece que TracingConsumer es un contenedor de consumidor simple: https://github.com/openzipkin/brave/blob/363ceb4c922305ffb4a68ac47dc152e1d15da0fb/instrumentation/kafka-clients/src/main/java/brave/ kafka/clientes/TracingConsumer.java#L69-L79
Debería poder crear una subclase de DefaultKafkaConsumerFactory ; anula los métodos createConsumer : el contenedor de escucha usa...
this.consumer = KafkaMessageListenerContainer.this.consumerFactory.createConsumer( this.consumerGroupId, this.containerProperties.getClientId(), KafkaMessageListenerContainer.this.clientIdSuffix, consumerProperties); ... llame a super.createConsumer(...) y envuélvalo en un TracingConsumer .
Si está utilizando 2.5.3 o posterior, puede agregar un ConsumerPostProcessor al DKCF.
Así es como lo hace el detective: