2234 字
约 7 分钟
1
自然语言到 Spark Job 的端到端执行链路

自然语言到 Spark Job 的端到端执行链路

  1. 要解决的问题自然语言输入是不确定的,Spark 作业执行却要求确定的表、字段、资源、状态和失败语义。系统不能把模型生成的任意代码直接执行,也不能把“生成了 SQL”包装成“作业已经成功”。这一链路的目标是把不确定输入逐步收敛成结构化计划、白名单工具调用和可审计产物。主要约束包括:普通聊天...

字数: 2118 | 语雀原文


1. 要解决的问题

自然语言输入是不确定的,Spark 作业执行却要求确定的表、字段、资源、状态和失败语义。系统不能把模型生成的任意代码直接执行,也不能把“生成了 SQL”包装成“作业已经成功”。这一链路的目标是把不确定输入逐步收敛成结构化计划、白名单工具调用和可审计产物。

主要约束包括:

  • 普通聊天、文件处理和数据工程长任务必须隔离,避免错误路由。
  • Agent 只能调用系统注册的工具,不能根据字符串动态导入或执行任意函数。
  • 数据表和字段尽量来自 Catalog,而不是完全依赖模型记忆。
  • 主模型、Spark 或 Docker 不可用时,必须暴露降级或失败状态。
  • 每个运行使用稳定的 run_id 串联计划、步骤结果和 artifact。

2. 核心对象与职责

对象 所在层 职责
SparkOsApp / terminal chat Interface 获取输入、展示流式输出,不承载规划逻辑
WorkbenchService Application 统一入口,编排三类输入路径
TurnRouter Application 提取 @文件并识别文件任务类型
ModelRouter Infrastructure 按角色调用模型并处理 fallback
AgentRuntime Application 判断 Agent 意图、生成步骤并顺序执行
ConfigCatalog Infrastructure 提供真实数据集、路径、格式和字段语义
SkillRegistry Application 加载外置 Skill 描述并检查是否存在
SparkToolExecutor Application 将白名单 tool name 映射到可信 handler
AgentPlan Domain 保存用户目标、数据集、步骤、假设和 warning
AgentRunResult Domain 保存状态、已完成步骤、artifact 和 warning

3. 输入路由架构

flowchart TD
    U["用户输入"] --> W["WorkbenchService.stream_input()"]
    W --> R["TurnRouter.route()"]
    R -->|"包含 @文件"| F["FileTaskService.run()"]
    R -->|"无文件引用"| C{"AgentRuntime.can_handle()?"}
    C -->|"是"| A["AgentRuntime.stream()"]
    C -->|"否"| M["ModelRouter.stream_chat()"]
    F --> FO["训练数据或向量知识库 artifact"]
    A --> AO["Agent 步骤结果与 artifact"]
    M --> MO["模型 token 流或本地 fallback"]
    FO --> O["统一 Iterable[str] 输出"]
    AO --> O
    MO --> O

路由优先级很重要:@文件任务先于 Agent 关键词判断。否则“构建向量知识库 @doc.md”可能因为出现“数据”或“向量”被数据工程 Agent 提前截获。

4. Agent 完整执行链路

sequenceDiagram
    autonumber
    actor User as 用户
    participant UI as TUI/Terminal
    participant WB as WorkbenchService
    participant AR as AgentRuntime
    participant Catalog as ConfigCatalog
    participant Registry as SkillRegistry
    participant Tools as SparkToolExecutor
    participant Job as JobOrchestrator
    participant Store as SQLiteJobStore

    User->>UI: 输入 Spark/ETL/质量任务
    UI->>WB: stream_input(user_input)
    WB->>AR: can_handle(user_input)
    AR-->>WB: true
    WB->>AR: stream(user_input)
    AR->>AR: plan(user_input)
    AR->>Catalog: search(user_input)
    Catalog-->>AR: List[DatasetProfile]
    AR->>AR: _select_skills() + _build_steps()
    AR->>Registry: list()/has(skill_name)
    Registry-->>AR: Skill 存在性与 warning
    loop 每个 AgentPlanStep
        AR->>Tools: execute(plan, step)
        alt 普通生成步骤
            Tools->>Tools: handler 生成 payload
            Tools->>Tools: _write_json(run_id, filename)
            Tools-->>AR: AgentStepResult
        else spark-job 步骤
            Tools->>Job: submit(run_id, job_type, payload)
            Job->>Store: save(queued/最终状态)
            Job-->>Tools: JobRecord
            Tools-->>AR: 执行配置与指标
        end
    end
    AR->>AR: 汇总 status/results/artifacts/warnings
    AR-->>WB: format_stream(result)
    WB-->>UI: Iterable[str]
    UI-->>User: 展示步骤、产物和最终状态

5. 关键实现细节

5.1 WorkbenchService.stream_input() 是统一边界

它不关心 Textual 或普通终端,只返回 Iterable[str]。因此 UI 可以替换成 HTTP SSE,而不需要把 Agent 路由重新实现一遍。三条路径最终统一成字符串迭代器,但内部语义不同:模型聊天可能是真 token 流,文件任务是完成后切块,Agent 当前是执行完成后格式化。

5.2 TurnRouter 的文件引用解析

正则 _FILE_PATTERN 支持无引号、单引号和双引号路径。相对路径以 workspace root 解析,并把 raw/path/exists 保存到 FileReference。任务类型按训练数据关键词和向量关键词识别,无法识别时返回 TaskType.UNKNOWN,由服务提示用户补充上下文。

