如何在GCP Data Fusion pipeline中运行代码扁平化员工父子层级关系
GCP Data Fusion 集成员工层级扁平化代码操作指南
方案1:使用SQL代码实现(推荐,性能更优)
- 首先在Data Fusion画布中添加源节点,选择你存储员工数据的对应数据源(比如Cloud SQL、BigQuery、GCS CSV等),完成数据源连接配置,确保能正常读取
员工、直属上级两个字段。 - 添加
BigQuery Execute节点(如果你的数据已经在BigQuery里,或者你可以临时把数据导入BigQuery处理),把你已有的递归层级扁平化SQL代码粘贴到节点的SQL编辑框中,注意表名替换为你源节点输出的临时表名/实际表名。如果你的数据不在BigQuery,也可以选择
Database类的执行节点,匹配你的数据源类型即可。 - SQL执行完成后,添加Sink节点(比如BigQuery、GCS等),配置输出路径存储扁平化后的层级数据即可。
方案2:使用Python代码实现
- 在Data Fusion画布中添加源节点,完成数据源配置,确保正常读取
员工、直属上级字段。 - 添加
Python Transform节点(也叫PySpark Transform节点,部分版本显示为Python Evaluator),打开节点配置页:- 在
Schema配置项中,先定义你扁平化后需要输出的所有字段(比如员工ID、员工姓名、一级上级、二级上级、直到顶层上级等) - 把你已有的层级扁平化Python代码粘贴到代码编辑区域,注意代码需要适配PySpark语法(因为Data Fusion的Python节点底层跑的是PySpark任务,不能直接用普通本地Python的pandas语法,如果你的原有代码是pandas实现的,需要修改为PySpark DataFrame的操作逻辑)
- 代码逻辑最后需要返回处理好的PySpark DataFrame,字段和你之前定义的输出Schema保持一致。
- 在
- 配置完成后连接Python节点到Sink节点,配置输出路径即可。
测试验证注意事项
- 首次运行可以先配置小批量测试数据,验证层级展开结果是否符合预期,再全量运行。
- 如果你的组织层级超过10层,建议优先使用SQL递归方案,性能比Python处理高30%以上。
内容的提问来源于stack exchange,提问作者Nadia Amina
相关产品推荐
相关产品推荐

