跳转至

端到端数据工程

本场景教程使用开源的 DAComp 数据工程数据集,串起一条完整的本地 Datus 工作流:理解数仓分层设计、在本地 DuckDB workbench 文件中交互式建表、 生成 ETL、产出 marts 数据、提交 Airflow 天级任务,并把结果写入 Superset 创建仪表盘。

本教程从哪里开始

本教程从源数据开始,新建数据管道和 Dashboard。如果你已有 Superset Dashboard,希望把它转换成分析 Agent,请阅读将 Dashboard 变成 Copilot。第一次使用 Datus 时,建议先完成安装并完成第一次提问

本文通过 Datus plugin 使用 Airflow 和 Superset。Datus datasource 负责 SQL 执行与数据传输,plugin 则通过 Airflow 和 Superset API 发现、创建、运行和 检查资源。

本地开源 quickstart 不需要 Iceberg、MinIO 或 S3。SaaS Studio tour 使用托管的 DuckDB + Iceberg lakehouse;对应的 namespace 模型见文末 「SaaS Studio Tour 变体」。

步骤 0:下载 quickstart 数据

DAComp 不包含datus-agent 仓库中。本文使用一个从 DAComp Lever 示例整理出来的小型 quickstart 数据包,不需要下载完整 DAComp 压缩包。

先创建并进入工作目录:

mkdir -p ~/datus-quickstart-data
cd ~/datus-quickstart-data

然后直接执行下面这段 bash,会下载并解压 quickstart 数据包和本地 Docker stack,创建可写的 DuckDB workbench,导出 DACOMP_HOME / DATUS_QUICKSTART_STACK,最后打印两个环境变量供后续步骤使用:

curl -L -o datus-de-lever-quickstart-v1.zip \
  https://github.com/Datus-ai/datus-quickstart-data/releases/download/data-engineering-v1/datus-de-lever-quickstart-v1.zip
curl -L -o datus-data-engineering-quickstart-stack-v1.zip \
  https://github.com/Datus-ai/datus-quickstart-data/releases/download/data-engineering-v1/datus-data-engineering-quickstart-stack-v1.zip

unzip -o datus-de-lever-quickstart-v1.zip
unzip -o datus-data-engineering-quickstart-stack-v1.zip

export DACOMP_HOME="$(pwd)/datus-de-lever-quickstart"
export DATUS_QUICKSTART_STACK="$(pwd)/data-engineering-quickstart-stack"
cp "$DACOMP_HOME/lever_start.duckdb" "$DACOMP_HOME/lever_workbench.duckdb"
cd "$DACOMP_HOME"

echo "export DACOMP_HOME=$DACOMP_HOME"
echo "export DATUS_QUICKSTART_STACK=$DATUS_QUICKSTART_STACK"

后续步骤默认这个目录下至少有这些文件:

  • docs/data_contract.yaml
  • config/layer_dependencies.yaml
  • lever_start.duckdb

步骤 1:理解数仓分层

这个 DAComp 示例已经给出了一套典型的分层数仓设计:

层级 表数量 作用
staging 24 清洗原始 ATS 数据,统一类型和格式
intermediate 17 做实体关联和可复用业务逻辑
marts 14 产出可直接分析、报表和出图的结果层

最关键的两个设计文件是:

  • docs/data_contract.yaml:描述字段清洗、校验和标准化规则
  • config/layer_dependencies.yaml:描述层级顺序与表依赖关系

在开始写 DDL 和 ETL 之前,先把这两份文件过一遍,后面给 agent 的提示词就能更贴近原始设计。

步骤 2:启动本地 quickstart 环境

下载的 stack 中已经包含本文会用到的本地 demo 服务。

Superset 中名为 examples 的 Database 使用 postgres:5432/superset_examples 连接 PostgreSQL。Superset plugin 会根据 这个不包含凭据的连接标识解析对应的 Datus datasource。启动 Superset 前, 需要把同一个 endpoint 暴露给主机,并确保 Compose service 名称可以在主机 上解析:

(
set -e
cd "$DATUS_QUICKSTART_STACK/superset"

cat > docker-compose.override.yml <<'YAML'
services:
  postgres:
    ports:
      - "5432:5432"
YAML

if grep -qE '(^|[[:space:]])postgres([[:space:]]|$)' /etc/hosts && \
   ! grep -qxE '[[:space:]]*127\.0\.0\.1[[:space:]]+postgres[[:space:]]*' /etc/hosts; then
  echo 'Conflicting /etc/hosts entry for postgres; replace it with: 127.0.0.1 postgres' >&2
  exit 1
fi
grep -qxE '[[:space:]]*127\.0\.0\.1[[:space:]]+postgres[[:space:]]*' /etc/hosts || \
  echo '127.0.0.1 postgres' | sudo tee -a /etc/hosts

docker compose up -d
)

