如何借助Fixer.io API与Cloud Fusion实现Postgres到BigQuery的货币转换数据迁移
整体实现方案总览
全流程可以在你已部署好的Cloud Fusion实例中通过定时批处理流水线实现,核心链路为「PostgreSQL数据同步→Fixer API当日汇率拉取→金额字段统一转换→BigQuery落地」,适配逐表迁移、每日汇率更新的需求。
步骤1:配置PostgreSQL数据源同步
- 首先在Cloud Fusion中创建每日定时触发的批处理流水线,触发时间设置在业务低峰期即可
- 调用内置的
PostgreSQL Batch Source插件,完成数据库连接配置,可按需为每张待迁移表配置独立源节点,也可以通过动态参数实现批量逐表拉取 - 同步后先做基础清洗:过滤金额字段为空、币种标识非法的无效数据,避免影响后续转换逻辑
步骤2:集成Fixer API拉取当日汇率
在流水线中新增汇率拉取节点,全流程每日仅需调用1次API即可,不需要逐行数据重复调用,避免触发接口限流:
- 调用Cloud Fusion内置的
HTTP插件,填入你的Fixer API密钥,指定基准币种为USD、查询币种包含EUR,拉取当日的EUR兑USD汇率值 - 拉取到的汇率结果可以暂存为流水线全局变量,也可以写入一个独立的BigQuery汇率中间表,方便所有待迁移表的转换逻辑复用同一汇率值
注意:如果Fixer API调用失败,可以配置3次重试逻辑,重试仍失败的话自动读取前1日的有效汇率兜底,避免流水线整体中断
步骤3:汇率转换逻辑实现
- 用Cloud Fusion的
Wrangler节点或者Spark SQL Execute节点实现转换,核心逻辑:如果币种标识为USD则直接保留原始金额,如果为EUR则用原始金额 * 当日EUR兑USD汇率 - 建议新增单独的
amount_usd字段存储转换后的统一USD金额,同时保留原始金额、原始币种字段,方便后续对账和逻辑回溯 - 如果有多张表需要做相同转换,可以把上述逻辑封装为通用UDF,所有表节点直接复用即可,不需要重复开发
步骤4:数据写入BigQuery落地
- 调用
BigQuery Sink插件,配置好目标数据集、对应表名,建议按日分区存储,方便后续数据回溯 - 写入前可新增校验步骤:对比原始总金额、转换后总金额的差值,确认汇率转换逻辑没有出现计算错误
额外优化建议
- 可以把每日拉取的历史汇率都存储到专门的汇率表中,如果后续需要回溯历史数据的转换逻辑,不需要重新调用Fixer API拉取历史汇率
- 如果后续新增其他币种,只需要修改Fixer API的查询币种参数,同时更新转换逻辑即可,整体架构不需要做大的调整
内容的提问来源于stack exchange,提问作者user14380696
相关产品推荐
相关产品推荐

