Spark 2.3使用DStream从Kafka读取Avro记录转JSON的问题求助
问题分析
你遇到的错误原因有两点:
org.apache.spark.sql.avro.AvroDeserializer确实是Spark 2.4.0及以上版本才引入的类,Spark 2.3中不存在这个类;- 该类是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
相关产品推荐
相关产品推荐