主机的 5432 端口必须可用。这个流程使用 postgres:5432,从而让 Datus 和 Superset 返回相同的连接标识。

启动 Airflow:

(
set -e
cd "$DATUS_QUICKSTART_STACK/airflow"
docker compose up -d
)

本地默认访问方式:

  • Superset:http://127.0.0.1:8088,用户名 admin,密码 admin
  • Airflow:http://127.0.0.1:8080,用户名 admin,密码 admin
  • PostgreSQL serving database:postgres:5432/superset_examples,用户名/密码 为 superset/superset

这套 quickstart 的 Superset compose 已经带了本地演示用的元数据库和管理员默认值。

Airflow compose 会把 ${DACOMP_HOME} 挂载到容器中,并暴露一个名为 duckdb_dacomp_lever 的 Airflow connection,指向 /workspace/lever_workbench.duckdb

即使这些本地 demo 凭据是公开默认值,也不要把它们直接写进 agent.yml。 请在每个运行 Datus 的 shell 中导出:

export AIRFLOW_PASSWORD=admin
export SUPERSET_PASSWORD=admin
export SUPERSET_PG_PASSWORD=superset

步骤 3:安装并配置 plugin

把两个 plugin 都从 Datus Plugins Git 仓库安装到 Datus 所在的同一个环境中。 Superset plugin 内置了后续步骤使用的 superset-dashboard-authoring skill:

datus plugin install "git:https://github.com/Datus-ai/Datus-Plugins.git#subdirectory=datus-airflow-plugin"
datus plugin install "git:https://github.com/Datus-ai/Datus-Plugins.git#subdirectory=datus-superset-plugin"
datus plugin info airflow
datus plugin info superset

如需从已记录的 Git 来源更新两个 plugin:

datus plugin upgrade airflow
datus plugin upgrade superset

把下面这段配置合并到 ~/.datus/conf/agent.yml 现有的 agent: 下面。保留已有的 agent.providers 配置;/model 会使用这些凭据。路径会直接使用步骤 0 里导出的 DACOMP_HOMEDATUS_QUICKSTART_STACK 环境变量。

agent:
  filesystem:
    allow_write:
      - "${DATUS_QUICKSTART_STACK}/airflow/dags"

  services:
    datasources:
      lever_duckdb:
        type: duckdb
        uri: "duckdb:///${DACOMP_HOME}/lever_workbench.duckdb"
        default: true
      superset_serving:
        type: postgresql
        host: postgres
        port: 5432
        database: superset_examples
        schema: public
        username: superset
        password: ${SUPERSET_PG_PASSWORD}

  plugins:
    airflow:
      local:
        default: true
        api_base_url: http://127.0.0.1:8080/api/v1
        api_version: auto
        username: admin
        password: ${AIRFLOW_PASSWORD}
        verify_ssl: true
        timeout: 30
        dags_folder: "${DATUS_QUICKSTART_STACK}/airflow/dags"
        dag_id_prefix: daily_lever_
        allow_commands: dags,tasks,version,health

    superset:
      local:
        default: true
        api_base_url: http://127.0.0.1:8088
        auth_mode: login
        username: admin
        password: ${SUPERSET_PASSWORD}
        provider: db
        verify_ssl: "true"
        timeout: "30"

filesystem.allow_write 允许 agent 把 DAG 发布到 Airflow 挂载的主机目录。 dags_folder 告诉 agent 应把运行副本发布到哪里。DAG 发现、触发、运行状态 检查和日志读取由主 agent 通过 Airflow plugin 完成。

先为当前项目启用两个 profile,然后启动聊天会话:

cd "$DACOMP_HOME"
datus plugin enable airflow --profile local
datus plugin enable superset --profile local
datus --datasource lever_duckdb

不要自行运行具体 plugin 命令,让主 agent 验证两个服务:

使用已启用的 local profiles,先通过 Airflow plugin 查询服务端版本和健康状态,再通过 Superset plugin 查询健康状态和可用数据库。只执行只读检查,并报告所有连接或身份验证错误。

始终在启动 Datus 前完成 plugin 配置和启用。plugin skill 和环境上下文会在 session 启动时准备好;修改 profile 后请重启 session。这里选择的 lever_duckdb datasource 指向可写的 workbench 文件。

quickstart 通过 Airflow 的 AIRFLOW_CONN_DUCKDB_DACOMP_LEVER 环境变量注入 duckdb_dacomp_lever。task 运行时可以通过 BaseHook 读取环境变量 connection, 但 Airflow REST connection endpoint 不会返回它。步骤 6 会通过实际运行 DAG 验证这个 connection。

