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

Apache Beam中SerializableFunction替代DoFn@Setup的序列化方案咨询

问题背景

使用Apache Avro v1.11.0自动生成类Foo,将Foo.Builder作为Bar类的属性。在Apache Beam v2.51.0管道中:

  • 当Bar在DoFn内使用并通过@Setup方法初始化时,运行正常;
  • 改用MapElements+SerializableFunction(或ProcessFunction),并将Bar作为实例属性(private Bar bar = new Bar();)初始化时,抛出NotSerializableException,错误根源为com.my_project.avro.outgoing_data.Foo$Builder不可序列化。

已尝试的方案:

  • 强制转换Foo.Builder为Serializable失败;
  • 在SerializableFunction构造函数中初始化Bar,仍报相同错误;
  • 将Bar设为SerializableFunction的静态成员,可行,但希望使用实例属性。

完整错误栈:

Exception in thread "main" java.lang.IllegalArgumentException: unable to serialize DoFnWithExecutionInformation{doFn=org.apache.beam.sdk.transforms.MapElements$MapWithFailures$2@536b71b4, mainOutputTag=Tag<org.apache.beam.sdk.transforms.MapElements$MapWithFailures$MapWithFailuresDoFn$1.:398#2249e1908bcf01f3>, sideInputMapping={}, schemaInformation=DoFnSchemaInformation{elementConverters=[], fieldAccessDescriptor=*}}
at org.apache.beam.sdk.util.SerializableUtils.serializeToByteArray(SerializableUtils.java:59)
...
Caused by: java.io.NotSerializableException: com.my_project.avro.outgoing_data.Foo$Builder
at java.base/java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1185)
...

解决方案

方案1:用@Transient标记+延迟初始化Foo.Builder

在Bar类中,将Foo.Builder标记为@Transient避免序列化该属性,然后在第一次使用时初始化,模拟DoFn的@Setup在工作节点初始化的逻辑:

public class Bar implements Serializable {
    @Transient
    private Foo.Builder fooBuilder;

    public Foo.Builder getFooBuilder() {
        if (fooBuilder == null) {
            fooBuilder = Foo.newBuilder();
        }
        return fooBuilder;
    }

    // 业务方法实现
}

序列化Bar时会跳过fooBuilder,反序列化后首次调用getFooBuilder()时,会在当前工作节点重新初始化构建器。

方案2:延迟初始化Bar实例

在SerializableFunction中,将Bar标记为@Transient,并通过SerializableSupplier延迟创建实例,确保Bar(及内部的Foo.Builder)在序列化完成后才初始化:

public class MyMapFunction implements SerializableFunction<Input, Output> {
    private final SerializableSupplier<Bar> barSupplier = () -> new Bar();
    @Transient
    private Bar bar;

    @Override
    public Output apply(Input input) {
        if (bar == null) {
            bar = barSupplier.get();
        }
        // 使用bar处理输入数据
        return ...;
    }
}

方案3:修改Avro生成配置,让Builder实现Serializable

Avro默认生成的Builder类不实现Serializable,可以通过配置让生成的类自动实现该接口。以Maven插件为例,添加<enableBuilderSerialization>true</enableBuilderSerialization>参数:

<plugin>
    <groupId>org.apache.avro</groupId>
    <artifactId>avro-maven-plugin</artifactId>
    <version>1.11.0</version>
    <executions>
        <execution>
            <phase>generate-sources</phase>
            <goals>
                <goal>schema</goal>
            </goals>
            <configuration>
                <sourceDirectory>${project.basedir}/src/main/avro</sourceDirectory>
                <outputDirectory>${project.build.directory}/generated-sources/avro</outputDirectory>
                <enableBuilderSerialization>true</enableBuilderSerialization>
            </configuration>
        </execution>
    </executions>
</plugin>

重新生成Foo类后,Foo.Builder会实现Serializable,直接序列化Bar实例即可正常运行。

方案4:用DoFn包装MapElements逻辑

如果想保留MapElements的链式语法,同时利用DoFn的@Setup机制,可以直接传入自定义DoFn:

PCollection<Input> input = ...;
PCollection<Output> output = input.apply(MapElements.via(new DoFn<Input, Output>() {
    private Bar bar;

    @Setup
    public void setup() {
        bar = new Bar();
    }

    @ProcessElement
    public void processElement(ProcessContext c) {
        // 使用bar处理输入元素
        c.output(...);
    }
}));

这种方式既保持了MapElements的简洁性,又避免了序列化问题。


内容的提问来源于stack exchange,提问作者Daniel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 17:40:58