Apache Beam中用户会话窗口结合有状态处理抛出异常问题咨询
问题排查:Apache Beam会话窗口结合ValueState报错解决方案
问题背景
要实现用户会话窗口级别的事件计数,业务流程如下:
- 读取输入数据(测试场景从文件读取,实际生产为流处理)
- 解析JSON格式的输入事件
- 为事件绑定对应时间戳
- 转换为
<SessionId, InputEvent>格式的键值对PCollection - 按sessionId设置间隔为2分钟的会话窗口
- 在ParDo算子中累加ValueState计数器,打印日志验证计数结果
错误信息
运行时抛出如下异常:
java.lang.UnsupportedOperationException: MergingWindowFn is not supported for stateful DoFns, WindowFn is: org.apache.beam.sdk.transforms.windowing.Sessions@1d4df at org.apache.beam.repackaged.direct_java.runners.core.StatefulDoFnRunner.rejectMergingWindowFn (StatefulDoFnRunner.java:112) at org.apache.beam.repackaged.direct_java.runners.core.StatefulDoFnRunner.<init> (StatefulDoFnRunner.java:107) at org.apache.beam.repackaged.direct_java.runners.core.DoFnRunners.defaultStatefulDoFnRunner (DoFnRunners.java:157) at org.apache.beam.runners.direct.ParDoEvaluator.lambda$defaultRunnerFactory$0 (ParDoEvaluator.java:111) at org.apache.beam.runners.direct.ParDoEvaluator.create (ParDoEvaluator.java:156) at org.apache.beam.runners.direct.ParDoEvaluatorFactory.createParDoEvaluator (ParDoEvaluatorFactory.java:152) at org.apache.beam.runners.direct.ParDoEvaluatorFactory.createEvaluator (ParDoEvaluatorFactory.java:123) at org.apache.beam.runners.direct.StatefulParDoEvaluatorFactory.createEvaluator (StatefulParDoEvaluatorFactory.java:109) at org.apache.beam.runners.direct.StatefulParDoEvaluatorFactory.forApplication (StatefulParDoEvaluatorFactory.java:89) at org.apache.beam.runners.direct.TransformEvaluatorRegistry.forApplication (TransformEvaluatorRegistry.java:178) at org.apache.beam.runners.direct.DirectTransformExecutor.run (DirectTransformExecutor.java:122) at java.util.concurrent.Executors$RunnableAdapter.call (Executors.java:511) at java.util.concurrent.FutureTask.run (FutureTask.java:266) at java.util.concurrent.ThreadPoolExecutor.runWorker (ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run (ThreadPoolExecutor.java:624) at java.lang.Thread.run (Thread.java:748)
问题原因
Beam中会话窗口(Sessions)属于合并型窗口(MergingWindowFn),运行时会动态合并多个重叠的小窗口为一个完整的会话窗口,而有状态DoFn的ValueState是和初始创建的窗口绑定的,窗口合并时State无法自动迁移合并,因此Beam原生不支持合并型窗口下直接使用有状态DoFn,这是抛出异常的根本原因。
可行解决方案
方案1:用Combine.perKey聚合替代自定义ValueState(推荐)
统计每个会话每个key的事件数是典型的聚合场景,直接用Beam自带的聚合算子即可自动兼容会话窗口的合并逻辑,无需自定义State。
修改后的核心代码如下:
pipeline // 原有步骤保持不变:读数据、解析JSON、加时间戳 .apply("ReadInputData", TextIO.read().from(options.getInputPath())) .apply("ParseJson", ParseJsons.of(InputEvents.class)) .setCoder(SerializableCoder.of(InputEvents.class)) .apply("AddTimestamp", WithTimestamps.of( (InputEvents events) -> Instant.parse(events.getTimestamp(), DateTimeFormat.forPattern("yyyy-MM-dd HH:mm:ss zzz")) ) ) // 转KV时value直接存1,减少后续计算开销 .apply("MapEventsToKV", MapElements.via( new SimpleFunction<InputEvents, KV<String, Integer>>() { @Override public KV<String, Integer> apply(InputEvents input) { return KV.of(input.getSessionId(), 1); } })) // 会话窗口设置不变 .apply("SessionWindows", Window.<KV<String, Integer>>into( Sessions.withGapDuration(Duration.standardMinutes(2)) .withTimestampCombiner(TimestampCombiner.END_OF_WINDOW) ) // 替换自定义State ParDo为Count.perKey聚合 .apply("CountEventsPerSession", Count.perKey()) // 打印日志输出结果 .apply("LogResult", ParDo.of(new DoFn<KV<String, Long>, Void>() { @ProcessElement public void processElement(@Element KV<String, Long> countResult, IntervalWindow window) { LOG.info("Window ==> {} :: counterValue ==> {}", window.toString(), countResult.getValue()); } })); return pipeline.run();
方案2:必须使用自定义State处理复杂逻辑的场景
如果业务逻辑需要自定义State做复杂计算,可以先通过GroupByKey将会话窗口内同一个key的所有事件合并为Iterable输出,此时窗口已经完成合并不会再变更,后续再在ParDo中使用State即可。
验证结果
使用给出的测试数据运行修改后的代码,可得到预期输出:
Window ==> 2021-08-29T10:54:54.000Z..2021-08-29T10:58:14.000Z :: counterValue ==> 3 Window ==> 2021-08-29T10:55:59.000Z..2021-08-29T10:58:35.000Z :: counterValue ==> 2
内容的提问来源于stack exchange,提问作者anandt
相关产品推荐
相关产品推荐

