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验证拦截。
问题拆解
- PySpark的
to_avro会直接生成原始Avro二进制数据,但Confluent的KafkaAvroSerializer要求的是带有Schema ID的Confluent封装格式(即先从Registry获取Schema ID,再将ID与原始Avro数据打包)。 - 用
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-serializerjar包,版本需与Kafka、Schema Registry一致。 - 若Registry开启严格验证,需保证消息结构与注册的Schema完全兼容(字段名、类型一致)。
- 确认Registry的subject名称正确,默认规则为
{topic}-value,若Schema注册在其他subject下需手动指定。
内容的提问来源于stack exchange,提问作者yAsH
相关产品推荐
相关产品推荐

