PySpark from_avro是否兼容Kafka Avro的Schema Registry魔数?解析报错怎么解决?
问题:PySpark解析Confluent Schema Registry的Avro流数据报错
场景描述
Kafka中存储的是Avro格式的流数据,我通过Confluent Schema Registry管理数据的Schema。尝试使用PySpark拉取数据,并借助Schema Registry中的Schema解析Avro字节数据时,持续抛出解析错误。
相关代码
import json import os from schema_registry.client import SchemaRegistryClient from pyspark.sql.avro.functions import from_avro os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.1,io.delta:delta-core_2.12:1.1.0,org.apache.spark:spark-avro_2.12:3.2.1 --conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension --conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog pyspark-shell' topic_name = "my_topic" bootstrap_server = "pkc-XXXXy.us-east-1.aws.confluent.cloud:9092" kafka_options = { "kafka.sasl.jaas.config": 'org.apache.kafka.common.security.plain.PlainLoginModule required username="MY_USERNAME" password="MY_PASSWORD";', "kafka.sasl.mechanism": "PLAIN", "kafka.security.protocol" : "SASL_SSL", "kafka.bootstrap.servers": bootstrap_server, "group.id": "group_wow2222", "subscribe": topic_name, "startingOffsets": "earliest", "maxOffsetsPerTrigger": 5, } log_streaming_df = spark \ .readStream \ .format("kafka") \ .options(**kafka_options) \ .load() sr_client = SchemaRegistryClient( { "url": "https://psrc-XXXX.us-east-2.aws.confluent.cloud", "basic.auth.credentials.source": "USER_INFO", "basic.auth.user.info": "MY_USERNAME_FOR_SR:MY_PW_FOR_SR" } ) schema_obj = sr_client.get_schema(f"{topic_name}-value", version="latest") log_streaming_df = log_streaming_df.select( from_avro( "value", json.dumps(schema_obj.schema.raw_schema) ) ).alias("logs") log_streaming_df.writeStream \ .format("console") \ .option("checkpointLocation", "./my_checkpoint") \ .outputMode("append") \ .start()
报错信息
Aborting task (0 + 1) / 2] org.apache.spark.SparkException: Malformed records are detected in record parsing. Current parse Mode: FAILFAST. To process malformed records as null result, try setting the option 'mode' as 'PERMISSIVE'. at org.apache.spark.sql.avro.AvroDataToCatalyst.nullSafeEval(AvroDataToCatalyst.scala:113) at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source) at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43) at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1.hasNext(WholeStageCodegenExec.scala:759) at org.apache.spark.sql.execution.datasources.v2.DataWritingSparkTask$.$anonfun$run$1(WriteToDataSourceV2Exec.scala:412) at org.apache.spark.util.Utils$.tryWithSafeFinallyAndFailureCallbacks(Utils.scala:1496) at org.apache.spark.sql.execution.datasources.v2.DataWritingSparkTask$.run(WriteToDataSourceV2Exec.scala:457) at org.apache.spark.sql.execution.datasources.v2.V2TableWriteExec.$anonfun$writeWithV2$2(WriteToDataSourceV2Exec.scala:358) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90) at org.apache.spark.scheduler.Task.run(Task.scala:131) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:506) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1462) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:509) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.avro.AvroRuntimeException: Malformed data. Length is negative: -1 at org.apache.avro.io.BinaryDecoder.readString(BinaryDecoder.java:308) at org.apache.avro.io.ResolvingDecoder.readString(ResolvingDecoder.java:208) at org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:469) at org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:459) at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:191) at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:160) at org.apache.avro.generic.GenericDatumReader.readField(GenericDatumReader.java:259) at org.apache.avro.generic.GenericDatumReader.readRecord(GenericDatumReader.java:247) at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:179) at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:160) at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:153) at org.apache.spark.s
验证测试
将流数据先保存为Parquet格式到磁盘,再读取为静态DataFrame,去除前5字节(1字节魔数+4字节Schema Registry的Schema ID)后使用fastavro解析可成功,相关代码如下:
df = spark.read.parquet("./temp/*.parquet") # 提前从Kafka流中保存的数据 import pandas as pd from fastavro import schemaless_reader, parse_schema from io import BytesIO from pyspark.sql.functions import pandas_udf, col from pyspark.sql.types import StringType schema = parse_schema(schema_obj.schema.raw_schema) @pandas_udf(StringType()) def simple_udf(v: pd.Series) -> pd.Series: return v.apply(lambda x: json.dumps(schemaless_reader(BytesIO(x[5:]), schema))) df = df.withColumn("real_value", simple_udf(col("value")))
核心疑问
推测PySpark的from_avro未处理该5字节前缀导致解析报错,请问是否正确?如果正确,无需使用pandas_udf的简易解决方法是什么?
解答
你的推测完全正确。Confluent的Avro序列化格式会在原始Avro数据前添加1字节魔数(固定为0x00)+4字节Schema ID的前缀,而PySpark原生的from_avro仅支持标准Avro格式,不会处理这个Confluent扩展的前缀,因此直接解析会抛出格式错误。
以下是两种无需使用pandas_udf的简易解决方法:
方法1:使用Confluent官方Spark Avro库
Confluent提供了适配其Schema Registry的Spark工具库,可自动处理前缀并完成Schema解析:
- 修改PySpark提交参数,替换原生avro包为Confluent的适配包(版本需与Confluent平台版本匹配,示例用7.3.0):
--packages io.confluent:kafka-avro-serializer:7.3.0,io.confluent:kafka-schema-registry-client:7.3.0,org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.1,io.delta:delta-core_2.12:1.1.0
- 调整代码使用Confluent的反序列化器:
import json from confluent.kafka.schema_registry.avro import AvroDeserializer from confluent.kafka.serialization import SerializationContext, MessageField from pyspark.sql.functions import udf from pyspark.sql.types import StructType # 初始化Confluent Avro反序列化器 avro_deserializer = AvroDeserializer( schema_registry_client=sr_client, schema_str=json.dumps(schema_obj.schema.raw_schema) ) # 定义普通UDF处理反序列化 def deserialize_avro(value): if value is None: return None try: return avro_deserializer(value, SerializationContext(topic_name, MessageField.VALUE)) except Exception: return None # 生成适配Spark数据类型的UDF deserialize_udf = udf(deserialize_avro, StructType.fromJson(schema_obj.schema.raw_schema)) # 解析Kafka消息 log_streaming_df = log_streaming_df.select( deserialize_udf("value").alias("logs") )
方法2:手动裁剪前缀后用原生from_avro解析
如果不想引入额外依赖,可直接用Spark的substring函数裁剪掉前5字节前缀,再用原生from_avro解析:
from pyspark.sql.functions import substring, col from pyspark.sql.avro.functions import from_avro # 裁剪前5字节(Spark的substring为1索引,从第6位开始取到末尾) trimmed_df = log_streaming_df.withColumn( "trimmed_value", substring(col("value"), 6, -1) ) # 用原生from_avro解析裁剪后的标准Avro数据 log_streaming_df = trimmed_df.select( from_avro( "trimmed_value", json.dumps(schema_obj.schema.raw_schema), mode="PERMISSIVE" # 可选:容错模式,解析失败时返回null ).alias("logs") )
内容的提问来源于stack exchange,提问作者user3595632
相关产品推荐
相关产品推荐

