【问题标题】:How to send record on topic when window is closed in kafka streams在kafka流中关闭窗口时如何发送有关主题的记录
【发布时间】:2018-02-28 10:51:35
【问题描述】:

实际上,我已经为此苦苦挣扎了几天。我正在使用来自 4 个主题的记录。我需要通过 TimedWindow 聚合记录。时间到了,我想向接收器主题发送批准的消息或未批准的消息。这可能与 kafka 流有关吗?

它似乎将每条记录都放入新主题,即使窗口仍然打开,这真的不是我想要的。

下面是简单的代码:

 builder.stream(getTopicList(), Consumed.with(Serdes.ByteArray(), 
 Serdes.ByteArray()))
.flatMap(new ExceptionSafeKeyValueMapper<String, 
 FooTriggerMessage>("", Serdes.String(),
       fooTriggerSerde))
 .filter((key, value) -> value.getTriggerEventId() != null)
 .groupBy((key, value) -> value.getTriggerEventId().toString(),
       Serialized.with(Serdes.String(), fooTriggerSerde))

.windowedBy(TimeWindows.of(TimeUnit.SECONDS.toMillis(30))
.advanceBy(TimeUnit.SECONDS.toMillis(30)))

.aggregate(() -> new BarApprovalMessage(), /* initializer */
       (key, value, aggValue) -> getApproval(key, value, aggValue),/*adder*/
       Materialized
               .<String, BarApprovalMessage, WindowStore<Bytes, byte[]>>as(
                       storeName) /* state store name */
               .withValueSerde(barApprovalSerde))
.toStream().to(appProperties.getBarApprovalEngineOutgoing(), 
Produced.with(windowedSerde, barApprovalSerde));

到目前为止,每条记录都被下沉到outgoingTopic,我只希望它在窗口关闭时发送一条消息,可以这么说。

这可能吗?

【问题讨论】:

标签: apache-kafka-streams windowed


【解决方案1】:

【讨论】:

  • 酷!那么我的 hack 就不再需要了 :) 将其标记为解决方案!
【解决方案2】:

如果其他人需要答案,我会回答我自己的问题。在转换阶段,我使用上下文来创建调度程序。该调度程序采用三个参数。标点的间隔,使用的时间(挂钟或流时间)和供应商(达到时间时调用的方法)。我使用挂钟时间并为每个唯一的窗口键启动了一个新的调度程序。我在 KeyValue 存储中添加每条消息并返回 null。然后,在每 30 秒调用一次的方法中,我检查窗口是否关闭,并遍历密钥库中的消息,聚合并使用 context.forward 和 context.commit。中提琴!在 30 秒的窗口中收到 4 条消息,产生了 1 条消息。

【讨论】:

  • 愿意为此分享任何示例代码吗?非常有帮助,我正在为 Session Windows 开发一个类似的用例。
  • 我写了一篇关于它的博客文章。它是瑞典语的,但我猜代码部分还是有意义的。如果没有我可以尝试发布它! google.se/amp/s/cygni.se/oppna-fonster-med-kafka-streams/amp
  • 非常感谢!大部分都是有道理的,因为我的用例有点不同(由于会话可能不完整,所以无法安排固定等待)。但是,在您的示例中,“调度程序”来自哪里?我不确定类型是什么,也没有看到它在任何地方声明。
  • sentIdsList 的问题实际上是相同的,我假设您将其传递给已发送的地图。那是一家国营商店,还是来自不同的地方?
  • Scheduler 持有 Cancelables。我在变压器中声明它们。使用 context.schedule 以给定的间隔进行标点。 Sentidslist 是一个全局变量。但我认为你不需要关于那个的信息:)
【解决方案3】:

我遇到了这个问题,但我解决了这个问题,在固定窗口之后添加了 grace(0) 并使用了 Suppressed API

public void process(KStream<SensorKeyDTO, SensorDataDTO> stream) {

        buildAggregateMetricsBySensor(stream)
                .to(outputTopic, Produced.with(String(), new SensorAggregateMetricsSerde()));

    }

private KStream<String, SensorAggregateMetricsDTO> buildAggregateMetricsBySensor(KStream<SensorKeyDTO, SensorDataDTO> stream) {
        return stream
                .map((key, val) -> new KeyValue<>(val.getId(), val))
                .groupByKey(Grouped.with(String(), new SensorDataSerde()))
                .windowedBy(TimeWindows.of(Duration.ofMinutes(WINDOW_SIZE_IN_MINUTES)).grace(Duration.ofMillis(0)))
                .aggregate(SensorAggregateMetricsDTO::new,
                        (String k, SensorDataDTO v, SensorAggregateMetricsDTO va) -> aggregateData(v, va),
                        buildWindowPersistentStore())
                .suppress(Suppressed.untilWindowCloses(unbounded()))
                .toStream()
                .map((key, value) -> KeyValue.pair(key.key(), value));
    }


    private Materialized<String, SensorAggregateMetricsDTO, WindowStore<Bytes, byte[]>> buildWindowPersistentStore() {
        return Materialized
                .<String, SensorAggregateMetricsDTO, WindowStore<Bytes, byte[]>>as(WINDOW_STORE_NAME)
                .withKeySerde(String())
                .withValueSerde(new SensorAggregateMetricsSerde());
    }

在这里你可以看到结果

【讨论】:

  • 这是否仍然适用于事件时间。如果我们想用挂钟时间来抑制呢?
猜你喜欢
  • 2023-01-24
  • 2022-12-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-11-29
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多