Dataflow任务运行失败:非可移植Dataflow不支持BundleFinalizer
我们运行一个简单的Dataflow任务,从Spanner数据库读取数据并展示变更记录的Mod Type,但部署后始终失败,报错如下:
Error message from worker: java.lang.UnsupportedOperationException: BundleFinalizer unsupported by non-portable Dataflow.
org.apache.beam.runners.dataflow.worker.SplittableProcessFnFactory$SplittableDoFnRunnerFactory.lambda$createRunner$2(SplittableProcessFnFactory.java:172)
org.apache.beam.runners.dataflow.worker.repackaged.org.apache.beam.runners.core.OutputAndTimeBoundedSplittableProcessElementInvoker$1.bundleFinalizer(OutputAndTimeBoundedSplittableProcessElementInvoker.java:206)
org.apache.beam.sdk.io.gcp.spanner.changestreams.dofn.ReadChangeStreamPartitionDoFn$DoFnInvoker.invokeProcessElement(Unknown Source)
org.apache.beam.runners.dataflow.worker.repackaged.org.apache.beam.runners.core.OutputAndTimeBoundedSplittableProcessElementInvoker.invokeProcessElement(OutputAndTimeBoundedSplittableProcessElementInvoker.java:125)
org.apache.beam.runners.dataflow.worker.repackaged.org.apache.beam.runners.core.SplittableParDoViaKeyedWorkItems$ProcessFn.processElement(SplittableParDoViaKeyedWorkItems.java:567)
相关源代码:
public static void main(String[] args) { PipelineOptions pipelineOptions = PipelineOptionsFactory.fromArgs(args).withValidation().create(); Pipeline pipeline = Pipeline.create(pipelineOptions); SpannerConfig spannerConfig = SpannerConfig.create().withProjectId("prj-test-1234") .withInstanceId("dbinstance1").withDatabaseId("dbtestnew"); PCollection<String> changeRecord = pipeline .apply(SpannerIO.readChangeStream().withSpannerConfig(spannerConfig) .withChangeStreamName("dbteststream").withMetadataDatabase("testmetadata")) .apply(ParDo.of(new DoFn<DataChangeRecord, String>() { @ProcessElement public void process(ProcessContext context) { System.out.println("context " + context.element().getModType()); context.output(context.element().getModType().name()); } })); pipeline.run(); }
问题原因
SpannerIO.readChangeStream依赖Beam的BundleFinalizer特性处理变更流的状态持久化,但旧的非Portable Dataflow Runner不支持该特性,导致作业启动失败。
解决方法
1. 启用Dataflow Portable作业提交模式
你需要通过以下方式开启Portable模式:
- 命令行参数:提交作业时添加
--experiments=use_portable_job_submission,同时指定--runner=DataflowRunner - 代码配置:在代码中转换为
DataflowPipelineOptions并设置相关参数:// 将通用PipelineOptions转换为Dataflow专用选项 DataflowPipelineOptions dfOptions = pipelineOptions.as(DataflowPipelineOptions.class); // 指定使用DataflowRunner dfOptions.setRunner(DataflowRunner.class); // 启用Portable作业提交实验特性 dfOptions.setExperiments(Collections.singletonList("use_portable_job_submission"));
2. 确认Beam SDK版本
确保项目使用的Beam SDK版本≥2.30.0,该版本开始正式支持Spanner变更流的Portable运行时兼容。
3. 权限与镜像检查
- 确保Dataflow作业使用的服务账号拥有
roles/spanner.databaseReader(读取Spanner数据)和roles/dataflow.jobRunner(提交Dataflow作业)权限 - 若使用自定义容器镜像,需确保镜像包含Beam Portable运行时依赖
内容的提问来源于stack exchange,提问作者Sabari

