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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:59:37