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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 15:36:08