如果 CLI 提示还没有配置模型,继续之前先在 CLI 内运行:

/model

选择 provider/model,并按提示填写凭据。/model 会把 provider 凭据写入 ~/.datus/conf/agent.ymlagent.providers,并把当前项目使用的 provider/model 写入 ./.datus/config.yml

步骤 4:创建必要的 staging 表

自然语言 agent 任务不要以 CREATECOPY 这类 SQL 动词开头;CLI 会根据这些 开头关键字判断是否直接执行 SQL。

先要求 agent 创建目标 schema:

请在当前 DuckDB 数据库中创建目标 schema:staging、intermediate 和 marts。保持已有的 raw schema 不变。

这条教程只构建一条窄但完整的依赖链:marts.lever__requisition_enhanced。 字段选择、字段重命名和业务逻辑以 docs/data_contract.yaml 为准。

先要求 agent 检查物理源表字段,避免把 source-to-target 重命名误判为字段缺失:

检查 raw.requisition、raw.user、raw.requisition_posting 和 raw.requisition_offer 的 schema 和样例行。生成 SQL 前,根据物理列确认以下源到目标重命名:raw.requisition 的 id 改为 requisition_id、name 改为 requisition_name、creator_id 改为 creator_user_id、owner_id 改为 owner_user_id、hiring_manager_id 改为 hiring_manager_user_id;raw.user 的 id 改为 user_id、name 改为 user_name、external_directory_id 改为 external_directory_user_id。对于源表中已经存在的列,不要创建 NULL 占位列。

再要求 agent 根据 lever__requisition_enhancedintermediate.int_lever__requisition_userssource_models 创建必需的 staging 表。agent 会把任务分发到建表流程:

读取 ./docs/data_contract.yaml,并创建 marts.lever__requisition_enhanced 所需的 staging 表:基于 raw.requisition 创建 staging.stg_lever__requisition,基于 raw.user 创建 staging.stg_lever__user,基于 raw.requisition_posting 创建 staging.stg_lever__requisition_posting,基于 raw.requisition_offer 创建 staging.stg_lever__requisition_offer。字段设计和源到目标映射以 contract 为准。

这四张 staging 表就是 requisition enhanced 示例需要的最小 raw-to-staging 输入。

步骤 5:生成 intermediate 和 marts 表

先生成 intermediate 表。它应该按照 docs/data_contract.yamlint_lever__requisition_users 的定义,把 requisition 字段和 user 字段关联起来。

创建 intermediate 表:

读取 ./docs/data_contract.yaml,基于 staging.stg_lever__requisition 和 staging.stg_lever__user 创建 intermediate.int_lever__requisition_users。字段设计、关联关系和源到目标映射以 contract 为准。

再生成面向分析的 marts 表。契约中定义 marts.lever__requisition_enhanced 是一张按 requisition_id 一行的表,依赖:

  • intermediate.int_lever__requisition_users
  • staging.stg_lever__requisition_posting
  • staging.stg_lever__requisition_offer

创建 marts 表:

读取 ./docs/data_contract.yaml,基于 intermediate.int_lever__requisition_users、staging.stg_lever__requisition_posting 和 staging.stg_lever__requisition_offer 创建 marts.lever__requisition_enhanced。业务逻辑以 contract 为准:保留全部基础 requisition 行,按 requisition_id 统计关联的 posting 和 offer 数量,将缺失数量填充为 0,并添加 has_posting 和 has_offer 标记。

这条链路的基本顺序始终是:

staging -> intermediate -> marts

生成完成后,验证每一层以及仪表盘所需的维度:

SELECT 'stg_user' AS model, COUNT(*) AS row_count FROM staging.stg_lever__user
UNION ALL
SELECT 'stg_requisition', COUNT(*) FROM staging.stg_lever__requisition
UNION ALL
SELECT 'stg_requisition_posting', COUNT(*) FROM staging.stg_lever__requisition_posting
UNION ALL
SELECT 'stg_requisition_offer', COUNT(*) FROM staging.stg_lever__requisition_offer
UNION ALL
SELECT 'int_requisition_users', COUNT(*) FROM intermediate.int_lever__requisition_users
UNION ALL
SELECT 'marts_requisition_enhanced', COUNT(*) FROM marts.lever__requisition_enhanced;

SELECT
  COUNT(*) AS total_rows,
  COUNT(status) AS rows_with_status,
  COUNT(team) AS rows_with_team,
  COUNT(location) AS rows_with_location,
  SUM(count_postings) AS posting_links,
  SUM(count_offers) AS offer_links