5.3 Agent 意图与计划

AgentRuntime.can_handle() 当前基于 _AGENT_INTENT_KEYWORDSplan() 的顺序是:

  • _select_skills() 根据请求关键词生成 Skill 名称。
  • CatalogPort.search() 查找相关数据集。
  • _build_steps() 为每个 Skill 创建 AgentPlanStep
  • _skill_warnings() 检查外置 Skill 是否加载。
  • 生成 AgentPlan,包含假设、warning 和结构化步骤。 这条路径是确定性规则规划,不是 LLM 自主工具选择。它的优势是稳定、可复现;缺点是语义泛化有限。

5.4 白名单执行

SparkToolExecutor.execute() 内部维护固定字典:

handlers = {
    "spark_sql": self._spark_sql,
    "etl_pipeline": self._etl_pipeline,
    "data_cleaning": self._data_cleaning,
    "data_quality": self._data_quality,
    "feature_engineering": self._feature_engineering,
    "spark_job": self._spark_job,
    "dag_diagnosis": self._dag_diagnosis,
    "graph_compute": self._graph_compute,
}

未知工具直接抛出 ValueError。这比 getattr() 或动态 import 更安全,也便于代码审查和单元测试。

5.5 步骤失败采用短路语义

AgentRuntime.run() 逐步执行。任一步抛出异常时:

  • 保留此前已完成的 results
  • 收集已生成的 artifact;
  • 添加包含 Skill、异常类型和信息的 warning;
  • 返回 AgentRunStatus.FAILED,不继续执行依赖步骤。 这相当于简单的 fail-fast 工作流。当前没有补偿事务或步骤级重试。

6. 模型与执行环境降级

flowchart TD
    Q["调用模型角色"] --> P{"主 provider 成功?"}
    P -->|"是"| PR["返回主模型结果"]
    P -->|"否"| E["保存 _last_error"]
    E --> F{"配置 fallback_provider?"}
    F -->|"是且不同于主 provider"| L["调用本地规则模型"]
    F -->|"否"| X["重新抛出异常"]
    L --> LR["返回 fallback 结果并在 TUI 显示 warning"]

    S["spark-job 步骤"] --> D{"Docker runner 可用?"}
    D -->|"是且数据集可执行"| DS["容器内真实 Spark"]
    D -->|"否"| PS{"PySpark + Java 可用?"}
    PS -->|"是"| LS["宿主 SparkSession 执行"]
    PS -->|"否且 require_spark=true"| SF["写失败执行配置"]
    PS -->|"否且 require_spark=false"| PL["计划/本地编排模式"]

模型 fallback 和 Spark fallback 是两类不同问题:前者处理认知服务不可用,后者处理计算运行时不可用。两者都必须显式暴露,不应静默把执行语义改变掉。

7. Artifact 设计

每次运行以 artifacts/agent-runs/<run-id>/ 为根目录,按阶段编号写入 JSON:

阶段 文件
查询 01_distributed_query.json
ETL 02_etl_pipeline.json
清洗 03_data_cleaning.json
质量 04_data_quality.json
特征 05_feature_engineering.json
执行 06_execution_config.json
诊断 07_dag_diagnosis.json
图计算 08_graph_result.json

编号不是调度依据,而是便于人工浏览。真正关联依赖的是 run_idAgentPlanStep.depends_on

8. 实现难点与取舍

不确定语义与确定执行之间的边界

完全由规则选择 Skill 稳定但泛化弱;完全交给 LLM 灵活但不可控。合理演进是让模型生成受 JSON Schema 约束的候选计划,再由 Catalog、权限、成本和依赖校验器审核。

“生成”与“执行”的状态区分

SQL、质量规则或执行配置生成成功,不等于 Spark 成功。代码通过 spark_available、execution mode、Job status 和 metrics 区分这些状态,这是避免 Agent 假成功的关键。

长链路部分成功

当前失败时保留已完成 artifact,但不能从断点恢复。生产化需要持久化 step status、输入摘要、输出版本、幂等键和依赖图,以便重启后从最后成功节点继续。

安全边界仍不完整

白名单工具阻止任意函数执行,但 SQL 仍需 AST 解析、只读权限、表字段白名单、分区扫描门槛、资源预算和超时取消。

9. 源码索引与验证

  • sparkos/application/workbench.py:统一输入边界。
  • sparkos/application/turn_router.py@文件路由。
  • sparkos/application/agent_runtime.py:规则规划与 fail-fast 执行。
  • sparkos/application/spark_tools.py:白名单 handler 与 artifact。
  • sparkos/infrastructure/llm/model_router.py:模型角色和 fallback。
  • sparkos/cli.py:组件装配。
  • tests/test_workbench.py:聊天、文件、Agent、图任务和 artifact 测试。

10. 如何向面试官概括

我没有让 LLM 直接生成并执行代码,而是把自然语言任务收敛成结构化计划和白名单工具调用。系统区分聊天、文件和数据工程任务,通过 Catalog 限制表字段,通过 run_id 和 artifact 保留执行证据,并对模型与 Spark 运行时分别设计显式降级。当前规划主要是确定性规则,下一步才是受 Schema 约束的 LLM Planner。

自然语言到 Spark Job 的端到端执行链路
http://www.clxhxhhr.top/posts/784/
作者
clxstart
发布于
2026-09-18
许可协议
CC BY-NC-SA 4.0
评论
0 条
还没有评论,先写一条吧。