Flink中JsonDeserializationSchema.deserialize触发NullPointerException问题排查
Flink调用JsonDeserializationSchema.deserialize抛出NullPointerException问题分析
问题场景
尝试用Flink读取.jsonl格式文件,通过FileSource读取文本行后,在flatMap算子中直接调用JsonDeserializationSchema.deserialize方法转换为POJO时,抛出NullPointerException;改用ObjectMapper直接解析则正常运行。
数据文件内容
{"name": "jimmy", "age": 10} {"name": "tommy", "age": 11} {"name": "marry", "age": 12}
报错代码片段
public class FlinkApp { public static void main(String[] args) throws Exception { var env = StreamExecutionEnvironment.getExecutionEnvironment(); build(env); env.execute(); } private static void build(StreamExecutionEnvironment environment) { Path in_path = new Path("hdfs://.../people.jsonl"); FileSource<String> source = FileSource .forRecordStreamFormat(new TextLineInputFormat(), in_path) .build(); var dsSource = environment .fromSource(source, WatermarkStrategy.noWatermarks(), "text source"); JsonDeserializationSchema<People> jsonFormat = new JsonDeserializationSchema<>(People.class); dsSource .flatMap((FlatMapFunction<String, People>) (s, c) -> c.collect(jsonFormat.deserialize(s.getBytes()))) .returns(People.class) .addSink(new StdoutSink<>("json")); } } class StdoutSink<T> implements SinkFunction<T> { private String name; StdoutSink(String name) { this.name = name; } @Override public void invoke(T value, Context context) throws Exception { String b = String.format("%s: %s%n", name, value); System.out.println(b); } }
报错堆栈
java.lang.NullPointerException: null at org.apache.flink.formats.json.JsonDeserializationSchema.deserialize(JsonDeserializationSchema.java:69) ~[flink-json-1.16.1.jar:1.16.1] at my.flink.play.FlinkApp.lambda$build$29a9905b$1(FlinkApp.java:34) ~[flink-java-play-all.jar:?] at org.apache.flink.streaming.api.operators.StreamFlatMap.processElement(StreamFlatMap.java:47) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:82) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:57) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:29) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.SourceOperatorStreamTask$AsyncDataOutputToOutput.emitRecord(SourceOperatorStreamTask.java:313) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.api.operators.source.SourceOutputWithWatermarks.collect(SourceOutputWithWatermarks.java:110) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.api.operators.source.SourceOutputWithWatermarks.collect(SourceOutputWithWatermarks.java:101) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.connector.file.src.impl.FileSourceRecordEmitter.emitRecord(FileSourceRecordEmitter.java:45) ~[flink-connector-files-1.16.1.jar:1.16.1] at org.apache.flink.connector.file.src.impl.FileSourceRecordEmitter.emitRecord(FileSourceRecordEmitter.java:35) ~[flink-connector-files-1.16.1.jar:1.16.1] at org.apache.flink.connector.base.source.reader.SourceReaderBase.pollNext(SourceReaderBase.java:143) ~[flink-connector-files-1.16.1.jar:1.16.1] at org.apache.flink.streaming.api.operators.SourceOperator.emitNext(SourceOperator.java:385) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.io.StreamTaskSourceInput.emitNext(StreamTaskSourceInput.java:68) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:542) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:831) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:780) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:935) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:914) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:728) ~[flink-dist-1.16.1.jar:1.16.1] at org.apache.flink.runtime.taskmanager.Task.run(Task.java:550) ~[flink-dist-1.16.1.jar:1.16.1] at java.lang.Thread.run(Thread.java:829) ~[?:?]
问题原因
JsonDeserializationSchema的核心实现依赖**initialize方法完成内部ObjectMapper的初始化**:
- 该类的构造方法仅保存了目标POJO的类型信息,并未初始化用于JSON解析的
ObjectMapper实例; initialize方法会被Flink的连接器(如KafkaSource)在启动阶段自动调用,完成ObjectMapper的创建与配置;- 当直接在算子中实例化
JsonDeserializationSchema并调用deserialize时,initialize方法未被执行,导致内部ObjectMapper为null,进而抛出NullPointerException。
对应Flink 1.16.1版本的关键源码逻辑:
public class JsonDeserializationSchema<T> extends AbstractDeserializationSchema<T> { private final Class<T> type; private transient ObjectMapper objectMapper; public JsonDeserializationSchema(Class<T> type) { this.type = checkNotNull(type); } @Override public void initialize(InitializationContext context) throws Exception { objectMapper = createObjectMapper(); configureObjectMapper(objectMapper); } @Override public T deserialize(byte[] message) throws IOException { return objectMapper.readValue(message, type); // 此处objectMapper未初始化则抛出NPE } }
解决方案
方案1:手动调用initialize方法
在使用JsonDeserializationSchema前,手动调用initialize方法完成初始化(无需上下文时可传入null):
JsonDeserializationSchema<People> jsonFormat = new JsonDeserializationSchema<>(People.class); jsonFormat.initialize(null); // 手动触发初始化 dsSource .flatMap((FlatMapFunction<String, People>) (s, c) -> c.collect(jsonFormat.deserialize(s.getBytes()))) .returns(People.class) .addSink(new StdoutSink<>("json"));
方案2:直接使用ObjectMapper解析
如你已经尝试的方式,直接实例化ObjectMapper完成JSON解析,避开JsonDeserializationSchema的初始化依赖:
ObjectMapper mapper = new ObjectMapper(); dsSource .flatMap((FlatMapFunction<String, People>) (s, c) -> c.collect(mapper.readValue(s.getBytes(), People.class))) .returns(People.class) .addSink(new StdoutSink<>("json"));
方案3:使用FileSource配合JSON格式的RowData反序列化
如果希望通过连接器层面完成反序列化,可使用JsonRowDataDeserializationSchema配合FileSource,再将RowData转换为POJO:
// 定义JSON对应的RowType RowType rowType = RowType.of( DataTypes.STRING().notNull(), // name字段 DataTypes.INT().notNull() // age字段 ); // 创建支持JSON解析的FileSource FileSource<RowData> source = FileSource.forRecordStreamFormat( new JsonRowDataDeserializationSchema(rowType, new Configuration()), in_path) .build(); var dsSource = environment.fromSource(source, WatermarkStrategy.noWatermarks(), "json source"); // 将RowData转换为People POJO dsSource.map(row -> new People(row.getString(0), row.getInt(1))) .addSink(new StdoutSink<>("json"));
内容的提问来源于stack exchange,提问作者Bing-hsu Gao
相关产品推荐
相关产品推荐

