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

Apache Beam转换在GCP Dataflow窗口中产生额外空元素原因排查

GCP Dataflow窗口去重时出现额外空消息问题

我从GCP Pub/Sub接收并分析JSON格式的流数据,包含320K条唯一消息。但在去重窗口处理时,出现了额外8.5K条空消息,且空元素数量不固定。

如何发现额外元素

  • 检查每个集合的值并写入日志,可看到全部320K条非空消息及额外8.5K条空消息。
  • 通过Dataflow UI查看(附图: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 04:18:19