使用PySpark将CSV转换为Avro缺失依赖的解决方案问询
解决方法
你遇到的报错有两个核心原因:一是Spark启动时没有加载avro依赖包,二是你用的PySpark 2.4.3版本本身就没有pyspark.sql.avro模块,这个模块是Spark 3.0版本之后才正式加入PySpark官方封装的。按以下步骤修改代码即可:
1. 启动SparkSession时内置加载avro依赖
在SparkSession.builder中添加spark.jars.packages配置,Spark会自动拉取对应版本的avro依赖包,无需手动调用spark-submit传参。注意依赖包的版本要和你使用的Spark版本、Scala编译版本匹配:
- Spark 2.4.3默认使用Scala 2.11编译,对应依赖为
org.apache.spark:spark-avro_2.11:2.4.3 - 如果你的Spark是用Scala 2.12编译的,将依赖中的
2.11改为2.12即可
2. 适配Spark 2.4.3的Avro函数调用
2.4.x版本的Avro转换函数没有Python封装,直接用Spark SQL表达式调用即可,无需导入pyspark.sql.avro相关模块。
修改后的完整代码
from pyspark.sql import SparkSession from pyspark.sql.functions import expr spark = SparkSession \ .builder \ .appName("Avro testing") \ # 配置自动加载avro依赖 .config("spark.jars.packages", "org.apache.spark:spark-avro_2.11:2.4.3") \ .getOrCreate() schema = open("schema.avsc", "r").read() # 读取CSV df = spark.read.csv(path="/data/data.csv", header=True) # 用SQL表达式调用avro转换函数,和你原有逻辑保持一致 output = df\ .select(expr(f"from_avro(value, '{schema}')").alias("user"))\ .where('user.favorite_color == "red"')\ .select(expr("to_avro(user.name)").alias("value"))
补充说明
- 如果你是在内网环境无法访问公共maven仓库,可以提前把spark-avro的jar包下载到本地,把配置改为
.config("spark.jars", "/本地路径/spark-avro_2.11-2.4.3.jar")即可 - 注意你原有逻辑中直接读取CSV后取
value列是不符合常规CSV读取逻辑的:带header读取CSV后,DataFrame的列名对应CSV的表头字段,不存在value列。如果你的需求是直接将整行CSV数据按schema转Avro,需要先将整行合并为二进制字符串再调用from_avro,或者直接把CSV的列对应到schema结构后调用to_avro
内容的提问来源于stack exchange,提问作者gbeaven
相关产品推荐
相关产品推荐

