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

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解析:

  1. 修改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
  1. 调整代码使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 09:25:36