FROM marts.lever__requisition_enhanced;

每个 model 都必须非空,rows_with_statusrows_with_teamrows_with_locationposting_linksoffer_links 都必须大于 0。使用 version 1 quickstart 数据包时,marts 表应有 146 行。如果维度意外全部为 NULL,请返回前面的 schema 检查,修正源字段映射后再继续。

保存并验证用于刷新同一条契约生成链路的 SQL;每天早上 8 点的调度将在步骤 6 中创建:

汇总已经成功创建 staging、intermediate 和 marts schema,以及四张 staging 表、intermediate.int_lever__requisition_users 和 marts.lever__requisition_enhanced 的准确 SQL 语句。按依赖顺序将它们写入 ./jobs/daily_lever_requisition_enhanced.sql,不要用新生成的 SQL 替换已经验证过的语句。在 lever_duckdb 上执行一次保存后的文件,并确认仍能得到相同的非零校验结果。

步骤 6:发布并运行天级 Airflow DAG

Airflow plugin 可以查询 DAG、检查源码和导入错误、触发运行,并读取 run 状态、task 状态和日志。agent 通过 filesystem 工具发布新 DAG,再使用 plugin 完成验证和运行。

对于本地 stack,发布就是把生成的文件写入 allowlist 中、并挂载到 /opt/airflow/dags 的主机目录。要求 agent 编写、发布并验证 DAG:

使用 local profile 的 Airflow plugin,并遵循其 airflow skill。创建 ./dags/daily_lever_requisition_enhanced.py,DAG ID 为 daily_lever_requisition_enhanced,schedule 为 0 8 * * *,关闭 catchup,并使用固定且带时区的开始日期。DAG 运行时读取 /workspace/jobs/daily_lever_requisition_enhanced.sql,通过 BaseHook 获取 duckdb_dacomp_lever Airflow connection,根据 connection 的 schema 或 host 还原 DuckDB SQLAlchemy URL,并在显式提交的事务中执行已经验证的 SQL。保留项目内的源文件,然后使用 filesystem tools 将完全相同的内容写入 local profile 配置的 dags_folder,并确认两个文件一致。等待 Airflow plugin 能查询到该 DAG,检查 import errors 和 DAG 详情,然后触发一次并等待运行完成。等待结束后再次读取最新 run,显示最终的 dag_run_id 和 state。如果运行失败,先检查 task states 和 logs,再报告错误。

发布和触发操作可能需要确认。同一个 agent prompt 已包含必要的回查;如需重复 检查,请让主 agent 通过 Airflow plugin 列出匹配 DAG、import errors、DAG 详情 和最新 run。

你应该会看到:

  • 维护中的源码位于 $DACOMP_HOME/dags/daily_lever_requisition_enhanced.py
  • ${DATUS_QUICKSTART_STACK}/airflow/dags 下出现内容完全相同的运行副本
  • 同一个文件在 Airflow 容器内显示为 /opt/airflow/dags/daily_lever_requisition_enhanced.py
  • Airflow 返回 dag_id、成功的 dag_run_id 和运行状态

步骤 7:把 marts 表同步到 Superset serving DB

上面的 marts 表是通过 lever_duckdb datasource 生成的。创建仪表盘之前,需要先把它复制到 superset_serving Postgres datasource。 这里的 lever_duckdbsuperset_serving 都是 agent.yml 里的 Datus datasource 名称,不是 DuckDB 或 Postgres 内部真实的 database/catalog 名。

请将 lever_duckdb datasource 中的源表 marts.lever__requisition_enhanced 复制到 superset_serving datasource,目标表为 public.lever__requisition_enhanced;如果目标表已经存在则替换。然后验证源表和目标表的行数。

如果 public.lever__requisition_enhanced 还不存在,传输工具会根据源查询结果列 自动创建目标表。version 1 数据包的源表和目标表都应该返回 146 行。

完成后,这张表就位于 Superset 中注册为 examples 的 PostgreSQL 数据库。 两边都使用 postgres:5432/superset_examples 标识数据库,因此 plugin 可以 把它唯一解析到 superset_serving Datus datasource。

步骤 8:通过 plugin 创建 Superset Dashboard

当表已经存在于 superset_serving,要求 agent 使用 plugin 的 authoring skill:

