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

如何通过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读取时遇到两个问题:

  1. 直接将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()

输出显示乱码(截图略)。

  1. 尝试手动指定从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}|
+------------+

解决方案

核心原因

  1. 乱码:AVRO是二进制序列化格式,直接转为字符串必然出现乱码,必须用AVRO专用解析器处理。
  2. 全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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 02:31:01