如何在Data Fusion流水线中实现基于控制表的迭代执行?
在Google Cloud Data Fusion中实现基于控制表行参数的批量BigQuery查询并追加到临时表
核心实现步骤
你已经完成了第一步(从Control_Table取数并设为行参数),接下来按以下配置完成流水线:
添加Loop组件实现逐行遍历
- 将第一个BQ Execute的输出端口连接到Loop组件的输入端口。
- 无需额外配置参数映射,只要Control_Table的字段名(
system_name、start_date、end_date)和后续SQL的参数名一致,Loop会自动将每一行的字段值传递为对应参数。
配置Loop内部的BQ Execute组件
- 在Loop组件内部添加第二个BQ Execute节点,编写参数化SQL语句:
SELECT col_a, col_b from dw_table WHERE system_name = ${system_name} and transaction_date between ${start_date} and ${end_date}; - 确保该组件的参数引用和Control_Table的字段名完全匹配,Data Fusion会自动替换每个循环迭代中的参数值。
- 在Loop组件内部添加第二个BQ Execute节点,编写参数化SQL语句:
实现结果追加到临时表
有两种可选方案:- 方案一:用Merge组件统一收集后写入
- 在Loop外部添加Merge组件,将Loop内部BQ Execute的输出连接到Merge的输入端口。
- Merge会汇总所有循环产生的结果集,再连接到BigQuery Write组件,配置写入模式为Append,目标表指定为你的临时表。
- 方案二:Loop内部直接追加写入
- 在Loop内部的BQ Execute之后直接连接BigQuery Write组件,配置写入模式为Append,目标表设为临时表。这种方式会在每次循环迭代后直接追加数据到临时表,无需额外Merge。
- 方案一:用Merge组件统一收集后写入
关键注意事项
- 临时表需提前创建,或在BigQuery Write中开启“自动创建表”选项(需确保输出字段
col_a、col_b的类型和临时表结构一致)。 - 确认Control_Table中
start_date、end_date的字段类型与dw_table.transaction_date的类型匹配,避免SQL执行时的类型错误。 - 测试阶段建议在Control_Table中仅保留1-2行数据,验证参数替换和数据写入是否正常后再扩大数据量。
内容的提问来源于stack exchange,提问作者Ravi
相关产品推荐
相关产品推荐

