如何在BigQuery中仅于行新增或更新时提取JSON至不同列
BigQuery增量提取JSON列(仅处理新增/更新行)
要实现只处理新增或更新的行、避免全表扫描来降低成本,核心思路是通过增量过滤+变更追踪缩小处理范围,下面是几种实用方案:
1. 先给源表加追踪时间戳
如果源表还没有记录行变更时间的字段,先添加一个last_updated列,确保每次插入/更新时自动更新:
-- 添加时间戳列 ALTER TABLE `your-project.your-dataset.source_table` ADD COLUMN last_updated TIMESTAMP DEFAULT CURRENT_TIMESTAMP() NOT NULL; -- 初始化现有行的时间戳(仅执行一次) UPDATE `your-project.your-dataset.source_table` SET last_updated = CURRENT_TIMESTAMP() WHERE last_updated IS NULL;
2. 增量处理的具体实现
方案一:定时增量查询(适合批量场景)
每次运行任务时,只处理上次任务之后变更的行,需要维护一个记录上次运行时间的控制表:
-- 读取上次运行时间(假设控制表已创建,结构为CREATE TABLE control_table(last_run_time TIMESTAMP);) DECLARE last_run_time TIMESTAMP; SET last_run_time = (SELECT last_run_time FROM `your-project.your-dataset.control_table` LIMIT 1); -- 用MERGE同步增量行到目标表(支持插入新行+更新已有行) MERGE INTO `your-project.your-dataset.target_table` t USING ( SELECT id, -- 源表主键,用于匹配目标表行 JSON_VALUE(json_column, '$.field1') AS extracted_field1, JSON_VALUE(json_column, '$.field2') AS extracted_field2, -- 按需提取其他JSON字段 last_updated FROM `your-project.your-dataset.source_table` WHERE last_updated > last_run_time ) s ON t.id = s.id WHEN MATCHED THEN UPDATE SET extracted_field1 = s.extracted_field1, extracted_field2 = s.extracted_field2, last_updated = s.last_updated WHEN NOT MATCHED THEN INSERT ( id, extracted_field1, extracted_field2, last_updated ) VALUES ( s.id, s.extracted_field1, s.extracted_field2, s.last_updated ); -- 更新控制表的上次运行时间 MERGE INTO `your-project.your-dataset.control_table` ct USING (SELECT CURRENT_TIMESTAMP() AS new_run_time) nr ON 1=1 WHEN MATCHED THEN UPDATE SET last_run_time = nr.new_run_time WHEN NOT MATCHED THEN INSERT (last_run_time) VALUES (nr.new_run_time);
方案二:分区表过滤(适合按时间分区的源表)
如果源表是按时间分区的(比如按_PARTITIONTIME或自定义时间分区),可以直接按分区范围过滤,只处理最近的分区:
MERGE INTO `your-project.your-dataset.target_table` t USING ( SELECT id, JSON_VALUE(json_column, '$.field1') AS extracted_field1, JSON_VALUE(json_column, '$.field2') AS extracted_field2 FROM `your-project.your-dataset.source_table` -- 只处理最近24小时的分区,可根据业务调整 WHERE _PARTITIONTIME BETWEEN TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 24 HOUR) AND CURRENT_TIMESTAMP() ) s ON t.id = s.id WHEN MATCHED THEN UPDATE SET extracted_field1 = s.extracted_field1, extracted_field2 = s.extracted_field2 WHEN NOT MATCHED THEN INSERT (id, extracted_field1, extracted_field2) VALUES (s.id, s.extracted_field1, s.extracted_field2);
方案三:CDC变更捕获(适合实时场景)
如果源表的数据是实时写入/更新的,可以开启BigQuery的变更数据捕获(CDC)功能,直接订阅变更日志:
-- 第一步:开启源表的CDC,生成变更日志表 CREATE CHANGE DATA CAPTURE FOR TABLE `your-project.your-dataset.source_table` INTO `your-project.your-dataset.cdc_log_table`; -- 第二步:从CDC日志提取插入/更新行并同步到目标表 MERGE INTO `your-project.your-dataset.target_table` t USING ( SELECT id, JSON_VALUE(json_column, '$.field1') AS extracted_field1, JSON_VALUE(json_column, '$.field2') AS extracted_field2 FROM `your-project.your-dataset.cdc_log_table` -- 只处理插入和更新操作,忽略删除 WHERE change_type IN ('INSERT', 'UPDATE') ) s ON t.id = s.id WHEN MATCHED THEN UPDATE SET extracted_field1 = s.extracted_field1, extracted_field2 = s.extracted_field2 WHEN NOT MATCHED THEN INSERT (id, extracted_field1, extracted_field2) VALUES (s.id, s.extracted_field1, s.extracted_field2);
关键注意点
- 必须确保源表有主键或唯一标识列,否则无法准确匹配目标表的行进行更新/去重。
- 批量场景下,控制表的维护很重要,避免重复处理同一批数据。
- 若JSON字段有嵌套结构,可使用
JSON_QUERY或JSON_EXTRACT_ARRAY等函数按需提取。
内容的提问来源于stack exchange,提问作者user21504785
相关产品推荐
相关产品推荐

