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

在Flink中消费Kafka的JSON数据时遇代码报错求助

问题排查与解决

错误核心原因

  1. 泛型类型不匹配:KafkaSource<String>声明为字符串类型,但JsonDeserializationSchema返回的是TradeEvent实体类,类型冲突导致编译错误。
  2. 缺少对应实体类:未定义与JSON字段映射的TradeEvent类,反序列化器无法找到目标类型。
  3. 输入JSON格式非法:提供的示例JSON字段间未添加逗号,属于无效JSON(大概率是输入遗漏,需确保生产者输出合法格式)。
  4. 潜在依赖缺失:未引入Flink JSON序列化相关依赖,可能导致运行时找不到JsonDeserializationSchema。

分步解决

1. 修正KafkaSource泛型类型

将KafkaSource<String>改为KafkaSource<TradeEvent>,确保泛型与反序列化器输出类型一致:

KafkaSource<TradeEvent> source = KafkaSource.<TradeEvent>builder()
        .setBootstrapServers("localhost:19092")
        .setTopics("1000pepeusdt","1000pepeusdtt")
        .setStartingOffsets(OffsetsInitializer.earliest())
        .setValueOnlyDeserializer(new JsonDeserializationSchema<>(TradeEvent.class))
        .build();

2. 创建TradeEvent实体类

根据JSON字段创建实体类,字段名需与JSON键完全一致(大小写敏感),可使用Lombok简化代码:

package org.example;

import lombok.Data;

@Data
public class TradeEvent {
    private String e;
    private Long E;
    private Long a;
    private String s;
    private String p;
    private String q;
    private Long f;
    private Long l;
    private Long T;
    private Boolean m;
}

若不使用Lombok,需手动添加所有字段的无参构造方法、getter和setter。

3. 修复生产者JSON格式

确保Kafka中的JSON数据为合法格式,字段间添加逗号:

{
"e":"aggTrade",
"E":1684547379069,
"a":30161410,
"s":"1000PEPEUSDT",
"p":"0.0015036",
"q":"665070",
"f":70012793,
"l":70012793,
"T":1684547378913,
"m":true
}

如果使用Maven,在pom.xml中添加:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-json</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>

Gradle则添加:

implementation "org.apache.flink:flink-json:${flinkVersion}"

替换${flink.version}为你使用的Flink版本。

5. 清理无效注释

移除代码中错误的注释**\the error is on this line**,确保代码编译通过。


完整修正后代码

package org.example;

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.formats.json.JsonDeserializationSchema;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;

public class DataStreamJob {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        KafkaSource<TradeEvent> source = KafkaSource.<TradeEvent>builder()
                .setBootstrapServers("localhost:19092")
                .setTopics("1000pepeusdt","1000pepeusdtt")
                .setStartingOffsets(OffsetsInitializer.earliest())
                .setValueOnlyDeserializer(new JsonDeserializationSchema<>(TradeEvent.class))
                .build();

        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source")
                .print();

        env.execute("Kafka Consumer Example");
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 04:30:02