如何通过Node.js按dag_id更新BigQuery sla_table指定字段
BigQuery按dag_id更新预期开始时间实现方案
你现有代码用insert接口只会新增表行,不符合「更新现有行、不插入新数据」的要求,直接按下面改就行:
核心修改点
- 改造
determineCron函数,不再仅打印日志,直接返回格式化后的标准时间值,适配BigQuery的TIMESTAMP字段类型 - 替换原有插入逻辑为BigQuery DML UPDATE语句,以
dag_id作为匹配条件,仅更新expected_start_date字段,不会修改其他字段、也不会新增行 - 补全变量声明,增加异常兜底,避免查询不到cron配置时代码中断
- 统一时间存储格式为UTC标准ISO字符串,规避时区偏差问题
修改后完整代码
// 提前引入cron解析器,和你现有依赖保持一致 const parser = require('cron-parser'); async function accessData(){ const access = await main(); const result = access.dag_runs.map(file => ({ start_date: file.start_date, end_date: file.end_date, state: file.state, dag_run_id: file.dag_run_id, dag_id: file.dag_id })) console.log(result); // 原有插入逻辑废弃,改为后续单条匹配更新 console.log(`Parsed ${result.length} dag run records`) } async function determineCron(runItem){ const dagID = runItem?.dag_id; if (!dagID) return null; console.log('Processing dag:', dagID); // 查询cron配置,修正bigquery返回值取值逻辑 const [cronConfigRows] = await bqConnection().query(` SELECT cron_time FROM \`np-inventory-planning-thd.IPP_SLA.expected_sla\` WHERE dag_id = "${dagID}" LIMIT 1 `); if (!cronConfigRows.length) { console.warn(`No cron config found for dag_id: ${dagID}`); return null; } const cronTime = cronConfigRows[0].cron_time; const interval = parser.parseExpression(cronTime); const nextRunTime = new Date(interval.next().toString()); // 返回BigQuery兼容的UTC时间字符串 return nextRunTime.toISOString(); } (async function () { const runs = [ { start_date: '2022-06-26T23:00:00.742495+00:00', end_date: '2022-06-27T14:10:23.108401+00:00', state: 'failed', dag_run_id: 'scheduled__2022-06-25T23:00:00+00:00', dag_id: 'EFS-Winning-Route-daily-batch' }, { start_date: '2022-07-07T21:20:00.566888+00:00', end_date: '2022-07-07T23:20:55.250911+00:00', state: 'failed', dag_run_id: 'scheduled__2022-07-06T21:20:00+00:00', dag_id: 'ft-parm-trumping-daily' }, { start_date: '2022-07-03T14:00:00.779718+00:00', end_date: '2022-07-03T14:00:41.250433+00:00', state: 'failed', dag_run_id: 'scheduled__2022-06-26T14:00:00+00:00', dag_id: 'ft-parm-trumping-weekly' }, { start_date: '2022-07-08T04:00:01.038023+00:00', end_date: '2022-07-08T05:08:59.597408+00:00', state: 'failed', dag_run_id: 'scheduled__2022-07-07T04:00:00+00:00', dag_id: 'IP_MASTER' }, { start_date: '2022-07-08T13:45:00.757997+00:00', end_date: '2022-07-08T14:02:48.050405+00:00', state: 'success', dag_run_id: 'scheduled__2022-07-07T13:45:00+00:00', dag_id: 'IPP_CYCLE_PARMS' }, { start_date: '2022-07-08T02:00:00.821824+00:00', end_date: '2022-07-08T02:02:06.027268+00:00', state: 'success', dag_run_id: 'scheduled__2022-07-07T02:00:00+00:00', dag_id: 'ipp-daily-backups' }, { start_date: '2022-07-07T23:00:01.313332+00:00', end_date: '2022-07-07T23:02:37.032427+00:00', state: 'failed', dag_run_id: 'scheduled__2022-07-06T23:00:00+00:00', dag_id: 'oclt-leadtime-daily' }, { start_date: '2022-07-08T13:00:00.471935+00:00', end_date: '2022-07-08T13:00:49.819534+00:00', state: 'success', dag_run_id: 'scheduled__2022-07-07T13:00:00+00:00', dag_id: 'ope-metrics' }, { start_date: '2022-07-07T21:30:00.682954+00:00', end_date: '2022-07-07T21:33:29.885878+00:00', state: 'failed', dag_run_id: 'scheduled__2022-07-06T21:30:00+00:00', dag_id: 'parm-lite-daily' }, { start_date: '2022-07-08T08:00:01.756909+00:00', end_date: '2022-07-08T08:35:55.043131+00:00', state: 'failed', dag_run_id: 'scheduled__2022-07-07T08:00:00+00:00', dag_id: 'parm-lite-data-migration' }, { start_date: '2022-07-03T09:00:00.825694+00:00', end_date: '2022-07-03T09:10:31.680879+00:00', state: 'failed', dag_run_id: 'scheduled__2022-06-26T09:00:00+00:00', dag_id: 'parm-lite-sunday-push-instance' } ]; for (let run of runs) { try { const expectedStart = await determineCron(run); if (!expectedStart) continue; // 执行更新操作,按dag_id匹配,不新增行 const [updateRes] = await bqConnection().query(` UPDATE \`np-inventory-planning-thd.IPP_SLA.sla_table\` SET expected_start_date = TIMESTAMP("${expectedStart}") WHERE dag_id = "${run.dag_id}" `); console.log(`Updated expected_start_date for ${run.dag_id}, affected rows: ${updateRes.numDmlAffectedRows}`); } catch (err) { console.error(`Process failed for dag ${run.dag_id}:`, err); } } })();
注意事项
- 提前确认
bqConnection使用的服务账号拥有sla_table的UPDATE权限,否则DML语句会执行失败 - 如果同一个
dag_id在表中对应多条记录,当前语句会更新所有匹配行;如果仅需更新最新记录,可以在WHERE条件后追加时间筛选逻辑,比如AND create_time = (SELECT MAX(create_time) FROM \np-inventory-planning-thd.IPP_SLA.sla_table` WHERE dag_id = "${run.dag_id}")` - 代码里统一将cron解析出的中部夏令时时间转成UTC格式存储,BigQuery会自动按字段配置处理时区展示,不会出现时间偏移问题
内容的提问来源于stack exchange,提问作者Josh Davis
相关产品推荐
相关产品推荐

