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

Flink中String转JSONObject失败:已解析但无法收集报空键错误

Flink中String转JSONObject收集时报错"Assigned key must not be null!"的解决方法

问题根源

你遇到的问题本质是Flink无法正确处理FastJSON JSONObject的序列化/反序列化逻辑。虽然你成功解析出JSONObject,但FastJSON的JSONObject并非Flink原生支持的友好类型:它本质是动态Map结构,没有实现Flink要求的序列化规范,且Flink的类型推断器无法正确识别其类型信息,最终导致内部处理时触发看似与key相关的异常(实际是序列化环节的错误)。

解决方案

方案一:转换为POJO类(推荐)

将JSON结构映射为标准POJO类,Flink对POJO类型有完善的支持,性能和可维护性都更优。

  1. 创建对应JSON结构的POJO类:
// 外层对象
public class BillData {
    private BillInfo bill_info;

    public BillInfo getBill_info() {
        return bill_info;
    }

    public void setBill_info(BillInfo bill_info) {
        this.bill_info = bill_info;
    }
}

// 嵌套的bill_info对象
public class BillInfo {
    private String ADD_TIME;
    private String ORDER_ID;
    private String ADDER_NO;
    private String UPDATER_NO;
    private String S_USER_ID;
    private String B_USER_ID;
    private String BILL_ID;
    private String ADDER_NAME;
    private String UPDATE_TIME;
    private String UPDATER_NAME;

    // 必须提供无参构造函数
    public BillInfo() {}

    // 所有字段的getter和setter方法
    public String getADD_TIME() { return ADD_TIME; }
    public void setADD_TIME(String ADD_TIME) { this.ADD_TIME = ADD_TIME; }
    public String getORDER_ID() { return ORDER_ID; }
    public void setORDER_ID(String ORDER_ID) { this.ORDER_ID = ORDER_ID; }
    // 其他字段的getter/setter同理补充
}
  1. 修改业务代码,直接解析为POJO:
SingleOutputStreamOperator<BillData> billDS = inputDS.map(value -> JSONObject.parseObject(value, BillData.class));
billDS.print();

方案二:显式指定类型信息(若需保留JSONObject)

如果必须使用FastJSON的JSONObject,可以通过显式指定类型信息或注册序列化器来解决:

方法1:指定返回类型信息

import org.apache.flink.api.java.typeutils.TypeExtractor;

SingleOutputStreamOperator<JSONObject> jsonObjDS = inputDS.process(new ProcessFunction<String, JSONObject>() {
    @Override
    public void processElement(String value, Context ctx, Collector<JSONObject> out) throws Exception {
        JSONObject jsonObject = JSONObject.parseObject(value);
        out.collect(jsonObject);
    }
}).returns(TypeExtractor.getForClass(JSONObject.class));

方法2:注册FastJSON序列化器

在初始化环境时,将JSONObject的序列化器注册到Flink配置中:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.getConfig().registerTypeWithKryoSerializer(JSONObject.class, com.alibaba.fastjson.serializer.JSONSerializer.class);
env.setParallelism(1);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 08:56:03