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

解决Kafka嵌套JSON反序列化为Flink可用POJO的报错问题

问题场景

从Kafka获取的嵌套JSON数据:

{
    "url": "/import/something",
    "body": {
        "user_id": "4110e4f5-09d6-45d1-806a-fbe311f7fa17",
        "session_id": "8819a20b-1b19-4563-86a5-66c063e0e10d",
        "media_id": "11fb4e1a-9bb9-4b3c-90a3-c1cb2152b035",
        "timestamp": "2023-12-31T19:17:43Z"
    },
    "headers": {
        "host": "something.com",
        "x-request-id": "041117da2717b5a3a1809ec342f855be",
        "x-real-ip": "199.22.29.9",
        "x-forwarded-for": "199.22.29.9"
    }
}

使用Flink的JsonDeserializationSchema反序列化为VodMessage POJO,代码如下:

JsonDeserializationSchema<VodMessage> jsonFormat = new JsonDeserializationSchema<>(VodMessage.class);

KafkaSource<VodMessage> source = KafkaSource
        .<VodMessage>builder()
        .setValueOnlyDeserializer(jsonFormat)
        .build();

出现报错:

Caused by: org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.exc.MismatchedInputException: Cannot deserialize value of type java.lang.String from Object value (token JsonToken.START_OBJECT)

问题原因

你的VodMessage类中,body和headers字段被定义为String类型,但JSON数据里这两个字段是嵌套对象(而非字符串)。Jackson在反序列化时,尝试将对象类型转换为字符串类型,导致类型不匹配,触发该报错。

解决方法

方案1:定义嵌套POJO类

为body和headers分别创建对应的实体类,用强类型映射嵌套结构:

1.1 定义VodBody类(对应JSON中的body字段)

public class VodBody {
    private String user_id;
    private String session_id;
    private String media_id;
    private String timestamp;

    // 必须提供无参构造器,Jackson反序列化依赖
    public VodBody() {}

    // 生成所有字段的Getter和Setter方法
    public String getUser_id() { return user_id; }
    public void setUser_id(String user_id) { this.user_id = user_id; }
    public String getSession_id() { return session_id; }
    public void setSession_id(String session_id) { this.session_id = session_id; }
    public String getMedia_id() { return media_id; }
    public void setMedia_id(String media_id) { this.media_id = media_id; }
    public String getTimestamp() { return timestamp; }
    public void setTimestamp(String timestamp) { this.timestamp = timestamp; }
}

1.2 定义VodHeaders类(对应JSON中的headers字段)

注意JSON中的带连字符的字段,可以用@JsonProperty注解指定映射关系:

import com.fasterxml.jackson.annotation.JsonProperty;

public class VodHeaders {
    private String host;
    
    @JsonProperty("x-request-id")
    private String xRequestId;
    
    @JsonProperty("x-real-ip")
    private String xRealIp;
    
    @JsonProperty("x-forwarded-for")
    private String xForwardedFor;

    public VodHeaders() {}

    // 生成所有字段的Getter和Setter方法
    public String getHost() { return host; }
    public void setHost(String host) { this.host = host; }
    public String getXRequestId() { return xRequestId; }
    public void setXRequestId(String xRequestId) { this.xRequestId = xRequestId; }
    public String getXRealIp() { return xRealIp; }
    public void setXRealIp(String xRealIp) { this.xRealIp = xRealIp; }
    public String getXForwardedFor() { return xForwardedFor; }
    public void setXForwardedFor(String xForwardedFor) { this.xForwardedFor = xForwardedFor; }
}

1.3 修改VodMessage类

将body和headers字段类型替换为对应的嵌套POJO:

public class VodMessage {
    private String url;
    private VodBody body;
    private VodHeaders headers;

    public VodMessage() {}

    // 生成所有字段的Getter和Setter方法
    public String getUrl() { return url; }
    public void setUrl(String url) { this.url = url; }
    public VodBody getBody() { return body; }
    public void setBody(VodBody body) { this.body = body; }
    public VodHeaders getHeaders() { return headers; }
    public void setHeaders(VodHeaders headers) { this.headers = headers; }
}

方案2:用Map接收嵌套对象(弱类型方式)

如果不需要强类型的嵌套实体,可以将body和headers定义为Map<String, String>类型,Jackson会自动将嵌套对象转换为键值对:

public class VodMessage {
    private String url;
    private Map<String, String> body;
    private Map<String, String> headers;

    public VodMessage() {}

    // 生成所有字段的Getter和Setter方法
    public String getUrl() { return url; }
    public void setUrl(String url) { this.url = url; }
    public Map<String, String> getBody() { return body; }
    public void setBody(Map<String, String> body) { this.body = body; }
    public Map<String, String> getHeaders() { return headers; }
    public void setHeaders(Map<String, String> headers) { this.headers = headers; }
}

关键注意事项

  • 所有POJO类必须提供无参构造器,Jackson反序列化过程依赖该构造器创建实例
  • 确保所有字段都有对应的Getter和Setter方法,或使用Lombok的@Data、@NoArgsConstructor注解简化代码编写
  • 当JSON字段名与Java字段名不一致时,必须用@JsonProperty注解明确指定映射关系

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 18:07:47