【问题标题】:Adding custom header using Spring Kafka使用 Spring Kafka 添加自定义标头
【发布时间】:2018-07-25 16:12:42
【问题描述】:

我计划使用 Spring Kafka 客户端在 Spring Boot 应用程序中使用 kafka 设置来使用和生成消息。我在 Kafka 0.11 中看到了对自定义标头的支持,详见 here。虽然它适用于本地 Kafka 生产者和消费者,但我看不到在 Spring Kafka 中添加/读取自定义标头的支持。

我正在尝试根据我希望存储在消息标头中的重试计数为消息实现 DLQ,而无需解析有效负载。

【问题讨论】:

    标签: spring apache-kafka spring-kafka spring-boot-test


    【解决方案1】:

    当我偶然发现这个问题时,我正在寻找答案。但是我使用的是ProducerRecord<?, ?> 类而不是Message<?>,因此标题映射器似乎不相关。

    这是我添加自定义标题的方法:

    var record = new ProducerRecord<String, String>(topicName, "Hello World");
    record.headers().add("foo", "bar".getBytes());
    kafkaTemplate.send(record);
    

    现在要读取标题(在使用之前),我添加了一个自定义拦截器。

    import java.util.List;
    import lombok.extern.slf4j.Slf4j;
    import org.apache.kafka.clients.consumer.ConsumerInterceptor;
    import org.apache.kafka.clients.consumer.ConsumerRecord;
    import org.apache.kafka.clients.consumer.ConsumerRecords;
    
    @Slf4j
    public class MyConsumerInterceptor implements ConsumerInterceptor<Object, Object> {
    
        @Override
        public ConsumerRecords<Object, Object> onConsume(ConsumerRecords<Object, Object> records) {
            Set<TopicPartition> partitions = records.partitions();
            partitions.forEach(partition -> interceptRecordsFromPartition(records.records(partition)));
    
            return records;
        }
    
        private void interceptRecordsFromPartition(List<ConsumerRecord<Object, Object>> records) {
            records.forEach(record -> {
                var myHeaders = new ArrayList<Header>();
                record.headers().headers("MyHeader").forEach(myHeaders::add);
                log.info("My Headers: {}", myHeaders);
                // Do with header as you see fit
            });
        }
    
        @Override public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) {}
        @Override public void close() {}
        @Override public void configure(Map<String, ?> configs) {}
    }
    

    最后一点是使用以下(Spring Boot)配置向 Kafka Consumer Container 注册此拦截器:

    import java.util.Map;
    import org.apache.kafka.clients.consumer.ConsumerConfig;
    import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    import org.springframework.kafka.core.ConsumerFactory;
    import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
    
    @Configuration
    public class MessagingConfiguration {
    
        @Bean
        public ConsumerFactory<?, ?> kafkaConsumerFactory(KafkaProperties properties) {
            Map<String, Object> consumerProperties = properties.buildConsumerProperties();
            consumerProperties.put(ConsumerConfig.INTERCEPTOR_CLASSES_CONFIG, MyConsumerInterceptor.class.getName());
            return new DefaultKafkaConsumerFactory<>(consumerProperties);
        }
    
    }
    

    【讨论】:

    • 如何将自定义标头发送到该批邮件? @BitfullByte
    • 我不确定@AngadBansode,为每条消息单独设置是否足够?如果没有,我建议打开单独的问题。
    • 谢谢,其实我有@Payload List 消息,想为每条消息添加自定义标头并使用它们吗?
    • 这看起来像是一条消息,给我一个字符串列表
    【解决方案2】:

    嗯,Spring Kafka 从 2.0 版 开始提供标头支持:https://docs.spring.io/spring-kafka/docs/2.1.2.RELEASE/reference/html/_reference.html#headers

    您可以拥有该KafkaHeaderMapper 实例并使用它来填充标题到Message,然后再通过KafkaTemplate.send(Message&lt;?&gt; message) 发送它。或者你可以使用普通的KafkaTemplate.send(ProducerRecord&lt;K, V&gt; record)

    当您使用KafkaMessageListenerContainer 接收记录时,可以通过注入RecordMessagingMessageListenerAdapterMessagingMessageConverter 提供KafkaHeaderMapper

    因此,任何自定义标头都可以通过任一方式传输。

    【讨论】:

      猜你喜欢
      • 2018-08-02
      • 1970-01-01
      • 1970-01-01
      • 2012-04-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-10-09
      • 1970-01-01
      相关资源
      最近更新 更多