【发布时间】:2018-02-18 07:39:23
【问题描述】:
在这个应用程序中,我们正在消费消息处理它,然后再次将其转发给其他消费者,所以任何人都可以分享一种可能的方式来做到这一点, 如果您也可以分享一些示例,那就太好了。
【问题讨论】:
在这个应用程序中,我们正在消费消息处理它,然后再次将其转发给其他消费者,所以任何人都可以分享一种可能的方式来做到这一点, 如果您也可以分享一些示例,那就太好了。
【问题讨论】:
Spring Kafka 项目有@SendTo 注释只是为了这个目的,这就是让你的消费者也产生消息read docs
或者,您可以使用单个消费者从firstTopic 接收消息,并在其中添加kafkaTemplate
这样,在处理后它将向secondTopic 发送一条消息。下面是一个简单的例子,它只包含一个生产者和一个消费者。
有关配置类的完整示例参考(仅适用于基本的生产者-消费者示例),请参阅here
制片人
package com.codenotfound.kafka.producer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
public class Sender {
private static final Logger LOGGER =
LoggerFactory.getLogger(Sender.class);
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void send(String topic, String payload) {
LOGGER.info("sending payload='{}' to topic='{}'", payload, topic);
kafkaTemplate.send("firstTopic", payload);
}
}
消费者
package com.codenotfound.kafka.consumer;
import java.util.concurrent.CountDownLatch;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
public class Receiver {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
private static final Logger LOGGER =
LoggerFactory.getLogger(Receiver.class);
@KafkaListener(topics = "firstTopic")
public void receive(String payload) {
LOGGER.info("received payload='{}'", payload);
//To do processing and get generate payload
String payload2 = someprocessingLogic(payload);
kafkaTemplate.send("secondTopic", payload);
}
}
【讨论】:
"topic" 作为主题名称,例如"firstTopic" 用于主题名称。