【发布时间】:2018-05-23 00:07:16
【问题描述】:
正在尝试使用 kafka 流代码进行某些操作,并希望在拆分数据后添加延迟或类似 threads.sleep() 1ms 之类的东西....我很困惑如何做到这一点..有人可以帮我做吗?
KStreamBuilder builder = new KStreamBuilder();
KStream<String, String> textlines = builder.stream("INTOPIC");
KStream<String, String> mstream = textlines
.mapValues(value -> value.replace("[",""));
.mapValues(value -> value.replace("]",""));
.mapValues(value -> value.replaceAll("\\},\\{" ,"\\}\\},\\{\\{"))
.flatMapValues(value -> Arrays.asList(value.split("\\},\\{")));
mstream.to("OUTTOPIC");
KafkaStreams streams = new KafkaStreams(builder, config);
streams.start();
所以在 .flatmapvalues 语句之后我需要添加一个 thread.sleep() 1ms 那么我的语句可以在那里..?
【问题讨论】:
-
你能详细说明你为什么需要睡觉吗? (这真的是 XY 问题吗?)
标签: apache-kafka apache-kafka-streams