如何在Kafka Streams聚合函数中获取窗口起始时间并传入add方法?
在Kafka Streams窗口聚合中获取窗口起始时间的实现方法
要在你的聚合函数中获取窗口起始时间,核心是利用Kafka Streams窗口聚合提供的Windowed<K>类型键——它同时包含了原始业务键和窗口元数据信息。具体修改步骤如下:
调整聚合函数的参数类型
将聚合lambda中的第一个参数从String类型的业务键,改为Windowed<String>类型。通过这个对象,你既可以用key()方法拿到原来的业务键(比如你的"IFTTT"),也能通过window()方法获取窗口的时间信息。获取窗口起始时间并传入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()) );
- 适配OHLC类的add方法
同步修改OHLC.add()方法的签名,新增窗口起始时间参数,示例如下:
public OHLC add(Trade trade, String key, long windowStart) { // 可将窗口起始时间用于OHLC对象的时间戳赋值等业务逻辑 this.timestamp = windowStart; // 原有业务处理逻辑... return this; }
这样就能在聚合过程中拿到窗口起始时间,同时保留原来的业务键。
内容的提问来源于stack exchange,提问作者mohammadjavadkh
相关产品推荐
相关产品推荐

