使用Spark结合AWS Glue读取Confluent编码Avro记录时出现格式错误
根本原因
Confluent Avro编码的消息在实际Avro二进制 payload 前固定携带5字节前缀:1字节魔术位(固定值为0)+ 4字节大端序存储的Schema ID,Spark原生from_avro函数默认解析无任何前缀的纯Avro二进制数据,忽略该前缀会直接导致解析失败。
可行解决方案
方案1:高版本Spark直接启用Confluent兼容配置
Spark 3.2及以上版本的from_avro原生支持Confluent Avro格式,只需在参数中补充Schema Registry地址即可自动识别前缀、拉取对应Schema解析:
output = df.select( from_avro( "value", options={ "mode":"PERMISSIVE", "confluent.schema.registry.url": "你的Confluent Schema Registry访问地址", # 如Schema Registry配置了SSL认证,可在此补充对应的ssl配置参数 } ).alias("data") )
如果使用AWS Glue内置Schema Registry,可直接使用Glue官方适配的解析逻辑,无需手动处理前缀:
df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "host1:port1") \ .option("subscribe", "topic") \ .option("kafka.group.id", group_id) \ .option("kafka.client.id", client_id) \ .option("kafka.ssl.truststore.location", truststore) \ .option("kafka.ssl.truststore.password", password) \ .option("kafka.ssl.keystore.location", keystore) \ .option("kafka.ssl.keystore.password", password) \ .option("startingOffsets", "earliest") \ .option("endingOffsets", "latest") \ .load() # 使用Glue Schema Registry解析Confluent Avro output = df.select( from_avro( "value", schema_registry_options={ "schemaName": "你的Schema名称", "registryName": "你的Glue Schema Registry名称", "region": "AWS区域" } ).alias("data") )
方案2:手动剥离5字节前缀(无需连接Schema Registry)
如果已经提前获取到对应Avro Schema,不想额外对接Schema Registry,可直接截取value字段从第6位开始的内容,再传入from_avro解析:
from pyspark.sql.functions import col, substring, length output = df.select( from_avro( substring(col("value"), 6, length(col("value")) - 5), jsonFormatSchema, options={"mode":"PERMISSIVE"} ).alias("data") )
注意:该方案仅适用于对应Kafka topic所有消息都使用同一个固定Schema的场景,Schema更新后需要同步修改代码中的jsonFormatSchema。
内容的提问来源于stack exchange,提问作者Ricardo Nunes
相关产品推荐
相关产品推荐

