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

使用Kafka ByteArrayDeserializer读取Avro消息时遇ClassCastException求助

问题分析与解决

核心错误原因

  1. 消费者泛型类型不匹配:你配置了value.deserializer为ByteArrayDeserializer,但创建createDirectStream时指定的泛型是[String, String],这会导致Kafka返回的字节数组被强制转成String,直接触发类型转换异常。
  2. 冗余的getBytes()调用:即使泛型正确,record.value()已经是Array[Byte],不需要再调用getBytes(),这会把字节数组当成String处理,同样引发错误。

修正后的代码

val ssc = new StreamingContext(spark.sparkContext, Seconds(1))

val kafkaParams: Map[String, Object] = Map(
  "bootstrap.servers" -> "kafka-server:9092",
  "key.serializer" -> classOf[StringSerializer],
  "value.serializer" -> classOf[StringSerializer],
  "key.deserializer" -> classOf[StringDeserializer],
  "value.deserializer" -> classOf[ByteArrayDeserializer],
  "auto.offset.reset" -> "earliest",
  "enable.auto.commit" -> (false: java.lang.Boolean),
  "security.protocol" -> "SSL",
  "ssl.truststore.location" -> "truststore",
  "ssl.truststore.password" -> "pass",
  "ssl.keystore.location" -> "keystore.jks",
  "ssl.keystore.password" -> "pass",
  "group.id" -> "group1"
)

val topics: Array[String] = Array("topics")

// 修正泛型类型为[String, Array[Byte]]
val kafkaDstream = KafkaUtils.createDirectStream(
  ssc,
  LocationStrategies.PreferConsistent,
  ConsumerStrategies.Subscribe[String, Array[Byte]](topics, kafkaParams)
)

val schema = parser.parse(new String(Files.readAllBytes(Paths.get("avro2.avsc"))))
val datumReader = new SpecificDatumReader[GenericRecord](schema)

val processedStream = kafkaDstream.map(record => {
  // 直接使用record.value()作为字节数组,无需getBytes()
  val x = new ByteArrayInputStream(record.value())
  val binaryDecoder = DecoderFactory.get.binaryDecoder(x, null)
  datumReader.read(null, binaryDecoder)
})

processedStream.map(rec => rec.get("taskId")).print

额外建议

  • 如果你的Kafka主题中的Avro记录是带Schema Registry的,建议使用io.confluent.kafka.serializers.KafkaAvroDeserializer,无需手动解析Schema,配置更简便。
  • 考虑迁移到Spark结构化流(Structured Streaming),它对Kafka和Avro的支持更完善,API也更简洁稳定。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 17:55:19