在Flink中消费Kafka的JSON数据时遇代码报错求助
问题排查与解决
错误核心原因
- 泛型类型不匹配:
KafkaSource<String>声明为字符串类型,但JsonDeserializationSchema返回的是TradeEvent实体类,类型冲突导致编译错误。 - 缺少对应实体类:未定义与JSON字段映射的
TradeEvent类,反序列化器无法找到目标类型。 - 输入JSON格式非法:提供的示例JSON字段间未添加逗号,属于无效JSON(大概率是输入遗漏,需确保生产者输出合法格式)。
- 潜在依赖缺失:未引入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 }
4. 添加Flink JSON依赖
如果使用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
相关产品推荐
相关产品推荐

