【问题标题】:delay function in kafka streamskafka流中的延迟函数
【发布时间】: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


【解决方案1】:

不确定您想要实现什么,但您似乎想减慢处理速度?比你可以在你的使用代码中加入睡眠。为此,您的 lambda 表达式必须在返回实际结果之前调用“sleep”。作为替代方案,您还可以添加一个额外的.foreach()peek() 呼叫并在那里睡觉。

【讨论】:

  • 究竟有什么用例可以减慢处理速度?如果这是一个依赖/竞争条件问题,我强烈建议在 NiFi 之类的下游有条件地进行处理。
  • 感谢您的回复...所以问题是我将每秒获取数组流 [{"json data"},{"json data"},{"json data"}]像这样......我想把这个分开。所以我写了上面的代码来分割数组中的记录......现在我需要这些记录一个接一个地出现,而不是一堆,所以我计划放慢这个过程,或者什么可能是最好的解决方案?这些记录应该从 confluent kafka 转到流反应器,其中流反应器将一个接一个地接受输入..它不能像我的代码一样接受数据生成
  • 如前所述。您可以在 flatMap 之后执行 foreach() 并在 foreach() 内睡觉
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-04-17
  • 1970-01-01
  • 2021-08-02
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多