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

Spark 2.3使用DStream从Kafka读取Avro记录转JSON的问题求助

问题分析

你遇到的错误原因有两点:

  1. org.apache.spark.sql.avro.AvroDeserializer确实是Spark 2.4.0及以上版本才引入的类,Spark 2.3中不存在这个类;
  2. 该类是Spark SQL模块专为DataFrame设计的反序列化器,并不符合Kafka Consumer Deserializer的要求——Kafka的反序列化器必须拥有公共无参构造方法,而这个类不满足,因此Kafka无法实例化它。

下面提供两种适配Spark 2.3 + DStream场景的可行方案:


方案一:使用原生Apache Avro库手动解析

直接读取Kafka中的二进制Avro数据,再用Avro原生API解析并转换为JSON字符串,适用于已知Avro Schema的场景。

步骤1:添加依赖

确保项目引入Apache Avro依赖(版本建议和Spark 2.3内置的Avro版本一致,即1.8.2):

<!-- Maven示例 -->
<dependency>
    <groupId>org.apache.avro</groupId>
    <artifactId>avro</artifactId>
    <version>1.8.2</version>
</dependency>

步骤2:修改代码实现

import org.apache.avro.Schema
import org.apache.avro.generic.{GenericDatumReader, GenericDatumWriter, GenericRecord}
import org.apache.avro.io.{DecoderFactory, EncoderFactory}
import org.apache.kafka.common.serialization.{ByteArrayDeserializer, StringDeserializer}
import org.apache.spark.streaming.kafka010.{ConsumerStrategies, KafkaUtils, LocationStrategies}
import java.io.StringWriter
import java.io.File

// 1. 加载Avro Schema(可从本地文件或字符串读取)
val avroSchema = new Schema.Parser().parse(new File("path/to/your/schema.avsc"))
val datumReader = new GenericDatumReader[GenericRecord](avroSchema)
val datumWriter = new GenericDatumWriter[GenericRecord](avroSchema)

// 2. 修改Kafka参数,将value反序列化器改为ByteArrayDeserializer(读取原始二进制数据)
val kafkaParams: Map[String, Object] = Map(
    "bootstrap.servers" -> "kafka-servers",
    "key.deserializer" -> classOf[StringDeserializer],
    "value.deserializer" -> classOf[ByteArrayDeserializer],
    "auto.offset.reset" -> "earliest",
    "enable.auto.commit" -> (false: java.lang.Boolean),
    "group.id" -> "group1"
)

// 3. 创建DStream并处理数据
val kafkaDstream = KafkaUtils.createDirectStream[String, Array[Byte]](
    ssc,
    LocationStrategies.PreferConsistent,
    ConsumerStrategies.Subscribe[String, Array[Byte]](topics, kafkaParams)
)

val processedStream = kafkaDstream.map { record =>
    // 解析Avro二进制数据为GenericRecord
    val decoder = DecoderFactory.get().binaryDecoder(record.value(), null)
    val avroRecord = datumReader.read(null, decoder)
    
    // 将GenericRecord转换为JSON字符串
    val writer = new StringWriter()
    val jsonEncoder = EncoderFactory.get().jsonEncoder(avroSchema, writer)
    datumWriter.write(avroRecord, jsonEncoder)
    jsonEncoder.flush()
    
    (record.key(), writer.toString)
}

// 后续输出或处理逻辑
processedStream.foreachRDD { rdd =>
    rdd.foreach { case (key, jsonStr) =>
        println(jsonStr)
    }
}

方案二:适配Confluent Schema Registry(如果使用Confluent Avro)

如果你的Kafka数据是通过Confluent平台生成的带Schema Registry的Avro数据,可以使用Confluent官方的kafka-avro-serializer库自动从Schema Registry拉取Schema并解析。

步骤1:添加依赖

引入Confluent的序列化器依赖(版本需和你的Confluent平台兼容,Spark 2.3建议选5.0.0版本):

<!-- Maven示例 -->
<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-avro-serializer</artifactId>
    <version>5.0.0</version>
</dependency>

步骤2:修改代码实现

import io.confluent.kafka.serializers.KafkaAvroDeserializer
import io.confluent.kafka.serializers.KafkaAvroDeserializerConfig
import org.apache.avro.generic.{GenericDatumWriter, GenericRecord}
import org.apache.avro.io.EncoderFactory
import org.apache.kafka.common.serialization.StringDeserializer
import org.apache.spark.streaming.kafka010.{ConsumerStrategies, KafkaUtils, LocationStrategies}
import java.io.StringWriter

// 1. 配置Kafka参数,指定Schema Registry地址
val kafkaParams: Map[String, Object] = Map(
    "bootstrap.servers" -> "kafka-servers",
    "key.deserializer" -> classOf[StringDeserializer],
    "value.deserializer" -> classOf[KafkaAvroDeserializer],
    "auto.offset.reset" -> "earliest",
    "enable.auto.commit" -> (false: java.lang.Boolean),
    "group.id" -> "group1",
    KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG -> "http://your-schema-registry:8081",
    KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG -> "false" // 使用GenericRecord而非具体类
)

// 2. 创建DStream并处理数据
val kafkaDstream = KafkaUtils.createDirectStream[String, Object](
    ssc,
    LocationStrategies.PreferConsistent,
    ConsumerStrategies.Subscribe[String, Object](topics, kafkaParams)
)

val processedStream = kafkaDstream.map { record =>
    val avroRecord = record.value().asInstanceOf[GenericRecord]
    
    // 转换为JSON字符串
    val writer = new StringWriter()
    val jsonEncoder = EncoderFactory.get().jsonEncoder(avroRecord.getSchema, writer)
    val datumWriter = new GenericDatumWriter[GenericRecord](avroRecord.getSchema)
    datumWriter.write(avroRecord, jsonEncoder)
    jsonEncoder.flush()
    
    (record.key(), writer.toString)
}

// 后续输出或处理逻辑
processedStream.foreachRDD { rdd =>
    rdd.foreach { case (key, jsonStr) =>
        println(jsonStr)
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 14:10:29