Spark UDF序列化异常:Task not serializable问题求助
解决Spark Kafka Avro反序列化的Task Not Serializable异常
问题根因
异常核心是io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient未实现Serializable接口,而你的代码将该实例传入AvroDeserializer并在UDF中引用。Spark执行分布式任务时,需要把UDF相关对象序列化后发送到Executor节点,无法序列化的CachedSchemaRegistryClient直接导致任务失败。
解决方案
方案一:Executor端延迟初始化SchemaRegistryClient
修改AvroDeserializer类,避免直接传递CachedSchemaRegistryClient实例,改为在Executor节点上懒加载客户端:
import com.databricks.spark.avro.SchemaConverters import io.confluent.kafka.schemaregistry.client.{CachedSchemaRegistryClient, SchemaRegistryClient} import io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer import org.apache.avro.Schema import org.apache.avro.generic.GenericRecord import org.apache.spark.sql.SparkSession val topic = "topic" val kafkaUrl = "kafkaUrl" val schemaRegistryUrl = "schemaRegistryUrl" class AvroDeserializer(schemaRegistryUrl: String) extends AbstractKafkaAvroDeserializer with Serializable { // @transient标记避免序列化,懒加载保证Executor端仅初始化一次 @transient private lazy val schemaRegistryClient: SchemaRegistryClient = new CachedSchemaRegistryClient(schemaRegistryUrl, 128) override def deserialize(bytes: Array[Byte]): String = { this.schemaRegistry = schemaRegistryClient val genericRecord = super.deserialize(bytes).asInstanceOf[GenericRecord] genericRecord.toString } } val kafkaAvroDeserializer = new AvroDeserializer(schemaRegistryUrl) object DeserializerWrapper extends Serializable{ val deserializer = kafkaAvroDeserializer } spark.udf.register("deserialize", (bytes: Array[Byte]) => DeserializerWrapper.deserializer.deserialize(bytes)) val kafkaDataFrame = spark.read.format("kafka") .option("kafka.bootstrap.servers", kafkaUrl) .option("subscribe", topic) .option("startingOffsets", "earliest") .load() val valueDataFrame = kafkaDataFrame.selectExpr("deserialize(value) AS message") valueDataFrame.show()
方案二:使用Confluent官方Spark Avro库(推荐)
直接使用Confluent提供的spark-avro库,无需手动编写反序列化UDF,简化代码且彻底规避序列化问题:
添加依赖(Maven)
<dependency> <groupId>io.confluent</groupId> <artifactId>kafka-avro-serializer</artifactId> <version>7.4.0</version> <!-- 替换为你的Confluent版本 --> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>spark-avro</artifactId> <version>7.4.0_3.3.0</version> <!-- 格式:Confluent版本_Spark版本 --> </dependency>
代码实现
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.col import io.confluent.spark.avro._ val topic = "topic" val kafkaUrl = "kafkaUrl" val schemaRegistryUrl = "schemaRegistryUrl" val spark = SparkSession.builder() .appName("AvroConsumer") .getOrCreate() val kafkaDataFrame = spark.read.format("kafka") .option("kafka.bootstrap.servers", kafkaUrl) .option("subscribe", topic) .option("startingOffsets", "earliest") .load() // 使用from_avro函数直接反序列化,自动从Schema Registry获取对应Schema val valueDataFrame = kafkaDataFrame.select( from_avro(col("value"), schemaRegistryUrl, s"${topic}-value").alias("message") ) valueDataFrame.show()
内容的提问来源于stack exchange,提问作者Rambo
相关产品推荐
相关产品推荐

