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

Kafka Streams中JoinWindows.of(Duration)弃用的替代方案咨询

KStream与Windowed KTable连接时JoinWindows.of(Duration)弃用的替代方案

问题描述

需要将KStream<String, String>与KTable<Windowed<String>, int[]>连接以获取过去一小时的结果,原代码使用JoinWindows.of(Duration)实现,但该方法已被弃用,且尝试改用leftJoin也未解决弃用提示问题。原代码如下:

Duration windowSize = Duration.ofMinutes(60);
Duration advanceSize = Duration.ofMinutes(1);
TimeWindows hoppingWindow = TimeWindows.ofSizeWithNoGrace(windowSize).advanceBy(advanceSize);
Duration joinWindowSizeMs = Duration.ofHours(1);

// Aggregate to get [sum, count] in the last time window
KTable<Windowed<String>, int[]> averageTemp = mainStreamStandard.groupByKey()
    .windowedBy(hoppingWindow)
    .aggregate( () -> new int[]{0 ,0}, (aggKey, newVal, aggValue) -> {
            aggValue[0] += Integer.valueOf(newVal.split(":")[1]);
            aggValue[1] += 1;  
            return aggValue;
        }, Materialized.with(Serdes.String(), new IntArraySerde()));
        
// Join weather stations with their [sum,count] and their respective red alert events
KStream<String, String> joined = mainStreamAlert.join(averageTemp,
    JoinWindows.of(joinWindowSizeMs),
    (leftValue, rightValue) -> "left/" + leftValue + "/right/" + rightValue[0]/rightValue[1]);

解决方案

1. 替换弃用的JoinWindows创建方法

JoinWindows.of(Duration)弃用的核心原因是它默认将时间差平均分配到记录时间的前后(比如1小时窗口会匹配记录时间前后各30分钟的条目),不符合多数场景的精确需求。推荐使用以下两种替代方式:

方式一:快速替换(使用ofTimeDifferenceWithNoGrace)

该方法直接指定允许的时间差范围,且不设置grace period,适合不需要处理迟到数据的场景:

KStream<String, String> joined = mainStreamAlert.join(averageTemp,
    JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofHours(1)),
    (leftValue, rightValue) -> "left/" + leftValue + "/right/" + (rightValue[0] / (double) rightValue[1])
);

方式二:精确控制时间范围(推荐)

如果需要明确只匹配KStream记录时间过去1小时内的窗口数据,可以通过before和after方法精确设置时间偏移,同时关闭grace period:

KStream<String, String> joined = mainStreamAlert.join(averageTemp,
    JoinWindows.of(Duration.ofHours(1))
        .before(Duration.ZERO) // 不匹配记录时间之后的窗口
        .after(Duration.ofHours(1)) // 匹配记录时间之前1小时内的窗口
        .withNoGracePeriod(),
    (leftValue, rightValue) -> "left/" + leftValue + "/right/" + (rightValue[0] / (double) rightValue[1])
);

2. 额外优化建议

  • 修复平均值计算精度:原代码中rightValue[0]/rightValue[1]是整数除法,会丢失小数精度,建议强制转换为double类型,如上述代码所示。
  • 替换int[]为自定义POJO:使用int[]作为聚合值可读性差,且序列化易出问题,建议定义SumCount类:
    public class SumCount {
        private int sum;
        private int count;
        // 构造方法、getter、setter
    }
    
    聚合逻辑修改为:
    KTable<Windowed<String>, SumCount> averageTemp = mainStreamStandard.groupByKey()
        .windowedBy(hoppingWindow)
        .aggregate(
            SumCount::new,
            (aggKey, newVal, aggValue) -> {
                aggValue.setSum(aggValue.getSum() + Integer.valueOf(newVal.split(":")[1]));
                aggValue.setCount(aggValue.getCount() + 1);
                return aggValue;
            },
            Materialized.with(Serdes.String(), Serdes.serdeFrom(new JsonSerializer<>(), new JsonDeserializer<>(SumCount.class)))
        );
    
  • 合理设置窗口Grace Period:原代码使用ofSizeWithNoGrace,如果存在迟到数据会被直接丢弃,建议根据业务需求设置grace period,比如:
    TimeWindows hoppingWindow = TimeWindows.ofSize(windowSize)
        .advanceBy(advanceSize)
        .grace(Duration.ofMinutes(5)); // 允许5分钟内的迟到数据
    

关于leftJoin的说明

leftJoin与join的区别是保留KStream中无匹配KTable条目的记录,但它同样需要使用新的JoinWindows创建方法,弃用提示的根源是JoinWindows的初始化方式,而非join类型。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 02:31:00