如何通过PySpark正确格式化Kafka Topic的AVRO数据?
问题:Spark读取Kafka中KSQL写入的AVRO数据异常
问题背景
我通过KSQL向Kafka Topic写入AVRO格式数据:
CREATE STREAM TEST01 (KEY_COL VARCHAR KEY, COL1 INT, COL2 VARCHAR) WITH (KAFKA_TOPIC='test01', PARTITIONS=1, VALUE_FORMAT='AVRO'); INSERT INTO TEST01 (KEY_COL, COL1, COL2) VALUES ('X',1,'FOO'); INSERT INTO TEST01 (KEY_COL, COL1, COL2) VALUES ('Y',2,'BAR');
使用PySpark读取时遇到两个问题:
- 直接将
value转为字符串输出,结果出现乱码:
from pyspark.sql.session import SparkSession spark = SparkSession \ .builder \ .appName("Kafka_Test") \ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0") \ .getOrCreate() df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "test01") \ .option("startingOffsets","earliest") \ .load() df.selectExpr("cast(value as string) as value").writeStream.outputMode("append").format("console").start()
输出显示乱码(截图略)。
- 尝试手动指定从Schema Registry获取的Schema解析,结果得到全null记录:
jsonSchema={"type":"record","name":"KsqlDataSourceSchema","namespace":"io.confluent.ksql.avro_schemas","fields":[{"name":"COL1","type":["null","int"],"default":None},{"name":"COL2","type":["null","string"],"default":None}],"connect.name":"io.confluent.ksql.avro_schemas.KsqlDataSourceSchema"} df.select(from_avro("value", json.dumps(jsonSchema)).alias("sample_a")) \ .writeStream.format("console").start()
输出:
+------------+ | sample_a| +------------+ |{null, null}| |{null, null}| |{null, null}| +------------+
解决方案
核心原因
- 乱码:AVRO是二进制序列化格式,直接转为字符串必然出现乱码,必须用AVRO专用解析器处理。
- 全null解析结果:KSQL写入的AVRO数据采用Confluent扩展的AVRO格式(包含1字节魔术位+4字节Schema ID前缀),而Spark原生
from_avro只支持标准AVRO格式,无法识别Confluent的前缀,导致解析失败。
具体解决步骤
1. 添加Confluent Spark AVRO依赖
Spark需要额外依赖来处理Confluent格式的AVRO数据,在构建SparkSession时添加对应包:
spark = SparkSession \ .builder \ .appName("Kafka_Test") \ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0,io.confluent:spark-avro_2.12:7.4.0") \ .getOrCreate()
注意:
spark-avro版本需与Confluent平台版本匹配,示例用7.4.0,可根据实际环境调整。
2. 配置Schema Registry地址
告诉Spark从哪里获取AVRO Schema:
spark.conf.set("spark.sql.avro.confluent.schema.registry.url", "http://localhost:8081")
替换为实际的Schema Registry地址和端口。
3. 使用Confluent格式解析数据
直接通过Topic对应的Schema Subject读取解析,无需手动指定Schema:
from pyspark.sql.avro.functions import from_avro # 通过Topic对应的Schema Subject解析(KSQL自动注册的Subject格式为<topic-name>-value) df_parsed = df.select( from_avro("value", "test01-value", options={"mode": "PERMISSIVE"}).alias("data") ) # 展开解析后的字段并输出 df_parsed.select("data.*").writeStream \ .outputMode("append") \ .format("console") \ .start() \ .awaitTermination()
4. 预期输出
运行后将正确解析出数据:
+----+----+ |COL1|COL2| +----+----+ | 1| FOO| | 2| BAR| +----+----+
额外说明
- 若需手动指定Schema,需先处理Confluent格式的前缀字节,不推荐此方式,优先直接从Schema Registry获取。
- 确保Spark、Kafka、Confluent Schema Registry的版本兼容,避免因版本冲突导致解析失败。
内容的提问来源于stack exchange,提问作者RushHour
相关产品推荐
相关产品推荐

