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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 02:15:35