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

Spring Boot集成Spark时ObjectMapper序列化失败问题求助

Spring Boot集成Spark时ObjectMapper序列化异常的解决方案

问题背景

在基于Spring Boot 2.6.7的应用中集成Spark 2.11,定义了如下Bean:

@Component
@RequiredArgsConstructor
public class SomeService implements FlatMapFunction<T, K> {

  private final ObjectMapper mapper;
  
}

注入的是Spring上下文提供的标准ObjectMapper实例,但运行时抛出org.apache.spark.SparkException: Task not serializable异常,序列化栈信息如下:

Caused by: java.io.NotSerializableException: org.springframework.http.converter.json.SpringHandlerInstantiator
Serialization stack:
- object not serializable (class: org.springframework.http.converter.json.SpringHandlerInstantiator, value: org.springframework.http.converter.json.SpringHandlerInstantiator@6e4912db)
- field (class: com.fasterxml.jackson.databind.cfg.BaseSettings, name: _handlerInstantiator, type: class com.fasterxml.jackson.databind.cfg.HandlerInstantiator)
- object (class com.fasterxml.jackson.databind.cfg.BaseSettings, com.fasterxml.jackson.databind.cfg.BaseSettings@155616d8)
- field (class: com.fasterxml.jackson.databind.cfg.MapperConfig, name: _base, type: class com.fasterxml.jackson.databind.cfg.BaseSettings)
- object (class com.fasterxml.jackson.databind.DeserializationConfig, com.fasterxml.jackson.databind.DeserializationConfig@66e72ca2)
- field (class: com.fasterxml.jackson.databind.ObjectMapper, name: _deserializationConfig, type: class com.fasterxml.jackson.databind.DeserializationConfig)
- object (class com.fasterxml.jackson.databind.ObjectMapper, com.fasterxml.jackson.databind.ObjectMapper@433ef204)
- field (class: com.smth.SomeService, name: mapper, type: class com.fasterxml.jackson.databind.ObjectMapper)

异常根源是SpringHandlerInstantiator不可序列化,目前临时解决方案是在构造函数中手动实例化ObjectMapper:

public SomeService() {
  this.mapper = new ObjectMapper();
}

需要基于Spring依赖注入的正确解决方法。

解决方法

方法1:自定义可序列化的ObjectMapper Bean

创建排除SpringHandlerInstantiator的ObjectMapper Bean,保留Spring默认配置的同时保证序列化能力:

@Configuration
public class JacksonConfig {

    @Bean
    public ObjectMapper sparkCompatibleObjectMapper(ObjectMapper springObjectMapper) {
        // 复制Spring默认ObjectMapper的所有配置
        ObjectMapper mapper = springObjectMapper.copy();
        // 移除不可序列化的SpringHandlerInstantiator
        mapper.setHandlerInstantiator(null);
        return mapper;
    }
}

然后在SomeService中注入这个自定义Bean:

@Component
@RequiredArgsConstructor
public class SomeService implements FlatMapFunction<T, K> {

  private final ObjectMapper sparkCompatibleObjectMapper;
  
}

方法2:使用@Transient标记并延迟初始化ObjectMapper

将ObjectMapper标记为@Transient,在Spark任务执行时初始化可序列化的实例,同时同步Spring的配置:

@Component
public class SomeService implements FlatMapFunction<T, K> {

  @Transient
  private ObjectMapper mapper;

  // 注入Spring默认的ObjectMapper用于复制配置
  private final ObjectMapper springObjectMapper;

  public SomeService(ObjectMapper springObjectMapper) {
    this.springObjectMapper = springObjectMapper;
  }

  @Override
  public Iterator<K> call(T t) throws Exception {
    if (mapper == null) {
      // 复制Spring配置并移除不可序列化组件
      mapper = springObjectMapper.copy();
      mapper.setHandlerInstantiator(null);
    }
    // 业务逻辑处理
    return ...;
  }
}

方法3:结合Spark广播变量传递ObjectMapper

在Spark任务提交阶段,将配置好的可序列化ObjectMapper包装为广播变量,避免重复序列化:

// 在提交Spark任务的Spring组件中
@Autowired
private ObjectMapper springObjectMapper;

public void submitSparkTask() {
    ObjectMapper sparkMapper = springObjectMapper.copy();
    sparkMapper.setHandlerInstantiator(null);
    Broadcast<ObjectMapper> mapperBroadcast = sparkContext.broadcast(sparkMapper);

    // 将广播变量传入SomeService
    JavaRDD<K> resultRdd = inputRdd.flatMap(new SomeService(mapperBroadcast));
}

// 修改SomeService接收广播变量
public class SomeService implements FlatMapFunction<T, K> {

  private final Broadcast<ObjectMapper> mapperBroadcast;

  public SomeService(Broadcast<ObjectMapper> mapperBroadcast) {
    this.mapperBroadcast = mapperBroadcast;
  }

  @Override
  public Iterator<K> call(T t) throws Exception {
    ObjectMapper mapper = mapperBroadcast.value();
    // 业务逻辑处理
    return ...;
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 16:35:23