如何在Databricks中用PySpark按API返回Schema生成Parquet文件
在Databricks中基于API返回的Schema将CSV转换为Parquet文件(PySpark实现)
核心思路
先读取CSV得到全字符串类型的DataFrame,再根据API返回的Schema定义,将每一列转换为目标数据类型,最后写入Parquet格式文件。
步骤实现
1. 解析API返回的Schema字符串为PySpark DataType
假设API返回的Schema是键值对结构(比如JSON格式),示例如下:
api_schema_response = { "emp_name": "string(50)", "emp_salary": "decimal(7,4)", "joining_date": "timestamp" }
编写函数将类型字符串转换为PySpark对应的DataType:
from pyspark.sql.types import StringType, DecimalType, TimestampType, DataType def parse_schema_type(type_str: str) -> DataType: type_str = type_str.lower().strip() # 处理string类型(PySpark StringType不限制长度,忽略括号内的数值) if type_str.startswith("string"): return StringType() # 处理decimal类型,提取精度和小数位 elif type_str.startswith("decimal"): precision_scale = type_str.replace("decimal(", "").replace(")", "").split(",") precision = int(precision_scale[0].strip()) scale = int(precision_scale[1].strip()) return DecimalType(precision=precision, scale=scale) # 处理timestamp类型 elif type_str == "timestamp": return TimestampType() # 可扩展支持int、date等其他类型 else: raise ValueError(f"不支持的数据类型: {type_str}")
2. 读取CSV文件得到全字符串DataFrame
在Databricks中读取CSV,强制所有列为string类型:
# 根据实际情况调整分隔符、表头参数 raw_df = spark.read.csv( path="/dbfs/path/to/your/input.csv", header=True, inferSchema=False, sep="," )
3. 转换DataFrame列类型
遍历API返回的Schema,逐个转换列的类型:
from pyspark.sql.functions import col converted_df = raw_df for col_name, type_str in api_schema_response.items(): target_type = parse_schema_type(type_str) # 用try_cast替代cast可避免转换失败导致任务中断 converted_df = converted_df.withColumn( col_name, col(col_name).try_cast(target_type) ) # 可选:验证转换后的Schema是否符合预期 converted_df.printSchema()
4. 写入Parquet文件
将转换后的DataFrame写入指定存储路径(支持DBFS或云存储):
converted_df.write.parquet( path="/dbfs/path/to/your/output.parquet", mode="overwrite", # 根据需求选择overwrite/append/ignore等模式 compression="snappy" # 常用压缩格式,可根据需要调整 )
注意事项
- 若API返回的Schema不是字典格式,需先将其解析为可遍历的键值对结构(比如从JSON字符串解析)
- 对于需要截断过长字符串的场景,可以在类型转换前添加字符串截断逻辑
- 转换后建议抽样查看数据,确认类型转换结果符合预期
内容的提问来源于stack exchange,提问作者Hillol Saha
相关产品推荐
相关产品推荐

