如何在Palantir Foundry中对多个同Schema数据集执行Union操作
Palantir Foundry 同Schema多数据集Union合并最佳方案
针对你已经明确待合并数据集范围、且所有数据集列名、列数据类型完全一致的场景,优先选以下两种方案即可,生产环境首推代码实现。
方案1:PySpark代码Transform(生产环境首选)
这个方案性能最好、逻辑最稳定,版本可追溯,不会出现可视化操作的误改问题。
操作步骤:
- 新建Python类型的Transform数据集,把所有需要合并的数据集都添加为输入依赖
- 用Spark原生的
unionByName方法做拼接,这个方法是按列名对齐匹配,比老的union(按列位置匹配)更可靠,哪怕后续某张表调整了列顺序,也不会出现数据串列的问题 - 输出合并后的结果即可
如果待合并数据集数量少(3-5个),直接链式调用就行,示例代码:
from transforms.api import transform_df, Input, Output @transform_df( Output("合并后数据集的存储路径"), df_a=Input("数据集A的存储路径"), df_b=Input("数据集B的存储路径"), df_c=Input("数据集C的存储路径"), ) def compute(df_a, df_b, df_c): # 保留所有行,不做去重,匹配全量拼接需求 return df_a.unionByName(df_b).unionByName(df_c)
如果待合并数据集数量多(超过10个),不用手写一长串union链,用reduce批量处理更简洁:
from functools import reduce from transforms.api import transform_df, Input, Output @transform_df( Output("合并后数据集的存储路径"), input_tables={ "ds_a": Input("数据集A的存储路径"), "ds_b": Input("数据集B的存储路径"), "ds_c": Input("数据集C的存储路径"), # 后续新增待合并数据集直接在这里加对应Input就行 } ) def compute(input_tables): all_dfs = list(input_tables.values()) return reduce(lambda df1, df2: df1.unionByName(df2), all_dfs)
方案2:可视化拖拽实现(临时分析首选)
如果只是做临时探索分析,不想写代码,可以用Contour或者管道构建器拖拽实现:
- 新建Contour分析或者可视化管道任务
- 把所有待合并数据集拖入画布,因为已经确认过所有表schema一致,不需要额外做列对齐、类型转换
- 依次添加Union节点,把上游数据集接入,选择「保留所有重复行」的配置
- 最后把结果保存为目标数据集即可
注意:这种方式不适合生产调度任务,后续维护的时候很容易因为拖拽误操作改坏拼接逻辑,出问题排查难度比代码高很多。
结果验证
用提供的三个测试数据集跑上述逻辑,输出结果完全符合预期:
| col1 | col2 |
|---|---|
| 1 | a |
| 2 | b |
| 2 | c |
| 3 | d |
| 1 | e |
| 1 | f |
注意事项
- 因为所有表schema完全一致,不需要额外加列选择、类型转换的逻辑,后续如果某张表改了列名、改了字段类型,
unionByName会直接报错终止任务,不会静默生成脏数据,天然带一层数据质量校验 - 如果后续有去重需求,直接在合并结果后加
.dropDuplicates()即可,当前全量拼接的需求不需要加这步 - 不要用逐行循环插入的方式实现合并,性能比Spark原生union差非常多,数据量上十万之后很容易出现任务超时、资源占满的问题
内容的提问来源于stack exchange,提问作者domdomegg
相关产品推荐
相关产品推荐

