尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Agent Skill 实战:用 BigQuery Continuous Queries 实现无界 SQL 流式实时计算

Agent Skill 实战:用 BigQuery Continuous Queries 实现无界 SQL 流式实时计算 Agent Skill 实战用 BigQuery Continuous Queries 实现无界 SQL 流式实时计算【免费下载链接】skillsAgent Skills for Google products and technologies项目地址: https://gitcode.com/GitHub_Trending/skills29/skillsBigQuery continuous queries连续查询是持续运行、无界unbounded的 SQL 语句它把 BigQuery 从批处理分析引擎扩展为事件驱动的实时数据加工引擎。本文基于 skills 仓库bigquery-basics技能包中的 continuous-queries.md 参考文档展开完整覆盖连续查询的两大结果输出方式、三类核心使用场景、APPENDS/CHANGES的语法要点与两条可直接复制运行的 SQL 示例并结合技能包内的 change-history.md 与 core-concepts.md 补充底层函数语义与适用限制。读完后你将能够编写写回 BigQuery 表的连续查询、编写导出到 Pub/Sub 的连续查询并正确理解其运行时长、Reservation 与算子支持范围三大约束。1. 什么是连续查询定位与结果输出方式连续查询是与普通查询根本不同的作业形态它不是提交一次、返回结果而是以无界方式持续运行的 SQL 语句对流入 BigQuery 的数据进行实时分析。在 core-concepts.md 的 Analytics Workflows 一节中官方将其归类为 Stream Processing (BigQuery continuous queries)长时运行的 SQL 语句在数据到达 BigQuery 时以准实时near real time的方式完成分析与转换从而支撑无界流式管道unbounded streaming pipelines。结合 BigQuery 的架构来看这一能力如 core-concepts.md 所述BigQuery 采用计算与存储分离的架构列式存储 可扩展的分布式分析引擎而 streaming 支持持续的数据摄入与分析。连续查询正是架设在持续摄入之上的持续计算层上游数据例如流式写入的出租车轨迹表不断落表连续查询则像一个永不停机的消费端逐批消费新增/变更的数据并立即产出结果。连续查询的结果可以通过两种方式输出这是原文档明确列出的两种出口输出方式语法载体典型去向写回 BigQuery 表INSERT语句另一张 BigQuery 目标表导出到外部系统EXPORT DATA语句Pub/Sub、Bigtable、Spanner从技能包的整体组织看这一能力是bigquery-basics技能的四个流式参考文档之一change-history.md 负责有界地查询历史变更APPENDS/CHANGES的批量用法而 continuous-queries.md 负责无界地持续消费变更两者共享同一组时间序列函数这是理解本文语法的关键线索。2. 三大核心使用场景连续查询把 BigQuery 转变为事件驱动的数据处理引擎从而解锁实时能力。原文档列举了三类典型场景逐一拆解如下2.1 事件驱动工作流与 Agentic 系统当入流数据中检测到复杂事件时触发下游应用或自治代理autonomous agents执行动作。原文档给出的示例路径是通过 Pub/Sub 集成把实时事件发送到下游 agentic 系统做进一步处理。即数据流中的某个模式命中 → 连续查询产出事件 → 写入 Pub/Sub topic → 代理系统订阅消费。在本 skills 仓库的语境下项目本身即 Google 产品 Agent Skills 集合这条链路尤其值得注意连续查询可以作为代理系统的感知层把原始事件流实时翻译成语义化事件。2.2 实时 AI 推理直接对实时数据流应用生成式 AI 模型即时生成文本或嵌入embeddings支撑个性化客户交互与实时异常检测。core-concepts.md 中将该场景表述为real-time AI inference (using Vertex AI)并且技能包中另有一个平行的 bigquery-ai-ml 技能专门覆盖AI_*系列 SQL 函数如AI_GENERATE、AI_GENERATE_EMBEDDING、AI_DETECT_ANOMALIES等。可以推断把连续查询与这些 SQL 内置 AI 函数组合是流上推理的标准做法——连续查询保证数据逐条实时到达AI 函数保证到达即计算。2.3 Reverse ETL将增强后的事件数据从 BigQuery 直接推送到 Spanner、Bigtable 等面向在线服务的操作型数据库实现低延迟的应用数据服务。原文档将其描述为Seamlessly push enhanced event data from BigQuery directly to operational databases。这正是EXPORT DATA输出方向的典型用法分析层BigQuery做完富化加工后实时反向回灌到在线层。3. 语法要点FROM 子句中的 APPENDS 与 CHANGES3.1 起点时间戳是必填要素运行连续查询时必须在FROM子句中使用APPENDS函数针对某些 Pub/Sub 导出场景则使用CHANGES来指明从最早处理哪一段数据开始。起点时间戳start timestamp定义了连续查询开始处理数据的时间点——它不是查询的结束边界而是无界消费的起点。以原文档示例中的CURRENT_TIMESTAMP() - INTERVAL 10 MINUTE为例其含义是从当前时刻往前推 10 分钟开始处理之后持续处理所有新到达的数据。这个偏移量通常用于覆盖流式写入的延迟窗口确保不遗漏提交稍有滞后的行。3.2 APPENDS 与 CHANGES 的语义差异要正确选择函数需要先理解两者的语义。change-history.md 给出了权威定义APPENDS返回指定时间范围内追加到表中的全部行只关注新增。它不要求开启任何表选项。CHANGES返回指定时间范围内表发生的所有变更包括插入INSERT、更新UPDATE和删除DELETE。使用CHANGES前必须先为表设置选项ALTER TABLE my_project.my_dataset.my_table SET OPTIONS (enable_change_history TRUE);这解释了连续查询文档中一个容易疑惑的细节为什么写回 BigQuery 表的示例用APPENDS而导出 Pub/Sub的示例用CHANGES——前者只关心新到达的行后者则需要在删除或更新产生的删除事件时把信号发到下游例如维护 Bigtable/Spanner 中的副本需要感知 DELETE 才能删掉旧记录。3.3 变更元数据列CHANGES返回的行附带变更元数据列定义来自 change-history.md元数据列类型含义_CHANGE_TYPESTRING变更类型INSERT、UPDATE、DELETE_CHANGE_TIMESTAMPTIMESTAMP产生该变更的事务提交时间_CHANGE_IS_FOR_UPDATEBOOL当 DELETE 事件由一次行更新产生时为TRUE否则为FALSE仅CHANGES提供另外时间戳参数既可以是YYYY-MM-DD HH:MM:SS格式的字符串也可以是TIMESTAMP对象连续查询示例中使用CURRENT_TIMESTAMP() - INTERVAL 10 MINUTE表达式属于后者。4. 示例一写回 BigQuery 表的连续查询INSERT APPENDS原文档给出的第一个示例持续过滤taxirides表中新追加的下车记录写入transformed_taxirides表。INSERT INTO myproject.real_time_taxi_streaming.transformed_taxirides SELECT timestamp, meter_reading, ride_status FROM APPENDS(TABLE myproject.real_time_taxi_streaming.taxirides, CURRENT_TIMESTAMP() - INTERVAL 10 MINUTE) WHERE ride_status dropoff;要点解读结构上仍是一条普通INSERT ... SELECT区别只在FROM子句APPENDS(TABLE 源表, 起点时间戳)替代了表名本身。这意味着连续查询与批查询的写法差异被压缩到了最小——你会写批处理INSERT就会写连续查询。APPENDS的第二个参数是起点时间戳。此处取当前时间减 10 分钟从该时刻起持续消费后续所有追加行。投影与过滤只取timestamp、meter_reading、ride_status三列并用ride_status dropoff做逐行过滤——这类简单投影和标量谓词正是连续查询明确支持的算子形态见第 6 节。5. 示例二导出到 Pub/Sub 的连续查询EXPORT DATA CHANGES原文档的第二个示例把结果导出到 Pub/Sub topic消息体是 JSON 字符串并在消息中额外携带乘客评论作为属性过滤条件只放行DELETE事件EXPORT DATA OPTIONS ( format CLOUD_PUBSUB, uri https://pubsub.googleapis.com/projects/myproject/topics/taxi-real-time-rides) AS ( SELECT TO_JSON_STRING( STRUCT( ride_id, timestamp, latitude, longitude)) AS message, TO_JSON( STRUCT( CAST(passenger_comment AS STRING) AS passenger_comment)) FROM CHANGES(TABLE myproject.real_time_taxi_streaming.taxi_rides, CURRENT_TIMESTAMP() - INTERVAL 10 MINUTE) WHERE _CHANGE_TYPE DELETE );要点解读EXPORT DATA是出口声明OPTIONS (format CLOUD_PUBSUB, uri ...)指定导出格式与目标 Pub/Sub topic 的完整 URIhttps://pubsub.googleapis.com/projects/project/topics/topic。文档说明EXPORT DATA还可以把结果导出到 Bigtable 或 Spanner对应 Reverse ETL 场景。Pub/Sub 消息体构造SELECT的第一列命名为message用TO_JSON_STRING(STRUCT(...))把ride_id、timestamp、latitude、longitude序列化为 JSON 作为消息载荷第二列TO_JSON(STRUCT(CAST(passenger_comment AS STRING) AS passenger_comment))构造消息属性attributes。注意这里显式CAST(... AS STRING)说明属性值需要以字符串形式承载。CHANGES_CHANGE_TYPE DELETE源表使用了CHANGES而非APPENDS并且WHERE子句精确过滤_CHANGE_TYPE DELETE——这正是第 3.2 节所述场景的落地仅当下游副本需要同步删除旧数据时才发出信号。使用CHANGES的前提是源表已开启enable_change_history见 change-history.md。注意两个示例中源表名的差异示例一的源表是taxirides示例二是taxi_rides带下划线。这来自原文档原样实际使用时请以你的源表名替换。6. 关键限制与运行约束务必在部署前确认原文档 Important Considerations Limitations 一节列出了三条硬约束它们决定了连续查询的运维形态逐条说明6.1 授权与运行时长由用户账号发起的连续查询最长运行2 天后会自动停止需要运行最长150 天必须使用服务账号service account。由此得出的实践结论任何希望长期在线的连续查询管道都应使用服务账号提交并以服务账号的生命周期150 天上限规划滚动重启/重新提交策略而不是依赖交互式用户会话。6.2 Reservation 要求运行连续查询必须使用Enterprise 或 Enterprise Plus 版的 capacity reservation并且该 reservation 要配置CONTINUOUS作业类型job type。这意味着按需付费On-demand或 Standard 版容量无法直接跑连续查询reservation 的容量规划需要为持续运行的流式作业预留 slot结合 core-concepts.md 的 Pricing 一节容量模式按专用 slot 计费。从技能包结构看iac-usage.md 说明 Terraform 的 Google Provider 支持管理 BigQuery reservations 资源因此 reservation 的CONTINUOUS作业类型分配同样可以纳入 IaC 管理保证环境一致性。6.3 受支持的算子范围连续查询支持的是有限的一组有状态stateful操作包括特定类型的JOIN、聚合aggregation和窗口函数windowing functions。许多标准 SQL 能力不受支持——除非它们作为受支持有状态操作的一部分出现典型的如SELECT DISTINCTPIVOTEXISTS等子查询实操含义设计连续查询时应以投影 逐行标量过滤 受支持的 JOIN/聚合/窗口为基本盘把复杂的去重、透视或半连接逻辑推迟到下游系统处理而不是塞进连续查询本体。7. 与 Change History 的边界有界查询 vs 无界消费仓库中的 change-history.md 与本文主题共享APPENDS/CHANGES函数但用途不同二者常被混淆这里做明确区分维度Change History有界Continuous Query无界时间范围双边界APPENDS(TABLE t, start, end)/CHANGES(TABLE t, start, end)查询固定区间单起点APPENDS(TABLE t, start)此后持续消费作业形态普通查询执行完即结束持续运行的长作业受 2 天 / 150 天上限约束典型目标增量维护 BigQuery 之外的表副本、审计变更实时管道写表、导出 Pub/Sub/Bigtable/Spanner、流上 AI 推理前置条件用CHANGES需enable_change_history TRUEAPPENDS无前置需 Enterprise/Enterprise Plus reservation CONTINUOUS作业类型change-history.md 中给出的有界查询示例可作为对照SELECT *, _CHANGE_TYPE AS change_type, _CHANGE_TIMESTAMP AS change_time FROM APPENDS(TABLE my_dataset.my_table, 2023-12-31 08:00:00, 2023-12-31 12:00:00);而连续查询把第二个边界参数去掉换成持续消费语义。理解了这一对照就能准确判断某需求该用一次性增量查询还是常驻连续查询。8. 在 skills 仓库中如何使用与验证本文全部事实依据来自 skills 仓库的bigquery-basics技能包skills/cloud/bigquery-basics/使用与验证路径如下技能入口SKILL.md 声明了技能元数据name: bigquery-basics类别BigDataAndAnalytics其 Reference Directory 将本文对应的 continuous-queries.md 与 change-history.md 并列为流式相关参考。基础准备按 SKILL.md 的 Setup 流程启用 APIgcloud services enable bigquery.googleapis.com、创建 dataset 与表其中每条gcloud命令都要求内联附加归因前缀CLOUDSDK_METRICS_REQUEST_ATTRIBUTIONagent-skills/1.0 (skill:bigquery-basics)且该前缀只用于终端直接执行的命令不得写入生成给用户仓库的脚本或 Terraform 模板。运行连续查询的前提确认项目具备 Enterprise/Enterprise Plus reservation 且作业类型为CONTINUOUS并优先使用服务账号身份提交查询以获得 150 天运行上限。深入文档原文档末尾指向了 BigQuery 官方 continuous queries 系列文档介绍、JOIN、窗口聚合、监控等章节用于进一步确认受支持算子的完整清单与监控手段仓库内则以 core-concepts.md、change-history.md、iam-security.md 作为配套参考服务账号的 IAM 授权细节可在后者中查阅。9. 小结连续查询 在FROM子句用APPENDS/CHANGES指定起点时间戳后、持续运行的无界 SQL输出仅有两种INSERT写 BigQuery 表、EXPORT DATA导出 Pub/Sub/Bigtable/Spanner。三类落地场景事件驱动 Agentic 工作流经 Pub/Sub、实时 AI 推理、Reverse ETL。选择APPENDS还是CHANGES取决于是否要感知 UPDATE/DELETE用CHANGES须先开启enable_change_history。三条硬约束用户账号 2 天 / 服务账号 150 天的运行上限、Enterprise 或 Enterprise Plus reservation 加CONTINUOUS作业类型、有限的有状态算子支持SELECT DISTINCT、PIVOT、EXISTS等不在其列。与 Change History 的边界同一组函数双边界是批量增量查询单起点是常驻流式管道按需求形态二选一。【免费下载链接】skillsAgent Skills for Google products and technologies项目地址: https://gitcode.com/GitHub_Trending/skills29/skills创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表