- android - 多次调用 OnPrimaryClipChangedListener
- android - 无法更新 RecyclerView 中的 TextView 字段
- android.database.CursorIndexOutOfBoundsException : Index 0 requested, 光标大小为 0
- android - 使用 AppCompat 时,我们是否需要明确指定其 UI 组件(Spinner、EditText)颜色
我有一个如下所示的模板:
@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=MessageBDto 并在监听器中接收 key 吗?
监听器似乎没有获取 key 和/或它似乎是 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 :
@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;
}
当我查看监听器的调试器时,它显示该键是 MessageBDto 的空实例???
版本:
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 INFO 16591 --- [o60384112-0-C-1] com.example.demo.So60384112Application : Foo [field=foo]:Bar [field=bar]
2020-02-24 18:00:47.915 INFO 16591 --- [ main] com.example.demo.So60384112Application : Foo [field=foo]:Bar [field=BAR]
要返回 key ,您必须返回 Message<?>
。不幸的是,您还必须为回复主题和相关性设置 header 。
关于java - ReplyingKafkaTemplate/KafkaTemplate 未发送/接收 key ,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/60384112/
我有一个存储结构向量的应用程序。这些结构保存有关系统上每个 GPU 的信息,如内存和 giga-flop/s。每个系统上有不同数量的 GPU。 我有一个程序可以同时在多台机器上运行,我需要收集这些数据
我很好奇 MPI 中缺少此功能: MPI_Isendrecv( ... ); 即,非阻塞发送和接收,谁能告诉我其省略背后的基本原理? 最佳答案 我的看法是 MPI_SENDRECV存在是为了方便那些想
当我用以下方法监听TCP或UDP套接字时 ssize_t recv(int sockfd, void *buf, size_t len, int flags); 或者 ssize_t recvfrom
SUM:如何在 azure 事件网格中推迟事件触发或事件接收? 我设计的系统需要对低频对象状态(创建、启动、检查长时间启动状态、结束)使用react。它看起来像是事件处理的候选者。我想用azure函数
我正在 MPI 中实现一个程序,其中主进程(等级 = 0)应该能够接收来自其他进程的请求,这些进程要求只有根才知道的变量值。如果我按等级 0 进行 MPI_Recv(...),我必须指定向根发送请求的
我正在学习DX12,并在此过程中学习“旧版Win32”。 我在退出主循环时遇到问题,这似乎与我没有收到WM_CLOSE消息有关。 在C++,Windows 10控制台应用程序中。 #include
SUM:如何在 azure 事件网格中推迟事件触发或事件接收? 我设计的系统需要对低频对象状态(创建、启动、检查长时间启动状态、结束)使用react。它看起来像是事件处理的候选者。我想用azure函数
我想编写方法来通过号码发送短信并使用编辑文本字段中的文本。发送消息后,我想收到一些声音或其他东西来提醒我收到短信。我怎样才能做到这一点?先感谢您,狼。 最佳答案 这个网站似乎对两者都有很好的描述:ht
所以我正在用 Java 编写一个程序,在 DatagramSocket 和 DatagramPacket 的帮助下发送和接收数据。问题是,在我发送数据/接收数据之间的某个时间 - 我发送数据的程序中的
我是 Android 编程新手,我正在用 Java 编写一个应用程序,该应用程序可以打开相机拍照并保存。我通过 Intents 做到了,但看不到 onActivityResult 正在运行。 我已经在
我有一个套接字服务器和一个套接字客户端。客户端只有一个套接字。我必须使用线程在客户端发送/接收数据。 static int sock = -1; static std::mutex mutex; vo
我正在尝试使用 c 中的套接字实现 TCP 服务器/客户端。我以这样的方式编写程序,即我们在客户端发送的任何内容都逐行显示在服务器中,直到键入退出。该程序可以运行,但数据最后一起显示在服务器中。有人可
我正在使用微 Controller 与 SIM808 模块通信,我想发送和接收 AT 命令。 现在的问题是,对于某些命令,我只收到了我应该收到的答案的一部分,但对于其他一些命令,我收到了我应该
我用c设计了一个消息传递接口(interface),用于在我的系统中运行的不同进程之间提供通信。该接口(interface)为此目的创建 10-12 个线程,并使用 TCP 套接字提供通信。 它工作正
我需要澄清一下在套接字程序中使用多个发送/接收。我的客户端程序如下所示(使用 TCP SOCK_STREAM)。 send(sockfd,"Messgfromlient",15,0);
我正在构建一个真正的基本代理服务器到我现有的HTTP服务器中。将传入连接添加到队列中,并将信号发送到另一个等待线程队列中的一个线程。此线程从队列中获取传入连接并对其进行处理。 问题是代理程序真的很慢。
我正在使用 $routeProvider 设置一条类似 的路线 when('/grab/:param1/:param2', { controller: 'someController',
我在欧洲有通过 HLS 流式传输的商业流媒体服务器。http://europe.server/stream1/index.m3u8现在我在美国的客户由于距离而遇到一些网络问题。 所以我在美国部署了新服
我有一个长期运行的 celery 任务,该任务遍历一系列项目并执行一些操作。 任务应该以某种方式报告当前正在处理的项目,以便最终用户知道任务的进度。 目前,我的django应用程序和celery一起坐
我需要将音频文件从浏览器发送到 python Controller 。我是这样做的: var xmlHttp = new XMLHttpRequest(); xmlHttp.open( "POST",
我是一名优秀的程序员,十分优秀!