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

如何在Google Cloud Platform中使用PySpark将CSV文件转换为Avro文件

GCP环境PySpark CSV转Avro可行方案

1. 前置依赖配置

PySpark默认未内置Avro格式支持,需要先匹配你的Spark版本引入对应spark-avro依赖,依赖的大版本必须和运行环境的Spark版本完全一致,Scala版本也要匹配:

  • 比如Spark 3.3.x、Scala 2.12环境对应的依赖为org.apache.spark:spark-avro_2.12:3.3.0
  • 若使用GCP Dataproc集群运行作业,可在提交作业时通过参数引入依赖,无需提前安装到集群

2. 完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType # 按需替换为实际字段类型

# 初始化SparkSession
spark = SparkSession.builder \
    .appName("CSV-to-Avro") \
    # 本地测试未预装依赖时取消注释下方配置,替换为对应版本依赖
    # .config("spark.jars.packages", "org.apache.spark:spark-avro_2.12:3.3.0") \
    .getOrCreate()

# 手动定义CSV Schema,避免自动推断的类型误差,按需调整字段和类型
csv_schema = StructType([
    StructField("id", IntegerType(), nullable=True),
    StructField("username", StringType(), nullable=True),
    StructField("trade_amount", DoubleType(), nullable=True)
])

# 读取CSV文件,支持本地路径、GCS路径(gs://bucket/xxx格式)
df = spark.read.csv(
    path="gs://你的存储桶/输入CSV路径/",
    schema=csv_schema,
    header=True, # CSV第一行是表头则设为True,否则设为False
    sep=",", # 按实际分隔符调整,制表符为"\t"
    encoding="utf-8"
)

# 写入Avro文件
df.write.format("avro") \
    .mode("overwrite") # 写入模式可选:overwrite/append/ignore/errorifexists
    .option("compression", "snappy") # 可选压缩格式:uncompressed/snappy/deflate/bzip2
    .save("gs://你的存储桶/输出Avro路径/")

3. 常见问题处理

  • 若使用Dataproc提交作业,可直接通过命令指定依赖,无需在代码中配置:
gcloud dataproc jobs submit pyspark 你的脚本名.py \
    --cluster=你的集群名称 \
    --region=集群所属区域 \
    --properties=spark.jars.packages=org.apache.spark:spark-avro_2.12:3.3.0
  • 读写GCS路径报错时,确认运行作业的服务账号已开通对应存储桶的读写权限
  • 写入Avro时报类型错误,提前对CSV数据做清洗,确保字段值和定义的Schema类型匹配

内容的提问来源于stack exchange,提问作者Jopsiton

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 12:24:03