ReplyingKafkaTemplate / KafkaTemplate不发送/接收密钥是什么?

来源:爱站网时间:2021-11-22编辑:网友分享
有许多的编程小白是不懂ReplyingKafkaTemplate / KafkaTemplate不发送/接收密钥是什么?不知道的朋友没关系,爱站技术小编用一篇文章给你们参考一下,希望这个能帮助到你们哦,大家还可以做好笔记哦。

问题描述


我有一个看起来像这样的模板:

@Autowired
private ReplyingKafkaTemplate xxx2ReplyingKafkaTemplate;

我的发送包装器方法如下:

public RequestReplyFuture sendAndReceiveMessageB(MessageBDto message) {
    ProducerRecord producerRecord = new ProducerRecord(KafkaTopicConfig.xxx2_TOPIC, new ItemId(message.getCount()), message);
    producerRecord.headers().add(new RecordHeader(KafkaHeaders.REPLY_TOPIC, KafkaTopicConfig.xxx2_REPLY_TOPIC.getBytes()));
    return this.xxx2ReplyingKafkaTemplate.sendAndReceive(producerRecord);
}

这是我的听众:

@SendTo
@KafkaListener(topics=KafkaTopicConfig.xxx2_TOPIC, containerFactory="xxx2ListenerContainerFactory")
public MessageBDto xxx2Listener(ConsumerRecord message) {
    System.out.println("xxx2(value): " + message.value().getMessage() + ", " + message.value().getCount());
    message.value().setCount(message.value().getCount() * 2);
    return message.value();
}

这不是应该发送Key = ItemId,Value = MessageBDto并在侦听器中接收密钥吗?

侦听器似乎没有获取密钥和/或它似乎是MessageBDto的另一个实例。

我误解了它应该如何工作?

编辑:

生产者豆:

@Bean
public ProducerFactory xxx2ProducerFactory() {
    return new DefaultKafkaProducerFactory(super.producerConfigs(),
                                                                new JsonSerializer(),
                                                                new JsonSerializer());
}

@Bean
public ConsumerFactory xxx2ConsumerFactory() {
    return new DefaultKafkaConsumerFactory(super.consumerConfigs(),
                                             trustingDeserializer(ItemId.class),
                                             trustingDeserializer(MessageBDto.class));
}

@Bean
public KafkaMessageListenerContainer dtms2MessageListenerContainer() {
    return new KafkaMessageListenerContainer(xxx2ConsumerFactory(),
                                               new ContainerProperties(KafkaTopicConfig.xxx2_REPLY_TOPIC));
}

@Bean
public ReplyingKafkaTemplate xxx2ReplyingKafkaTemplate() {
    return new ReplyingKafkaTemplate(xxx2ProducerFactory(), xxx2MessageListenerContainer());
}

private  JsonDeserializer trustingDeserializer(Class targetType) {
    JsonDeserializer deserializer = new JsonDeserializer(targetType);
    deserializer.addTrustedPackages("*");
    return deserializer;
}

消费者豆:

@Bean
public KafkaTemplate xxx2KafkaTemplate() {
    return new KafkaTemplate(xxx2ProducerFactory());
}

@Bean
public KafkaListenerContainerFactory xxxListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory();
    factory.setConsumerFactory(dtms2ConsumerFactory());
    factory.setReplyTemplate(xxx2KafkaTemplate());
    return factory;
}

