You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Kafka Streams聚合函数中获取窗口起始时间并传入add方法?

在Kafka Streams窗口聚合中获取窗口起始时间的实现方法

要在你的聚合函数中获取窗口起始时间,核心是利用Kafka Streams窗口聚合提供的Windowed<K>类型键——它同时包含了原始业务键和窗口元数据信息。具体修改步骤如下:

  1. 调整聚合函数的参数类型
    将聚合lambda中的第一个参数从String类型的业务键,改为Windowed<String>类型。通过这个对象,你既可以用key()方法拿到原来的业务键(比如你的"IFTTT"),也能通过window()方法获取窗口的时间信息。

  2. 获取窗口起始时间并传入add方法
    通过windowedKey.window().startTime().toEpochMilli()(适用于Kafka Streams 2.3+版本)或windowedKey.window().start()(旧版本)获取窗口起始时间的毫秒值,再将其传递给OHLC.add()方法。

修改后的完整代码如下:

stream
    .filter(((key, Trade) -> Trade.tradeTime != null && Trade.tradeTime > todayMillis ))
    .groupByKey(
        Grouped.with(Serdes.String(), JSONSerdes.Trade())
    )
    .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(Convertor.getCandleByResolution(resolution)), Duration.ofDays(1)))
    .aggregate(
        OHLC::new,
        ((windowedKey, value, aggregate) -> {
            // 获取窗口起始时间(毫秒级)
            long windowStart = windowedKey.window().startTime().toEpochMilli();
            // 若使用旧版Kafka Streams,替换为:long windowStart = windowedKey.window().start();
            // 传递业务key、Trade对象和窗口起始时间到add方法
            return aggregate.add(value, windowedKey.key(), windowStart);
        }),
        Materialized.<String, OHLC, WindowStore<Bytes, byte[]>>as(stateStoreName)
            .withKeySerde(Serdes.String())
            .withValueSerde(JSONSerdes.OHLC())
    );
  1. 适配OHLC类的add方法
    同步修改OHLC.add()方法的签名,新增窗口起始时间参数,示例如下:
public OHLC add(Trade trade, String key, long windowStart) {
    // 可将窗口起始时间用于OHLC对象的时间戳赋值等业务逻辑
    this.timestamp = windowStart;
    // 原有业务处理逻辑...
    return this;
}

这样就能在聚合过程中拿到窗口起始时间,同时保留原来的业务键。


内容的提问来源于stack exchange,提问作者mohammadjavadkh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 10:27:43