【发布时间】:2017-12-15 01:17:40
【问题描述】:
我目前正在尝试轻松地将消息从一个 Kafka 集群上的主题流式传输到另一个(远程 -> 本地集群)。
我们的想法是立即使用 Kafka-Streams,这样我们就不需要在本地集群上复制实际消息,而只需将 Kafka-Streams 处理的“结果”获取到我们的 Kafka-Topics 中。
假设 WordCount 演示在我自己的另一台 PC 上的一个 Kafka-Instance 上。我还在本地机器上运行了一个 Kafka-Instance。
现在我想让 WordCount 演示在包含应计算单词的句子的主题(“远程”)上运行。
然而,计数应该写入我本地系统上的主题而不是“远程”主题。
使用 Kafka-Streams API 可以实现类似的操作吗?
例如。
val builder: KStreamBuilder = new KStreamBuilder(remote-streamConfig, local-streamconfig)
val textLines: KStream[String, String] = builder.stream("remote-input-topic",
remote-streamConfig)
val wordCounts: KTable[String, Long] = textLines
.flatMapValues(textLine => textLine.toLowerCase.split("\\W+").toIterable.asJava)
.groupBy((_, word) => word)
.count("word-counts")
wordCounts.to(stringSerde, longSerde, "local-output-topic", local-streamconfig)
val streams: KafkaStreams = new KafkaStreams(builder)
streams.start()
非常感谢
- 蒂姆
【问题讨论】:
-
查看Replicator,Matthias 在上面的回答中提到了这一点。这很符合你的描述。
-
MirrorMaker 是您所要求的开源选项cwiki.apache.org/confluence/display/KAFKA/…
标签: apache-kafka apache-kafka-streams