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
相关产品推荐
相关产品推荐

