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

PySpark Structured Streaming写入Kafka时Avro Schema验证失败问题

问题描述

尝试用PySpark Structured Streaming生成Avro编码消息发送到Kafka,Schema已在Confluent Schema Registry注册,但开启Schema验证时触发InvalidRecordException异常。

代码:

from pyspark.sql import SparkSession
from pyspark.sql.functions import struct, col
from pyspark.sql.avro.functions import to_avro
from pyspark.sql.types import StructType, StructField, StringType

# Create SparkSession
spark = SparkSession.builder \
    .appName("KafkaAvroSchemaProducer") \
    .getOrCreate()

# Define the schema for the DataFrame
schema = StructType([
    StructField("name", StringType(), True)
])

# Define sample message
sample_message = {"name": "John"}

# Create DataFrame with the sample message
sample_df = spark.createDataFrame([sample_message], schema=schema)

# Convert DataFrame to Avro format
avro_df = sample_df.select(
    to_avro(struct("*")).alias("value")  # Convert to Avro format
).selectExpr("CAST(NULL AS STRING) AS key", "value")

# Kafka producer configuration
producer_config = {
    'kafka.bootstrap.servers': 'your_kafka_broker:9092',
    'topic': 'sample_topic',
    'schema.registry.url': 'http://your_schema_registry:8081'
}

# Write DataFrame to Kafka in Avro format
avro_df.write \
    .format("kafka") \
    .option("kafka.bootstrap.servers", producer_config['kafka.bootstrap.servers']) \
    .option("topic", producer_config['topic']) \
    .option("key.serializer", "org.apache.kafka.common.serialization.StringSerializer") \
    .option("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer") \
    .option("schema.registry.url", producer_config['schema.registry.url']) \
    .save()

# Stop the Spark session
spark.stop()

Schema Registry中的Schema:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "name", "type": "string"}
  ]
}

异常信息:

InvalidRecordException: The record is rejected by the record interceptor

提问:使用to_avro序列化数据为Avro格式的方式是否存在问题?是否需要额外配置才能通过Schema验证?


解决方法

核心问题是同时混用了PySpark内置的to_avro函数和Confluent的KafkaAvroSerializer,导致消息格式不匹配,触发Schema验证拦截。

问题拆解

  1. PySpark的to_avro会直接生成原始Avro二进制数据,但Confluent的KafkaAvroSerializer要求的是带有Schema ID的Confluent封装格式(即先从Registry获取Schema ID,再将ID与原始Avro数据打包)。
  2. 用to_avro序列化后再交给KafkaAvroSerializer,会导致二次封装,最终消息结构完全不符合Registry中注册的Schema,被验证拦截。

正确方案二选一:

方案1:用PySpark原生Avro格式写入(无需Confluent序列化器)

指定Registry地址,让Spark直接生成符合Confluent格式的Avro消息,可手动指定已注册的Schema ID或允许自动注册:

# 替换原有的avro_df转换和写入逻辑
avro_df = sample_df.selectExpr("CAST(NULL AS STRING) AS key", "struct(*) AS value")

avro_df.write \
    .format("avro") \
    .option("kafka.bootstrap.servers", producer_config['kafka.bootstrap.servers']) \
    .option("kafka.topic", producer_config['topic']) \
    .option("avro.schema.registry.url", producer_config['schema.registry.url']) \
    # 若Schema已注册,指定对应ID
    .option("avro.schema.id", "你的Schema ID") \
    # 若Registry允许自动注册,可指定subject(默认是{topic}-value)
    # .option("avro.schema.registry.subject", "sample_topic-value")
    .save()

方案2:依赖Confluent序列化器处理(不用PySpark的to_avro)

通过自定义UDF将数据转为Confluent Avro格式,再写入Kafka:

from pyspark.sql.functions import udf, lit
from pyspark.sql.types import BinaryType
from io.confluent.kafka.serializers.KafkaAvroSerializer import KafkaAvroSerializer
from io.confluent.kafka.serializers.AbstractKafkaAvroSerDeConfig import AbstractKafkaAvroSerDeConfig

# 定义序列化UDF
def serialize_to_confluent_avro(record):
    serializer_config = {
        AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG: producer_config['schema.registry.url']
    }
    serializer = KafkaAvroSerializer(serializer_config)
    # 确保record结构与Registry中的Schema完全匹配
    return serializer.serialize(f"{producer_config['topic']}-value", record.asDict())

serialize_udf = udf(serialize_to_confluent_avro, BinaryType())

# 转换DataFrame
avro_df = sample_df.select(
    lit(None).cast(StringType()).alias("key"),
    serialize_udf(struct("*")).alias("value")
)

# 写入Kafka(无需指定value.serializer,已处理为二进制)
avro_df.write \
    .format("kafka") \
    .option("kafka.bootstrap.servers", producer_config['kafka.bootstrap.servers']) \
    .option("topic", producer_config['topic']) \
    .save()

额外注意事项

  • 确保Spark依赖包含spark-avro和kafka-avro-serializer jar包,版本需与Kafka、Schema Registry一致。
  • 若Registry开启严格验证,需保证消息结构与注册的Schema完全兼容(字段名、类型一致)。
  • 确认Registry的subject名称正确,默认规则为{topic}-value,若Schema注册在其他subject下需手动指定。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 20:35:03