Apache Beam转换在GCP Dataflow窗口中产生额外空元素原因排查
GCP Dataflow窗口去重时出现额外空消息问题
我从GCP Pub/Sub接收并分析JSON格式的流数据,包含320K条唯一消息。但在去重窗口处理时,出现了额外8.5K条空消息,且空元素数量不固定。
如何发现额外元素
- 检查每个集合的值并写入日志,可看到全部320K条非空消息及额外8.5K条空消息。
- 通过Dataflow UI查看(附图:
)
代码实现
Repeatedly triggerPersonWindow = Repeatedly.forever(AfterProcessingTime.pastFirstElementInPane()); Window<KV<String, PubSubRelationsMessage>> personWindow = Window.<KV<String, PubSubRelationsMessage>>into(FixedWindows.of(Duration.standardMinutes(10))) .triggering(triggerPersonWindow) .withAllowedLateness(Duration.ZERO) .discardingFiredPanes(); PCollectionTuple separateRelationMessages = parsedRelationMessages.get(ParseSplitPersonCompanyMessage.parsedRelationMessageTag) .apply("Person Window", personWindow) .apply("PS Deduplication", GroupByKey.create()) .apply("Separate Messages", ParDo.of(new SeparateRelationMessages()) .withOutputTags(SeparateRelationMessages.updateDeleteFlowTag, TupleTagList.of(SeparateRelationMessages.addedTag).and(GlobalTupleTag.InputError).and(GlobalTupleTag.GeneralError)));
环境版本
- Apache Beam 2.37.0
- Gradle 7.2
- Java 17
已尝试的解决方法
- 升级Apache Beam、Gradle及Java版本
- 将Dataflow Worker数量减少至1个
以上方法均未解决问题,恳请帮助。
内容的提问来源于stack exchange,提问作者Itay A
相关产品推荐
相关产品推荐

