我正在尝试使用Kafka Streams 0.10.1在Scala中创建一个简单的聚合示例,尽管我似乎失败了一个简单的"计数"聚合(使用Kafka控制台生成器).有了这样的代码:
val inputStream: KStream[String, String] = builder.stream("inputTopic")
inputStream
.map(new KeyValueMapper[String, String, KeyValue[String, String]] {
override def apply(k: String, v: String): KeyValue[String, String] = {
new KeyValue[String, String](v, v)
}
})
.groupByKey()
.count(TimeWindows.of(10000L), "count-test-1")
.toStream()
.to("outputTopic")
Run Code Online (Sandbox Code Playgroud)
它失败了"无法刷新状态存储计数 - 测试-1",我在帖子的末尾包含了完整的堆栈跟踪.另一方面,如果我使用print()而不是(),它就像魅力一样,将结果打印到控制台/终端:
[KTABLE-TOSTREAM-0000000013]: [aa@1483089460000] , 1
[KTABLE-TOSTREAM-0000000013]: [bb@1483089460000] , 1
[KTABLE-TOSTREAM-0000000013]: [cc@1483089460000] , 2
[KTABLE-TOSTREAM-0000000013]: [dd@1483089460000] , 3
[KTABLE-TOSTREAM-0000000013]: [ee@1483089460000] , 4
Run Code Online (Sandbox Code Playgroud)
有谁知道这种行为可能是什么原因?
仅供参考,我使用的操作系统是Windows 10作为主机(也通过IntelliJ运行Scala应用程序)和用于Kafka(在Docker容器中)和生产者/消费者应用程序的Ubuntu 16.04 VM.但是,我可以确认在Ubuntu VM上运行应用程序时也会遇到问题.
非常感谢您的帮助,感谢任何见解:-)
完整的堆栈跟踪:
2016-12-30 08:57:43 INFO StreamThread:573 - stream-thread [StreamThread-1] Committing task 2_0
2016-12-30 08:57:43 …Run Code Online (Sandbox Code Playgroud)