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

Spark提交任务时JSON转HashMap遇序列化异常,本地运行正常

问题分析与解决方法

问题根源

虽然你的Test类实现了Serializable,但Spark的map算子闭包捕获了类中的objectMapper(或jsonParser)成员变量。这些JSON解析器内部依赖的组件(比如Jackson的ScalaTypeModifier)并没有实现序列化接口,当Spark要把闭包序列化发送到Executor节点时,就会触发NotSerializableException。

本地运行没问题是因为:本地模式下Task直接在Driver进程内执行,不需要序列化闭包跨节点传输,所以不会暴露这个序列化问题。

解决方法

方法1:在map算子内部创建解析器实例

把objectMapper或jsonParser的初始化逻辑放到map里面,让每个Task自己初始化解析器,避免序列化外部对象:

class Test extends Serializable {
    def action() = {
        val sc = SparkContext.getOrCreate(sparkConf)
        val rdd1 = sc.textFile("your_path")
        val rdd2 = rdd1.map ( logline => {
            // Jackson初始化放到map内部
            val objectMapper = new ObjectMapper().registerModule(DefaultScalaModule)
            val jsonObject = objectMapper.readValue(logline, classOf[HashMap[String,String]])
            MyDataSet(jsonObject.get("field1"), jsonObject.get("field2"), ...)              
        } )
    }
}

方法2:用transient标记解析器,延迟初始化

把解析器变量标记为@transient(不会被序列化),然后通过方法延迟初始化,确保每个Executor上的Task都能重新创建实例:

class Test extends Serializable {
    @transient private var objectMapper: ObjectMapper = _

    private def getObjectMapper(): ObjectMapper = {
        if (objectMapper == null) {
            objectMapper = new ObjectMapper().registerModule(DefaultScalaModule)
        }
        objectMapper
    }

    def action() = {
        val sc = SparkContext.getOrCreate(sparkConf)
        val rdd1 = sc.textFile("your_path")
        val rdd2 = rdd1.map ( logline => {
            val jsonObject = getObjectMapper().readValue(logline, classOf[HashMap[String,String]])
            MyDataSet(jsonObject.get("field1"), jsonObject.get("field2"), ...)              
        } )
    }
}

方法3:使用广播变量传递配置(可选)

如果担心频繁创建解析器的性能损耗,可以广播Jackson的配置模块,在Task内基于广播的配置创建解析器:

class Test extends Serializable {
    def action() = {
        val sc = SparkContext.getOrCreate(sparkConf)
        // 广播Scala模块配置
        val scalaModule = sc.broadcast(DefaultScalaModule)
        val rdd1 = sc.textFile("your_path")
        val rdd2 = rdd1.map ( logline => {
            val objectMapper = new ObjectMapper().registerModule(scalaModule.value)
            val jsonObject = objectMapper.readValue(logline, classOf[HashMap[String,String]])
            MyDataSet(jsonObject.get("field1"), jsonObject.get("field2"), ...)              
        } )
    }
}

注意:Gson的jsonParser问题同理,用上述任意一种方法处理即可。

内容的提问来源于stack exchange,提问作者J.soo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 18:10:31