资讯动态

Apache Airflow Grid 视图 TI 摘要流式化改造:以单条 NDJSON 流替代逐 Run 的 N+1 请求

发布时间:2026/9/10 13:21:29 来源:尧图企业网站定制
Apache Airflow Grid 视图 TI 摘要流式化改造以单条 NDJSON 流替代逐 Run 的 N1 请求【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读本文聚焦 Apache Airflow 3.x UI 中 Grid、Graph、Gantt、任务详情等视图的一次关键 API 行为改造用一条GET /ui/grid/ti_summaries/{dag_id}?run_ids...的 NDJSON 流式请求替代原先每个 Dag Run 一次请求的 N1 模式。读完本文你将掌握该流式端点Streaming Endpoint的协议格式、服务端逐行生成progressive emission机制、基于dag_version_id的序列化 DAG 共享缓存原理以及前端useGridTiSummariesStreamHook 如何消费该流并实现列逐条渐进渲染。该变更记录于 airflow-core/newsfragments/62369.significant.rst属于API 变更 行为变更无 CLI、配置、插件或依赖变更是理解 Airflow 新版 UI 数据面设计的一份重要技术快照。一、改造背景Grid 视图的 N1 请求问题在 Airflow 3 之前的 UI 架构中Grid网格、Graph图、Gantt甘特图、Task Detail任务详情等视图需要展示每个 Dag Run 下每个 Task Instance 的状态摘要。旧实现采取每个 Run 一次请求的方式对每个 Dag Run 分别调用单 Run 端点GET /ui/grid/ti_summaries/{dag_id}/{run_id}页面上同时展示 N 个 Run就会发起 N 个 HTTP 请求形成典型的N1 请求模式服务端每次请求都要重新反序列化该 Run 对应的序列化 DAG 结构重复工作量大所有列必须等全部请求完成后才能一次性渲染运行慢的 Run 会拖慢整页首屏。新方案本次变更的核心思路是把 N 次请求压缩为 1 次流式请求客户端一次传入多个run_ids服务端按 Run 逐个处理并即时输出客户端按行解析、逐列渲染。二、新端点协议GET /ui/grid/ti_summaries/{dag_id}?run_ids...2.1 请求与响应路径参数dag_id字符串必填查询参数run_ids字符串数组可选可重复传递如?run_idsrun_1run_idsrun_2响应媒体类型application/x-ndjson即NDJSONNewline-Delimited JSON流响应体每行是一个GridTISummaries对象的 JSON 序列化结果每个 Dag Run 对应一行。在 OpenAPI 规范中该端点被声明为get_grid_ti_summaries_stream200 响应描述为NDJSON stream — oneGridTISummariesJSON object per line, one per Dag run见 airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml。2.2 数据模型GridTISummaries每个 Run 的一行数据对应GridTISummaries模型见 airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/grid.pyclass GridTISummaries(BaseModel): DAG Run model for the Grid UI. run_id: str dag_id: str task_instances: list[LightGridTaskInstanceSummary]其中task_instances列表的元素类型为LightGridTaskInstanceSummary同文件 L27-L38class LightGridTaskInstanceSummary(BaseModel): Task Instance Summary model for the Grid UI. task_id: str task_display_name: str state: TaskInstanceState | None child_states: dict[TaskInstanceState | Literal[none], int] | None min_start_date: datetime | None max_end_date: datetime | None dag_version_number: int | None None has_note: bool False注意这是一份轻量摘要它只携带 Grid 渲染所需的聚合字段状态、子状态计数、最早开始/最晚结束时间、DAG 版本号、是否有备注并不携带任务的全部执行细节从而把单行体积控制在最小。一个典型的响应流形如{run_id: run_1, dag_id: example_dag, task_instances: [{task_id: t1, task_display_name: t1, state: success, child_states: null, min_start_date: ..., max_end_date: ..., dag_version_number: 3, has_note: false}]} {run_id: run_2, dag_id: example_dag, task_instances: [...]}2.3 旧端点的移除作为本次行为变更的一部分旧的单 Run 端点GET /ui/grid/ti_summaries/{dag_id}/{run_id}已被移除。所有需要 TI 摘要的 UI 视图统一改用流式端点并传入一个或多个run_ids。这意味着依赖旧端点的第三方脚本或插件需要迁移到新协议。三、服务端实现逐 Run 流式生成与连接管理新端点由get_grid_ti_summaries_stream实现位于 airflow-core/src/airflow/api_fastapi/core_api/routes/ui/grid.py。它返回一个 FastAPIStreamingResponse通过生成器_generate()逐 Run 产出数据。3.1 逐行生成与列渐进渲染生成器对run_ids列表逐个处理每个 Run 独立开启一个短生命周期 DB Sessioncreate_session(scopedFalse)查询该 Run 的全部 Task Instance 轻量字段task_id、state、dag_version_id、start_date、end_date、version_number、has_note并用yield_per1000分批拉取避免一次载入全部行调用_build_ti_summaries把原始行聚合成GridTISummaries结构一旦该 Run 的摘要就绪立即yield一行model_dump_json() \n。关键点在于服务端不等待所有 Run 处理完才返回而是每就绪一个 Run 就输出一行。客户端因此可以按行消费让 Grid 的列逐条出现columns appear progressively而不是等待整个响应完成后一次性渲染。3.2 慢客户端不占用数据库连接生成器在每次yield之间会显式关闭该 Run 的 Session见_generate内的with create_session(...)上下文管理使数据库连接在产出间隙被释放。这防止了慢客户端长时间不读取响应流时把数据库连接长期占住的问题——源码注释中明确引用了相关 issue 背景这是流式接口落地时的重要工程细节。3.3 序列化 DAG 结构按 dag_version_id 共享_build_ti_summaries会通过_get_serdag同文件 L97-L119解析该 Run 所用版本的序列化 DAG优先路径从 Run 的 Task Instance 上取dag_version_id然后从应用级共享的DBDagBag缓存中直接取该版本的序列化 DAGdag_bag.get_dag(dag_version_id, sessionsession)回退路径若dag_version_id为空对应 3.0 之前的升级场景则取该dag_id最早的DagVersion。由此多个共享同一dag_version_id的 Run 只会触发一次 DAG 反序列化——既避免了同一次请求内跨 Run 的重复反序列化也借助 app 级DBDagBag缓存避免了跨请求的重复反序列化。单元测试用查询计数验证了这一行为test_grid_ti_summaries_stream_deduplicates_serdag_loads断言两个 Run 共享同一版本时总查询数为 52 次鉴权 1 次共享的 serdag 查询 每 Run 各 1 次 TI 查询而非每个 Run 各做一次 serdag 查询的 6 次见 airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_grid.py。3.4 边角行为由测试印证静默跳过缺失 Run传入不存在的run_id时生成器对该 Run 的摘要为None则continue不会报错也不会输出该行test_grid_ti_summaries_stream_skips_missing_runs空 run_ids返回 200 与空响应体test_grid_ti_summaries_stream_empty_run_ids一行一 Run返回的行数严格等于实际存在的 Run 数test_grid_ti_summaries_stream_returns_all_runs。3.5 鉴权与聚合细节端点通过requires_access_dag同时要求TASK_INSTANCE与RUN两类DagAccessEntity的 GET 权限见 grid.py。在聚合层_build_ti_summariesL405-L464按task_id聚合状态、起止时间与备注标记再经由_find_aggregates将聚合结果映射回任务组Task Group树结构——因此返回的task_instances中同时包含 group、task、mapped_task 节点并做了group id 与 task id 冲突时优先保留 group 记录的消歧处理这也是 Grid 视图能展示任务组聚合状态的原因。四、前端消费useGridTiSummariesStream 与渐进渲染前端核心实现是useGridTiSummariesStreamHook位于 airflow-core/src/airflow/ui/src/queries/useGridTISummaries.ts。4.1 基于 fetch ReadableStream 的按行解析Hook 使用原生fetch请求${OpenAPI.BASE}/ui/grid/ti_summaries/${dagId}并把runIds映射为重复的run_ids查询参数。拿到response.body后通过getReader()读取ReadableStream分块用TextDecoder增量解码并把累积 buffer 按\n切分成行每解析出一批新行就以run_id为 key 合并进summariesByRunIdMapstring, GridTISummaries状态。由于服务端逐 Run 产出Hook 收到的每一行都对应一个新就绪的 RunGrid 视图得以在数据到达时逐列就地更新而不是整页空白等待。源码注释明确指出旧的摘要数据在流加载期间保持可见新行到达后原位更新列避免闪烁。4.2 连接生命周期AbortController 与清理每次runIds列表变化以稳定的runIdsKey runIds.join(,)为依赖或刷新 tick 变化时重新建立流式连接组件卸载或依赖变化时通过AbortController.abort()中断 fetch并对 reader 执行cancel()释放读取器避免后台残留流式连接。4.3 三重自动刷新机制Hook 为保持 Grid 数据新鲜实现了三类互补的刷新触发运行中周期重流只要还有 Run 处于 pending 状态就按useAutoRefresh的基础间隔定时递增refreshTick重开流且仅在页面处于前台focusManager.isFocused()时刷新与 React Query 的refetchIntervalInBackground: false语义保持一致状态迁移补一次从存在活跃 Run变为全部终态时由于定时器可能尚未触发过一次 tick快速任务在首个 tick 前就已结束Hook 会在该过渡瞬间额外重流一次确保最终 TI 状态落盘查询失效触发重流订阅 React Query 的 query cache当 Grid 相关的 keygetTaskInstances、getGridRuns、getDagRuns被显式invalidateQueries时以queueMicrotask合并同一执行 tick 内的多次失效为一次重流避免重复建立连接。这些逻辑共同保证了任务状态变化 → 视图及时刷新与避免无谓重连之间的平衡。4.4 全视图统一接入流式 Hook 已被 Grid、Graph、Gantt、任务详情等全部相关视图采用Grid.tsx把gridRuns的run_id列表传入 HookGraph.tsxGraph 视图同样消费summariesByRunIdGantt.tsxGantt 视图复用同一数据源MappedTaskInstance、GroupTaskInstance、TaskInstance 页面以及 useGridTISummaries.test.tsx 中的测试用例也覆盖了该 Hook 的调用形态。Graph 的测试Graph.test.tsx进一步验证了不同视图场景下向 Hook 传入states参数的行为差异。五、收益总结与迁移建议从仓库实现可以确认本次改造带来的直接收益包括维度改造前改造后请求数每 Run 一次请求N1单条流式请求携带全部 run_ids首屏体验全部 Run 就绪后一次性渲染每 Run 就绪即输出一行列渐进渲染序列化 DAG 反序列化每个 Run 各自加载同一dag_version_id跨 Run 共享且走 app 级DBDagBag缓存数据库连接占用请求期间持续占用每次 yield 间关闭 Session慢客户端不长期占用连接旧端点GET /ui/grid/ti_summaries/{dag_id}/{run_id}存在已移除对于希望验证该行为的开发者可直接运行服务端单元测试套件中的相关用例见 airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_grid.py或通过 Airflow 的 OpenAPI 文档_private_ui.yaml中get_grid_ti_summaries_stream操作手动调试该流式端点。迁移到旧端点的调用方时只需将逐 Run 请求改为一次性携带全部run_ids请求流式端点并按行解析application/x-ndjson响应即可。结语Grid TI 摘要的 NDJSON 流式化是 Airflow UI 数据面一次典型的以流式协议消除 N1重构服务端用生成器逐 Run 产出、按dag_version_id共享序列化 DAG、按 yield 边界释放数据库连接前端用ReadableStream按行消费、逐列更新并配合 AbortController 与多级刷新策略保证体验与资源消耗的平衡。理解这一模式不仅能帮你读懂 Airflow 新版 UI 的性能设计也为其他多实体摘要型列表页的流式改造提供了可直接借鉴的参考实现。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价