使用 local profile 的 Superset plugin,并遵循 superset-dashboard-authoring skill。这个快速上手 Dashboard 只创建三个 Chart:招聘需求总数 KPI、按 status 统计招聘需求、按 team 统计招聘需求。发现 Superset 中名为 examples 的 Database,并将其不含凭据的连接标识唯一解析到 superset_serving Datus datasource。先在该 Datus datasource 上验证 public.lever__requisition_enhanced 和计划使用的三个查询。将这张表注册为物理 Superset Dataset,然后创建 requisition operations Dashboard 及其三个 Chart。如果之前的尝试留下了匹配的 Dataset、Dashboard 或 Chart,请直接复用,不要重复创建。只把不含敏感信息的 Dataset、Dashboard 和 Chart 请求体保存在项目目录内的 JSON 文件中;不要持久化身份验证或登录请求体、token、cookie、密码或其他 secret。优先使用 typed CLI commands,仅在 typed request 被拒绝时才检查当前安装版本的 OpenAPI schema。每个 Chart 的 params 和 query_context JSON 字符串必须匹配。把三个 Chart 全部关联到 Dashboard,并更新完整的 position_json 布局,确保 Dashboard 不是空白页面。回读 Database、Dataset、Dashboard 和 Charts,确认 Database connection 仍标识为 postgres:5432/superset_examples,并执行 KPI 查询和一个分组 Chart 查询。返回 Database、Dataset、Dashboard、Chart IDs 以及 Dashboard URL。

数据准备是单独的 ETL / 调度步骤。创建仪表盘前,目标表或 SQL dataset 必须已经存在于 Superset 所识别的数据库中。Superset 的创建、更新和查询 操作会根据 plugin 权限规则要求确认。

同一个 agent prompt 会回读新建资源并执行代表性 chart 查询。如需重复验证,请把 返回的 ID 交给主 agent,让它通过 Superset plugin 检查 Database、table metadata、 Dashboard、Charts 和 chart data;不要复制本文中的示例 ID。

仪表盘应该包含 3 个 chart。total requisitions chart 查询应返回 146,按 status 和 team 分类的查询应返回多个分组,Database connection 应标识为 postgres:5432/superset_examples

这个缩小后的示例通常可以在主 agent 默认的 50 个 turn 内完成。如果仍看到 Max turns (50) exceeded,可以把下面的覆写合并到 ~/.datus/conf/agent.yml,将主 chat agent 的上限临时提高到 80:

agent:
  agentic_nodes:
    chat:
      max_turns: 80

保存后重启 Datus,再重新执行本步骤。成功后可以删除这个 max_turns 覆写,恢复 默认值;它只控制单次任务允许的最大工具推理轮数,不会让失败的任务自动续跑。

步骤 9:验证端到端结果

走完整条链路后,你应该能确认:

  • lever_workbench.duckdb 中已经有 stagingintermediatemarts schema
  • marts.lever__requisition_enhanced 是从 raw 数据经 staging 和 intermediate 层逐层加工得到的
  • $DACOMP_HOME/jobs$DACOMP_HOME/dags 中保存了已验证 SQL 和维护中的 DAG 源码
  • Airflow UI 中能看到成功运行的天级 DAG
  • 获得 Superset Database、Dataset、Dashboard、Chart ID 以及 dashboard URL

SaaS Studio Tour 变体

托管的 SaaS tour 使用同一条 Lever 工作流,但不使用本地 lever_workbench.duckdb 文件。平台会提供共享的 DuckDB + Iceberg lakehouse:

  • 共享只读 raw namespace:lake.demo_raw
  • 每个 workspace 独立可写 namespace:lake.ws_<workspace_id>
  • SaaS Airflow connection:duckdb_lever_workbench

托管平台会提供受管理的 Airflow/Superset plugin profile 和自己的 DAG 部署 通道。不要把本地的 filesystem.allow_write 或 Compose DAG 挂载路径带入 SaaS。plugin 操作仍负责发现、触发和验证托管资源;下面的 namespace 规则 保持不变。

每个用户都应该在独立 workspace 中运行 tour。backend 会按当前 workspace 渲染 seed 进去的 docs/data_contract.yaml,所以输出会写到 lake.ws_<workspace_id>,源数据继续来自 lake.demo_raw。prompt 和 SQL 应该使用完整限定名,例如:

lake.demo_raw.requisition
lake.ws_<workspace_id>.stg_lever__requisition
lake.ws_<workspace_id>.int_lever__requisition_users
lake.ws_<workspace_id>.marts_lever__requisition_enhanced

SaaS tour 中不要使用 raw.*staging.*intermediate.*marts.* 这类未限定的物理 schema 名。它们只表示逻辑层级;真实可写边界是 workspace namespace。

workspace namespace 发生变化时,需要重建 demo project 并重新生成 DAG, 确保 DAG 使用当前的 lake.ws_<workspace_id> namespace。

后续步骤