【发布时间】:2019-06-20 00:49:34
【问题描述】:
我正在研究 kafka-streams api 。基本上,Kafka-stream 从源主题获取数据,并在应用一些过滤器后将其写回目标 kafka 主题。
使用的依赖。 :
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>2.1.0</version>
</dependency>
下面是相同的代码。 :
{ ...
// create property
Properties property = new Properties();
property.setProperty(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,"127.0.0.1:9092");
property.setProperty(StreamsConfig.APPLICATION_ID_CONFIG,"kafka_streams_app");
property.setProperty(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class.getName());
property.setProperty(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class.getName());
//create topology
StreamsBuilder streamsBuilder = new StreamsBuilder();
//build topology
KStream<String,String> inputTopic = streamsBuilder.stream("source_topic");
//filtering data
KStream<String,String> filteredStream = inputTopic.filter(
(k,val)-> filterData(val)>10000
);
filteredStream.to("target_topic");
KafkaStreams streams = new KafkaStreams(
streamsBuilder.build(),
property
);
//start our stream app
streams.start();
...
}
这是我的应用架构:
生产者 API(在源主题中生产)=>
kafka-stream API(从源主题读取并将数据发送到目标主题) => kafka-consumer api(从目标主题读取)
我想要的是,当流将数据写入target topic 时,我想捕获事件是否成功。
有什么方法可以捕获该回调?谢谢
【问题讨论】:
-
@ValBonn 我正在尝试在流写入过滤主题时捕获事件。现在就像打印
data is written to topic作为成功消息一样。 -
我添加了一个答案——不可能。我仍然感兴趣,为什么你需要这个?您尝试使用回调实现什么目标?
标签: apache-kafka apache-kafka-streams