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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 07:07:02