【问题标题】:'for' loop support in Kafka KStreamsKafka Streams 中的“for”循环支持
【发布时间】: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


【解决方案1】:

问题是我在“for”循环中使用 value.split 来拆分我的数据....所以每当我的数据被拆分时,它应该休眠大约 10 毫秒...这是因为我需要我的数据来一个在另一个之后

根据您所说的要订购。要实现订购,您不需要sleep。它会起作用的。我假设您的代码所基于的 Kafka Streams WordCount 示例的工作方式相同:它也使用 flatMapValues,并且传递到平面地图的 lambda 将文本行拆分为单词。

除非我和其他人误解了您的问题(在这种情况下,您或许应该进一步澄清您的问题),否则我认为您不必要地使您的代码复杂化。

【讨论】:

  • link 你可以看看这个例子……这正是我需要做的……
  • 答案是一样的。只需使用flatMapValues,这将确保数据一个接一个。
猜你喜欢
  • 2013-08-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-06-12
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多