【发布时间】:2018-05-31 14:49:58
【问题描述】:
我需要知道如何在我的 kafka KStreams 行中使用“for”循环...下面是我的“for”循环,它需要包含在 KStreams 中
for (int i = 0; i < 6 ; i++) {
try {
textlines.flatMapValues(value -> Arrays.asList(value.split("\\},\\{")));
Thread.sleep(2000);
}catch (InterruptedException e){
e.printStackTrace();
}
}
我的 KStreams 看起来像
KStream<String, String> textlines = builder.stream("intopic");
KStream<String, String> mstream = textlines
.mapValues(value -> value.replace("[","" ) )
如何将上面的“for”循环添加到我的 KStreams 中
【问题讨论】:
-
这个 for 循环的确切目的是什么? KStream 对象只是一种构建拓扑的方法,该拓扑随后将在其他线程中运行(在 .start() 调用之后)。在您的代码中,您只是在拓扑中添加了 6 倍相同的处理器,而睡眠部分对流执行没有影响,只会延迟拓扑构建。
-
@nbchn 好的...问题是我在'for'循环中使用了 value.split 来拆分我的数据..所以每当我的数据被拆分时,它应该休眠大约 10 毫秒.. .这是因为我需要我的数据一个接一个地出现,如果您需要更多详细信息,请告诉我
标签: apache-kafka apache-kafka-streams