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

在Foundry中检查导入.parq文件Schema并生成列数统计数据框DF2

在Foundry中识别Parquet文件的Schema差异(列数维度)

解决方案代码(Foundry Transform 版本)

以下代码通过Foundry Transform实现需求,自动遍历目标数据集的所有Parquet文件,记录文件路径和对应列数到新数据集DF2:

import pyspark.sql.functions as F
from transforms.api import transform, Input, Output

@transform(
    output=Output("/your/output/path/DF2"),  # 替换为DF2的存储路径
    input_data=Input("/your/input/path/DF1")  # 替换为DF1的路径
)
def extract_schema_info(ctx, input_data, output):
    # 1. 生成所有Parquet文件的路径列表
    file_paths = [
        file.path for file in input_data.files() 
        if file.path.endswith((".parq", ".parquet"))
    ]
    
    # 2. 遍历文件并收集Schema信息
    schema_records = []
    for file_path in file_paths:
        # (a) 读取单个文件的Schema
        single_file_df = ctx.spark.read.parquet(f"{input_data.path}/{file_path}")
        # (b) 记录路径和列数
        schema_records.append((file_path, len(single_file_df.schema.fields)))
    
    # 3. 创建并输出结果数据框DF2
    df2 = ctx.spark.createDataFrame(schema_records, ["file_path", "column_count"])
    output.write_dataframe(df2)

Foundry Notebook 版本(交互式验证)

如果需要在Notebook中快速验证,可使用以下代码:

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

# 替换为你的DF1路径
input_dataset_path = "/your/input/path/DF1"
# 获取所有Parquet文件路径
file_paths = [
    file.path for file in spark._jvm.com.palantir.foundry.spark.dataset.Datasets.files(input_dataset_path)
    if file.path.endswith((".parq", ".parquet"))
]

# 收集Schema信息
schema_info = []
for path in file_paths:
    df = spark.read.parquet(f"{input_dataset_path}/{path}")
    schema_info.append((path, len(df.schema.fields)))

# 创建DF2并显示结果
df2 = spark.createDataFrame(schema_info, ["file_path", "column_count"])
df2.show(truncate=False)

# 可选:将DF2保存到Foundry
df2.write.mode("overwrite").parquet("/your/output/path/DF2")

扩展:识别更精细的Schema差异

如果需要对比完整的Schema(不仅是列数),可以扩展代码记录字段名和类型:

def schema_to_detail_string(schema):
    """将Schema转换为可对比的字符串格式"""
    return str([(field.name, field.dataType.simpleString()) for field in schema.fields])

# 在遍历中替换为:
single_file_schema = single_file_df.schema
schema_records.append((
    file_path, 
    len(single_file_schema.fields), 
    schema_to_detail_string(single_file_schema)
))

# 对应的DF2结构为:["file_path", "column_count", "schema_details"]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 20:06:24