解决Kafka嵌套JSON反序列化为Flink可用POJO的报错问题
Flink JSON反序列化嵌套JSON报错问题解决
问题场景
从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
相关产品推荐
相关产品推荐

