【发布时间】:2017-08-30 23:06:11
【问题描述】:
我开始尝试 Kafka Streams。我关注https://kafka.apache.org/0110/documentation/streams/quickstart。
我的沙盒是一个运行 Ubuntu 16.04.2 LTS、Kafka 0.11.0.0 和 Scala 2.11.11 的盒子。
如 Kafka Streams 快速入门指南中所述,以下是我遵循的步骤:
echo -e "all streams lead to kafka\nhello kafka streams\njoin kafka summit" > file-input.txt
bin/kafka-topics.sh --create \
--zookeeper localhost:2181 \
--replication-factor 1 \
--partitions 1 \
--topic streams-file-input
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic streams-file-input < file-input.txt
bin/kafka-run-class.sh org.apache.kafka.streams.examples.wordcount.WordCountDemo
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic streams-wordcount-output \
--from-beginning \
--formatter kafka.tools.DefaultMessageFormatter \
--property print.key=true \
--property print.value=true \
--property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \
--property value.deserializer=org.apache.kafka.common.serialization.LongDeserializer
当使用后一个命令查看streams-wordcount-output 时,我的标准输出显示如下:
all 1
streams 1
lead 1
to 1
kafka 1
hello 1
kafka 2
streams 2
join 1
kafka 3
summit 1
然后,在不中断 bin/kafka-console-consumer.sh 命令的情况下,我重新运行控制台生产者,如下所示:
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic streams-file-input < file-input.txt
令我惊讶的是,标准输出没有改变以反映这一新增功能所带来的变化。据我了解,file-input.txt 用于生成附加数据,因此字数应该已经刷新(现在所有标记都应计算两次)。 我的推理有什么问题?
【问题讨论】:
-
当然,
bin/kafka-run-class.sh org.apache.kafka.streams.examples.wordcount.WordCountDemo一直在运行?只是为了仔细检查,您还应该在streams-file-input主题上运行消费者,以确保您确实在那里添加新值... -
哦哦...我没有注意到 WordCountDemo 不再运行。再次运行它,输出看起来正确。谢谢 !但是, bin/kafka-run-class.sh org.apache.kafka.streams.examples.wordcount.WordCountDemo 在大约 5 秒后停止。据我了解,它应该永远运行。我错过了什么吗?
标签: apache-kafka apache-kafka-streams