日日操夜夜添-日日操影院-日日草夜夜操-日日干干-精品一区二区三区波多野结衣-精品一区二区三区高清免费不卡

公告:魔扣目錄網為廣大站長提供免費收錄網站服務,提交前請做好本站友鏈:【 網站目錄:http://www.ylptlb.cn 】, 免友鏈快審服務(50元/站),

點擊這里在線咨詢客服
新站提交
  • 網站:51998
  • 待審:31
  • 小程序:12
  • 文章:1030137
  • 會員:747

本文介紹了ReplyingKafkaTemplate/KafkaTemplate不是發送/接收密鑰的處理方法,對大家解決問題具有一定的參考價值,需要的朋友們下面隨著小編來一起學習吧!

問題描述

我有一個模板,如下所示:

@Autowired
private ReplyingKafkaTemplate<ItemId, MessageBDto, MessageBDto> xxx2ReplyingKafkaTemplate;

我的發送包裝方法如下所示:

public RequestReplyFuture<ItemId, MessageBDto, MessageBDto> sendAndReceiveMessageB(MessageBDto message) {
    ProducerRecord<ItemId, MessageBDto> 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<ItemId, MessageBDto> 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=MessageBD發送到偵聽器并在偵聽器中接收密鑰嗎?

偵聽器似乎未獲得密鑰和/或它似乎是MessageBDto的另一個實例。

我是否誤解了這應該如何工作?

編輯:

生產者Bean:

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

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

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

@Bean
public ReplyingKafkaTemplate<ItemId, MessageBDto, MessageBDto> xxx2ReplyingKafkaTemplate() {
    return new ReplyingKafkaTemplate<>(xxx2ProducerFactory(), xxx2MessageListenerContainer());
}

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

消費類豆類:

@Bean
public KafkaTemplate<ItemId, MessageBDto> xxx2KafkaTemplate() {
    return new KafkaTemplate<>(xxx2ProducerFactory());
}

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

當我查看監聽程序的調試器時,它顯示key是MessageBDto?

的空實例

版本:

ApacheZooKeeper 3.5.7
阿帕奇·卡夫卡2.12-2.4.0

Spring Boot 2.2.4

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
    <version>2.4.2.RELEASE</version> <!-- $NO-MVN-MAN-VER$ -->
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>2.4.0</version> <!--$NO-MVN-MAN-VER$-->
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.4.0</version>  <!--$NO-MVN-MAN-VER$-->
</dependency>

推薦答案

我不確定您的代碼出了什么問題。以下是一個工作示例…

@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<Foo, Bar> 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<Foo, Bar, Bar> replyer(ProducerFactory<Foo, Bar> pf,
            ConcurrentKafkaListenerContainerFactory<Foo, Bar> containerFactory) {

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

    @Bean
    public KafkaTemplate<Foo, Bar> kafkaTemplate(ProducerFactory<Foo, Bar> pf) {
        return new KafkaTemplate<>(pf);
    }

    @Bean
    public ApplicationRunner runner(ReplyingKafkaTemplate<Foo, Bar, Bar> template) {
        return args -> {
            RequestReplyFuture<Foo, Bar, Bar> future =
                    template.sendAndReceive(new ProducerRecord<>("so60384112", 0, new Foo("foo"), new Bar("bar")));
            ConsumerRecord<Foo, Bar> 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信息2020-[o60384112-0-C-1]com.example.demo.So60384112應用程序:FOO[FIELD=FOO]:BAR[FIELD=BAR]

2020-02-24 18:00:47.915信息16591-[Main]com.example.demo.So60384112應用程序:foo[field=foo]:bar[field=bar]

要返回密鑰,您必須返回Message<?>。遺憾的是,您還必須為回復主題和相關性設置標頭。

這篇關于ReplyingKafkaTemplate/KafkaTemplate不是發送/接收密鑰的文章就介紹到這了,希望我們推薦的答案對大家有所幫助,

分享到:
標簽:KafkaTemplate ReplyingKafkaTemplate 發送 密鑰 接收
用戶無頭像

網友整理

注冊時間:

網站:5 個   小程序:0 個  文章:12 篇

  • 51998

    網站

  • 12

    小程序

  • 1030137

    文章

  • 747

    會員

趕快注冊賬號,推廣您的網站吧!
最新入駐小程序

數獨大挑戰2018-06-03

數獨一種數學游戲,玩家需要根據9

答題星2018-06-03

您可以通過答題星輕松地創建試卷

全階人生考試2018-06-03

各種考試題,題庫,初中,高中,大學四六

運動步數有氧達人2018-06-03

記錄運動步數,積累氧氣值。還可偷

每日養生app2018-06-03

每日養生,天天健康

體育訓練成績評定2018-06-03

通用課目體育訓練成績評定