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
相关产品推荐
相关产品推荐

