Spark Streaming Java消费Kafka Avro数据乱码问题求助
解决Spark Streaming读取Avro格式Kafka数据乱码问题
你的问题核心是数据序列化格式不匹配——Kafka中存储的是Avro二进制格式的数据,但你的Java代码用了字符串反序列化器去解码,这必然会出现乱码。让我一步步帮你解决:
为什么会乱码?
你用kafka-avro-console-consumer能正常读取,是因为这个工具默认使用了Confluent的Avro反序列化器,它会自动从Schema Registry拉取对应的Avro Schema,把二进制数据解码成可读的JSON格式。而你的代码里用了StringDeserializer,它会把Avro二进制数据直接当成UTF-8字符串解析,那些奇怪的乱码就是这种错误解析的结果。
解决方案步骤
1. 添加必要的依赖
首先确保项目引入Confluent Avro序列化相关依赖(版本要和你的Confluent、Spark版本匹配):
<!-- Maven示例 --> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-avro-serializer</artifactId> <version>7.4.0</version> <!-- 替换为你的Confluent版本 --> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.12</artifactId> <version>3.3.0</version> <!-- 替换为你的Spark版本 --> </dependency>
2. 修改Kafka反序列化配置
更新kafkaParams,使用Avro反序列化器并指定Schema Registry地址:
Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", "localhost:9092"); kafkaParams.put("key.deserializer", StringDeserializer.class); // 你的key是null,用StringDeserializer没问题 kafkaParams.put("value.deserializer", io.confluent.kafka.serializers.KafkaAvroDeserializer.class); kafkaParams.put("group.id", "groupStreamId"); kafkaParams.put("auto.offset.reset", "latest"); kafkaParams.put("enable.auto.commit", false); // 必须添加Schema Registry地址,和console consumer保持一致 kafkaParams.put("schema.registry.url", "http://localhost:8081"); // 启用通用Avro读取器(不需要预先生成特定类) kafkaParams.put("specific.avro.reader", "false");
3. 在Spark Streaming中解析Avro数据
修改配置后,拉取的value会是GenericRecord类型(Avro通用记录),你可以从中提取字段:
JavaInputDStream<ConsumerRecord<String, GenericRecord>> stream = KafkaUtils.createDirectStream( streamingContext, LocationStrategies.PreferConsistent(), ConsumerStrategies.Subscribe(Collections.singletonList("postgres-ip_audit"), kafkaParams) ); // 处理每条记录 stream.foreachRDD(rdd -> { rdd.foreach(record -> { GenericRecord value = record.value(); // 提取字段注意类型转换,Avro的字符串会被封装成Utf8类 Integer id = (Integer) value.get("id"); String ip = ((org.apache.avro.util.Utf8) value.get("ip")).toString(); Long createTs = (Long) value.get("create_ts"); System.out.printf("id: %d, ip: %s, create_ts: %d%n", id, ip, createTs); }); });
额外注意事项
- 如果Schema Registry开启了认证,需要在
kafkaParams中添加对应的认证参数(比如basic.auth.user.info)。 - 如果你想使用特定Avro类(而非GenericRecord),可以用
avro-tools或Maven插件生成对应Java类,然后把specific.avro.reader设为true,反序列化器会自动将数据转为你的特定类。 - 若使用Spark Structured Streaming,核心逻辑一致,也可以用
from_avro函数直接解析数据。
内容的提问来源于stack exchange,提问作者Claudio Melis
相关产品推荐
相关产品推荐

