在Azure Synapse Analytics中用PySpark拆分CSV列并转GB单位
PySpark处理ADLS CSV:拆分列+存储单位转GB(覆盖原文件)
针对你的需求,以下是可直接使用的PySpark代码,附带详细注释,适合Azure Synapse新手:
1. 读取ADLS中的CSV文件
替换代码中的<存储账户名>、<容器名>、<文件路径>为你实际的信息:
# 导入Spark函数库 from pyspark.sql.functions import split, col, when, lit # 读取CSV(abfss是ADLS Gen2的标准访问协议) df = spark.read.csv( "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/<文件路径>/target_file.csv", header=True, # 指明CSV首行为表头 inferSchema=True # 自动推断列数据类型 )
2. 拆分第二列
假设第二列的表头是combined_col(替换成你实际的表头),这里以空格作为分隔符拆分(如果是其他分隔符,比如|,改成split(col("combined_col"), "\\|")):
# 拆分第二列为两个新列,示例命名为`size_value`和`unit` df = df.withColumn("size_value", split(col("combined_col"), " ")[0]) \ .withColumn("unit", split(col("combined_col"), " ")[1]) \ .drop("combined_col") # 删除原第二列
3. 转换存储单位为GB
处理常见存储单位(B/KB/MB/GB/TB),默认采用二进制换算(1GB=1024MB),如果需要十进制(1GB=1000MB),把代码中的1024替换为1000即可:
# 计算转换为GB后的数值(先确保size_value是数值类型,若为字符串则用cast("double")转换) df = df.withColumn( "size_gb", when(col("unit") == "TB", col("size_value").cast("double") * 1024) \ .when(col("unit") == "GB", col("size_value").cast("double")) \ .when(col("unit") == "MB", col("size_value").cast("double") / 1024) \ .when(col("unit") == "KB", col("size_value").cast("double") / (1024**2)) \ .when(col("unit") == "B", col("size_value").cast("double") / (1024**3)) \ .otherwise(lit(None)) # 未知单位返回空值 ) # 可选:删除原始的大小和单位列(如果不需要保留) df = df.drop("size_value", "unit")
4. 覆盖原文件写入ADLS
Spark默认会生成多个分区文件,用coalesce(1)合并为单个文件,再覆盖原路径:
# 写入CSV并覆盖原文件 df.coalesce(1) \ .write \ .mode("overwrite") \ .option("header", True) \ .csv("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/<文件路径>/target_file.csv")
关键注意事项
- 权限配置:确保Azure Synapse工作区的托管标识(MSI)在ADLS容器的访问控制中拥有
存储Blob数据参与者权限,否则会出现访问拒绝错误。 - 分隔符适配:如果第二列的分隔符不是空格,必须调整
split函数中的分隔符(特殊字符如|需要用\\|转义)。 - 数据类型校验:运行代码前可以用
df.printSchema()查看列类型,确保数值列的类型正确,避免转换错误。
内容的提问来源于stack exchange,提问作者Eli
相关产品推荐
相关产品推荐

