You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 22:28:13