在Apache Flink中从Kafka接收JSON并实现面向对象访问的问题
问题分析与解决方案
1. 为什么source.toString()返回对象哈希值?
source是KafkaSource实例,test是DataStream实例,它们的toString()方法默认返回类标识+内存地址(比如org.apache.flink.connector.kafka.source.KafkaSource@510f3d34),和你期望的Event数据属性完全无关。
要访问Event的LaPos等属性,必须在DataStream的算子中处理单个Event对象——比如用map算子提取属性,或者通过print()直接输出Event(你已经重写了Event的toString()返回LaPos)。
你代码里的test.print()是正确方式,前面的System.out.println(source.toString())和System.out.println(test.toString())完全无效,直接删掉即可。
2. JSON反序列化的潜在问题(你大概率会遇到)
当前Event POJO和JSON数据存在不匹配,会导致反序列化失败,即使代码能跑也拿不到正确的Event对象:
- 字段名大小写不匹配:JSON中的键是
"Timestamp"(首字母大写),但POJO里的字段是timestamp(首字母小写),Jackson默认不会自动匹配,需要添加@JsonProperty注解指定映射关系。 - Timestamp类型反序列化失败:JSON中的
Timestamp是字符串格式(如"2022-10-31 12:45:19.353"),默认的JsonDeserializationSchema无法直接转为java.sql.Timestamp,需要先以String接收再手动转换。
修正后的Event POJO示例
import com.fasterxml.jackson.annotation.JsonProperty; import java.sql.Timestamp; public class Event { public double LoPos; public double LaPos; @JsonProperty("Timestamp") // 匹配JSON中的大写键名 private String timestampStr; // 先以String接收时间字符串 public Timestamp timestamp; // 空构造函数必须保留,Flink反序列化依赖它 public Event() {} public Event(double loPos, double laPos, Timestamp timestamp) { LoPos = loPos; LaPos = laPos; this.timestamp = timestamp; } // 反序列化后自动将字符串转为Timestamp public void setTimestampStr(String timestampStr) { this.timestampStr = timestampStr; // 适配"yyyy-MM-dd HH:mm:ss.SSS"格式的时间字符串 this.timestamp = Timestamp.valueOf(timestampStr); } @Override public String toString() { return String.valueOf(LaPos); } }
3. 修正后的Flink消费代码
删掉无效的打印语句,保留正确的数据流处理:
KafkaSource<Event> source = KafkaSource.<Event>builder() .setBootstrapServers("localhost:9092") .setBounded(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(new JsonDeserializationSchema<>(Event.class)) .setTopics("testTopic2") .build(); DataStream<Event> test = env.fromSource(source, WatermarkStrategy.noWatermarks(), "test"); // 打印每个Event对象,会调用你重写的toString()输出LaPos test.print(); env.execute("Kafka JSON to POJO Job");
4. 额外验证点
确保Kafka Producer发送的是JSON字符串:当前代码发送的是ObjectNode,建议显式转为字符串后发送,避免序列化格式不匹配:
String jsonStr = test.writeValueAsString(jNode); ProducerRecord<String, String> rec = new ProducerRecord<>(topicName, jsonStr); producer.send(rec);
内容的提问来源于stack exchange,提问作者Mike
相关产品推荐
相关产品推荐

