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

Java转Kotlin后Apache Beam+Dataflow出现ReflectionCache序列化异常

问题根因

Apache Beam 运行时需要将所有流水线变换相关的对象序列化后分发到Worker节点执行。你代码中使用的注册了Kotlin模块的ObjectMapper实例,其内部依赖的com.fasterxml.jackson.module.kotlin.ReflectionCache未实现Java序列化接口,因此触发序列化异常。

解决方案

以下两种方案均可解决该问题,可根据场景选择:

方案1:自定义DoFn实现JSON解析(推荐)

通过Beam的@Setup生命周期注解在Worker节点本地初始化ObjectMapper,避免跨节点序列化:

class ParseExchangeObjectFn : DoFn<String, ExchangeObject>() {
    // 标记为transient避免参与序列化
    @Transient
    private lateinit var mapper: ObjectMapper

    // Worker节点初始化时执行,仅执行一次
    @Setup
    fun setup() {
        mapper = ObjectMapper().registerKotlinModule()
    }

    @ProcessElement
    fun processElement(@Element jsonLine: String, output: OutputReceiver<ExchangeObject>) {
        output.output(mapper.readValue(jsonLine, ExchangeObject::class.java))
    }
}

替换原有解析逻辑即可:

val exchangeObject = lines
    .apply("Convert Data to ExchangeObject model data", ParDo.of(ParseExchangeObjectFn()))
    .setCoder(SerializableCoder.of(ExchangeObject::class.java))

方案2:使用静态ObjectMapper实例

将ObjectMapper声明为类的静态成员,静态成员不会参与类实例的序列化流程:

class YourPipelineClass {
    companion object {
        private val objectMapper = ObjectMapper().registerKotlinModule()
    }

    // 流水线构建逻辑
    fun buildPipeline(options: YourOptions) {
        val lines = pipeline.apply("Reading Input Data ", TextIO.read().from(options.inputFile))
        val exchangeObject = lines
            .apply("Convert Data to ExchangeObject model data", ParseJsons.of(ExchangeObject::class.java).withMapper(objectMapper))
            .setCoder(SerializableCoder.of(ExchangeObject::class.java))
    }
}

注意:如果你的ExchangeObject数据类完全兼容Java Bean规范(包含无参构造函数、字段均为public或提供标准get/set方法),也可以直接去掉.withMapper(objectMapper)配置,使用ParseJsons默认提供的序列化能力即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 16:48:04