[当我在侦听器中查看侦听器时,表明该键是MessageBDto的空实例?

版本:

  • Apache Zookeeper 3.5.7
  • Apache Kafka 2.12-2.4.0
  • Spring Boot 2.2.4

    org.springframework.kafkaspring-kafka2.4.2.RELEASEorg.apache.kafkakafka-streams2.4.0org.apache.kafkakafka-clients2.4.0

思路:


我不确定您的代码是怎么回事。这是一个工作示例...

@SpringBootApplication
public class So60384112Application {


    private static final Logger LOG = LoggerFactory.getLogger(So60384112Application.class);


    public static void main(String[] args) {
        SpringApplication.run(So60384112Application.class, args).close();
    }

    @Bean
    public NewTopic topic() {
        return TopicBuilder.name("so60384112").partitions(1).replicas(1).build();
    }

    @Bean
    public NewTopic replies() {
        return TopicBuilder.name("so60384112replies").partitions(1).replicas(1).build();
    }

    @KafkaListener(id = "so60384112", topics = "so60384112")
    @SendTo
    public Message> listen(ConsumerRecord record) {
        LOG.info(record.key().toString() + ":" + record.value().toString());
        return MessageBuilder.withPayload(new Bar(record.value().getField().toUpperCase()))
                .setHeader(KafkaHeaders.MESSAGE_KEY, record.key())
                .setHeader(KafkaHeaders.CORRELATION_ID, record.headers().lastHeader(KafkaHeaders.CORRELATION_ID).value())
                .setHeader(KafkaHeaders.TOPIC, record.headers().lastHeader(KafkaHeaders.REPLY_TOPIC).value())
                .build();
    }

    @Bean
    public ReplyingKafkaTemplate replyer(ProducerFactory pf,
            ConcurrentKafkaListenerContainerFactory containerFactory) {

        containerFactory.setReplyTemplate(kafkaTemplate(pf));
        ConcurrentMessageListenerContainer container = containerFactory.createContainer("so60384112replies");
        container.getContainerProperties().setGroupId("so60384112replies");
        ReplyingKafkaTemplate replyer = new ReplyingKafkaTemplate(pf, container);
        return replyer;
    }

    @Bean
    public KafkaTemplate kafkaTemplate(ProducerFactory pf) {
        return new KafkaTemplate(pf);
    }

    @Bean
    public ApplicationRunner runner(ReplyingKafkaTemplate template) {
        return args -> {
            RequestReplyFuture future =
                    template.sendAndReceive(new ProducerRecord("so60384112", 0, new Foo("foo"), new Bar("bar")));
            ConsumerRecord record = future.get(10, TimeUnit.SECONDS);
            LOG.info(record.key().toString() + ":" + record.value().toString());
        };
    }

}

class Foo {

    private String field;

    public Foo() {
    }

    public Foo(String field) {
        this.field = field;
    }

    public String getField() {
        return this.field;
    }

    public void setField(String field) {
        this.field = field;
    }

    @Override
    public String toString() {
        return getClass().getSimpleName() + " [field=" + this.field + "]";
    }

}

class Bar  {

    private String field;

    public Bar() {
    }

    public Bar(String field) {
        this.field = field;
    }

    public String getField() {
        return this.field;
    }

    public void setField(String field) {
        this.field = field;
    }

    @Override
    public String toString() {
        return getClass().getSimpleName() + " [field=" + this.field + "]";
    }

}
spring.kafka.consumer.key-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.properties.spring.json.trusted.packages=*

spring.kafka.producer.key-serializer=org.springframework.kafka.support.serializer.JsonSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer

2020-02-24 18:00:47.904信息16591 --- [o60384112-0-C-1] com.example.demo.So60384112Application:Foo [field = foo]:Bar [field = bar]

2020-02-24 18:00:47.915信息16591 --- [main] com.example.demo.So60384112Application:Foo [field = foo]:Bar [field = BAR]

要返回键,您必须返回Message>。不幸的是,您也必须为回复主题和相关性设置标题。

以上内容就是爱站技术频道小编为大家分享的ReplyingKafkaTemplate / KafkaTemplate不发送/接收密钥是什么?看完以上分享之后,大家应该都知道是什么了吧。

上一篇:字符串数组在更新时删除旧值要怎么办

下一篇:怎么获取幸运数字?

您可能感兴趣的文章

相关阅读

热门软件源码

最新软件源码下载