Loading .gitignore +1 −0 Original line number Diff line number Diff line Loading @@ -52,3 +52,4 @@ cookies.txt # Tests tests/ tmp/ README.md +79 −38 Original line number Diff line number Diff line Loading @@ -5,17 +5,11 @@ ## 前提条件 - 已安装 Docker 和 Docker Compose - **目标机器账号权限要求**:用于部署的执行账号(无论是手动执行 `make` 的当前用户,还是 Ansible Inventory 中的 `ansible_user`)**必须具有免 `sudo` 执行 Docker 命令的权限**(即该用户已加入 `docker` 用户组),或者拥有完整的 `root/sudo` 权限。 - Git - Make - Linux 环境(因涉及权限与路径映射) ## 架构升级说明 1. **DAG Factory 模式**: `dags/common/dag_factory.py` 提取了通用的 `DockerOperator` 逻辑,支持数十个流程快速接入。 2. **统一镜像仓库变量**: 引入了 `HARBOR` 环境变量控制镜像前缀,可根据部署环境自动切换仓库。 3. **远程日志收集**: 增加了 Elasticsearch (39200)、Kibana (35601) 和 Filebeat 容器,实现了任务日志的统一索引、可视化和快速检索。 4. **统一调度网关 (API Gateway)**: 针对 10 万级并发与异构 JSON 检索需求,新增了基于 FastAPI 和 Postgres JSONB 的轻量级高性能网关,屏蔽了 Airflow 的复杂 API 和鉴权逻辑。 ## 部署步骤 本系统支持两种部署方式:**基于 Make 的单机/手动部署** 和 **基于 Ansible 的多节点自动化部署**。 Loading @@ -25,42 +19,87 @@ 如果您需要在真实的集群(1个 Master 节点 + 多个 Worker 节点)上部署该系统,我们提供了开箱即用的 Ansible Playbooks。由于我们有多个不同的部署环境(如 `p368`, `csu`, `zjlab`),我们为每个环境准备了专属的 Ansible Inventory 文件。 #### 1. 准备工作 - 确保在执行机上已安装 `ansible`。 - 确保执行机配置了到所有目标节点(Master & Workers)的 SSH 免密登录(具有 root 或 sudo 权限)。 - **安装 Ansible**:如果您的执行机(通常是您的本地电脑或 Master 节点)尚未安装 Ansible,请先通过以下命令安装: **Ubuntu / Debian:** ```bash sudo apt update sudo apt install -y software-properties-common sudo apt-add-repository --yes --update ppa:ansible/ansible sudo apt install -y ansible ``` **Python pip (跨平台推荐):** ```bash pip3 install ansible ``` - **配置 SSH 免密与权限**:确保执行机配置了到所有目标节点(Master & Workers)的 SSH 免密登录。此外,清单中配置的 `ansible_user` 用户必须在目标机器上具有执行 `docker` 和 `docker compose` 命令的权限(将其加入 docker 组或具有 sudo 权限)。 #### 2. 配置环境对应的 Inventory 文件 根据您要部署的环境,编辑 `ansible/` 目录下对应的清单文件(例如 `inventory.csu.ini`): 根据您要部署的环境,编辑 `ansible/inventories/` 目录下对应的清单文件(例如 `inventory.csu.ini`)。 **关于部署路径 (`deploy_dir`) 的重要说明:** 为了确保在各种环境(尤其是包含 NFS/NAS 共享存储,或者单机混合部署)下不会产生文件覆盖与权限冲突,**我们强烈推荐在 Inventory 中为每个节点(或组)显式指定独立的 `deploy_dir`**。 即使您的 Master 和 Worker 在不同的物理机上,如果在 `/opt/` 下挂载了同一块共享存储,不区分路径也会导致容器配置被互相覆盖。 推荐部署示例(为不同角色或节点指定不同路径与别名): ```ini [master] # 替换为您的 Master 节点 IP 192.168.25.18 ansible_user=root # node-master 为节点的别名,后续可通过 ansible node-master 进行针对性操作 node-master ansible_host=192.168.25.18 ansible_user=root deploy_dir=/opt/csst-airflow-master [worker] # 添加您所有的 Worker 节点 IP 192.168.25.19 ansible_user=root 192.168.25.20 ansible_user=root node-worker1 ansible_host=192.168.25.19 ansible_user=root deploy_dir=/opt/csst-airflow-worker1 node-worker2 ansible_host=192.168.25.20 ansible_user=root deploy_dir=/opt/csst-airflow-worker2 [all:vars] # 这里已经预设好了对应的环境名称,无需修改 env_name=csu # 远程服务器上的部署路径 # 全局默认部署路径(仅在未明确指定的节点上生效) deploy_dir=/opt/csst-airflow ``` #### 3. 配置环境变量 确保对应环境的 `.env` 文件(例如 `envs/.env.csu`)中的 `MASTER_IP` 配置正确,以便 Worker 节点能找到 Master 节点的服务。 为了让集群中的所有组件(特别是 Worker 节点)能够正确找到主节点和其他依赖,您需要为该环境配置专用的 `.env` 文件。 在 `envs/` 目录下,复制或编辑对应环境的 `.env` 文件(例如 `envs/csu.env`)。您**必须**修改 `MASTER_IP` 和 `HARBOR` 这两个关键变量: ```env # --- Airflow Core Config --- AIRFLOW_IMAGE_NAME=apache/airflow:3.0.5 AIRFLOW_UID=50000 # --- Cluster & Network Config --- # 必须修改:将其替换为 Master 节点的实际 IP 地址 MASTER_IP=192.168.25.18 # 必须修改:将其替换为您当前环境使用的镜像仓库地址 HARBOR=csu-harbor.csst.nao:10443 # --- API Gateway Config (其余端口和内部 IP 默认会自动引用 MASTER_IP,通常无需修改) --- POSTGRES_HOST=${MASTER_IP} AIRFLOW_HOST=${MASTER_IP} REDIS_HOST=${MASTER_IP} ELASTICSEARCH_HOST=${MASTER_IP} ``` #### 4. 一键部署集群 在 `docker-celery-3.0.5` 目录下执行(以 `csu` 环境为例): ```bash # 1. 部署所有节点的基础组件与服务容器 ansible-playbook -i ansible/inventory.csu.ini ansible/deploy.yml ansible-playbook -i ansible/inventories/inventory.csu.ini ansible/deploy.yml # 2. 同步 DAGs 代码到所有节点 (以后每次更新 DAGs 都可以单独执行此命令) ansible-playbook -i ansible/inventory.csu.ini ansible/sync_dags.yml # 2. 同步 DAGs 代码到所有节点 (仅更新 DAGs 业务代码时执行此命令即可) ansible-playbook -i ansible/inventories/inventory.csu.ini ansible/sync_dags.yml ``` #### 5. 软件升级与更新 当您在 Master 节点上通过 `git pull` 拉取了最新的框架代码、修改了 `docker-compose.yaml` 或者是更新了自定义的镜像版本后,您只需要执行以下一键更新命令: ```bash ansible-playbook -i ansible/inventories/inventory.csu.ini ansible/update.yml ``` *(该脚本会自动将最新的代码同步给所有集群节点,并让各节点拉取最新镜像、重启发生了变更的服务容器,而不会中断未发生变更的服务。)* --- Loading @@ -81,7 +120,7 @@ make fix-permission ```bash make list-envs ``` *(系统将自动读取 `envs/` 目录下的 `.env.*` 文件后缀)* *(系统将自动读取 `envs/` 目录下的 `*.env` 文件)* #### 3. 初始化 Airflow 数据库 Loading Loading @@ -117,17 +156,15 @@ make up-worker ENV=csu | 命令 | 描述 | |------|------| | `make all` | **一键完整部署**:停止、拉取、建目录、修权限、初始化、启动 | | `make list-envs` | 列出 `envs/` 目录下所有可用的环境名称 | | `make mkdir` | 创建所有必需的挂载目录 | | `make fix-permission` | 修复 `es_data` 目录的属主权限 (UID 1000) | | `make copy-env-*` | 复制对应的预设环境变量文件 | | `make init` | 初始化 Airflow 数据库 | | `make up` | 启动核心调度器、API 网关与基础设施 (不含 Worker) | | `make up-master` | 启动核心组件并附带开启 Flower 监控 | | `make up-worker` | 启动 Celery Worker 以执行具体的 Task | | `make down` | 停止所有服务 | | `make clean` | 停止服务并清理容器 | | `make ps` | 列出运行中的服务 | | `make fix-permission` | 修复 `es_data` 目录的属主权限 (UID 50000) | | `make init ENV=<env>` | 初始化 Airflow 数据库 | | `make up-master ENV=<env>` | 启动主控节点 (包含网关与所有监控大盘) | | `make up-worker ENV=<env>` | 启动计算节点 (仅包含 Worker 和监控探针) | | `make down ENV=<env>` | 停止并移除指定环境的容器 | | `make clean ENV=<env>` | 停止服务并清理容器 | | `make ps ENV=<env>` | 列出当前运行中的服务 | ## 访问服务与系统端口映射 Loading @@ -138,8 +175,11 @@ make up-worker ENV=csu | 服务组件 | 宿主机端口 | 用途说明 | 访问地址示例 | | :--- | :--- | :--- | :--- | | **API Gateway** | `38000` | **业务系统唯一对接入口**。提供高并发任务提交与基于 JSONB 的状态检索 | `http://<MASTER_IP>:38000` | | **Task Portal** | `38501` | **统一任务与监控控制台**。提供带进度预估的 JSON 筛选及 Grafana 资源大盘内嵌页 | `http://<MASTER_IP>:38501` | | **API Gateway** | `38000` | **业务系统唯一对接入口**。提供高并发任务提交、任务控制(取消/重试)与状态检索 | `http://<MASTER_IP>:38000` | | **Airflow Web UI** | `38080` | Airflow 原生控制台与官方 API (账号密码默认: `airflow`/`airflow`) | `http://<MASTER_IP>:38080` | | **Grafana** | `33000` | 时间序列数据监控与仪表盘可视化 (账号密码默认: `admin`/`admin`) | `http://<MASTER_IP>:33000` | | **Prometheus** | `39090` | 集群监控指标拉取服务器 | `http://<MASTER_IP>:39090` | | **Kibana** | `35601` | 集中式日志可视化中心 (直接在 Discover 页面选择 `airflow-*` 视图检索日志) | `http://<MASTER_IP>:35601` | | **Flower** | `35555` | Celery 集群状态监控面板 (**仅在执行 `make up-master` 时启动**) | `http://<MASTER_IP>:35555` | | **Elasticsearch** | `39200` | 存储运行日志的底层搜索引擎 API | `http://<MASTER_IP>:39200` | Loading @@ -147,9 +187,9 @@ make up-worker ENV=csu | **Redis** | `36379` | Celery 消息队列中间件 | `redis://<MASTER_IP>:36379` | ### 计算节点 (Worker Node) 执行 `make up-worker` 的机器: * **不暴露任何宿主机端口**。 * Worker 节点只需通过 `.env` 中的 `MASTER_IP` 主动连接到主控节点的 Redis 和 Postgres 即可静默消费任务。 执行 `make up-worker ENV=<env>` 的机器: * **不暴露核心业务端口**,仅暴露 `cAdvisor` (38081) 和 `Node Exporter` (通过 Host 网络) 供主节点拉取监控指标。 * Worker 节点只需通过 `*.env` 中的 `MASTER_IP` 环境变量,主动连接到主控节点的 Redis 和 Postgres 即可静默消费任务。 --- Loading @@ -164,8 +204,9 @@ make up-worker ENV=csu - 确保 Docker 正在运行 - 确保宿主机端口未被占用:`38080` (Airflow), `39200` (ES), `35601` (Kibana) - 若 Elasticsearch 启动失败 (Unhealthy),请查看权限问题,再次执行 `make fix-permission` 并 `docker compose restart elasticsearch` - 若 Elasticsearch 启动失败 (Unhealthy),请查看权限问题,再次执行 `make fix-permission` 并 `docker compose restart elasticsearch` (或 `docker-compose restart elasticsearch`) - 检查容器日志以了解错误: ```bash docker compose logs -f # 或者使用老版本:docker-compose logs -f ``` No newline at end of file api.md +190 −49 Original line number Diff line number Diff line Loading @@ -16,7 +16,7 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 ## 1. 批量任务触发 (Batch Submit Tasks) 向系统提交任务。您可以直接发送您业务中定义的 DAG 执行消息(**必须包含 `dag_id` 和 `dag_run_id`**),网关会直接复用该 `dag_run_id` 并在后台异步推送。整个 JSON 消息体会被完整存入底层数据库以供后续检索。如果重复提交相同 `dag_run_id`,网关会忽略重复记录,实现幂等。 向系统提交任务。为了支持批次聚合查看(如按观测批次),请求体采用嵌套 JSON 结构,包含 `dag_group_run` (该批次的公用信息) 和 `dag_run_list` (每个子任务的具体信息)。**每个子任务必须包含 `dag`(或 `dag_id`)和 `dag_run`(或 `dag_run_id`)**,网关会直接复用该 `dag_run` 作为底层任务 ID 并在后台异步推送。整个 JSON 消息体会被完整存入底层数据库以供后续检索。如果重复提交相同 ID,网关会忽略重复记录,实现幂等。 * **Endpoint**: `POST /api/tasks/batch-submit` * **请求头**: `Content-Type: application/json` Loading @@ -24,71 +24,169 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 ```bash curl -X POST "$AIRFLOW_API_GATEWAY/api/tasks/batch-submit" \ -H "Content-Type: application/json" \ -d '[ -d '{ "dag_group_run": { "dag_group": "default", "dag_group_run": "18a2662c9ea9948a0f802e007ecd1ce6ba4e927e", "batch_id": "default", "priority": 1, "created_time": "2026-04-02T06:21:41.675" }, "dag_run_list": [ { "dag_id": "csst-msc-l1-mbi", "dag_run_id": "123e4567-e89b-12d3-a456-426614174000", "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100131914", "detector": "09", "batch_id": "inttest", "docker_images": { "csst-msc-l1-mbi": "latest" } "obs_id": "10100547339", "detector": "24", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "priority": 1, "data_list": [ "69c3d451aed3b579a8c0b58a" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" }, { "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100547339", "detector": "25", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "a992e82967594fa4134a0f6bcb0b4172778d6780", "priority": 1, "data_list": [ "69c3d45519ed78b54758f180" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" } ]' ] }' ``` * **测试结果 (成功响应示例)**: ```json { "status": "accepted", "task_ids": [ "123e4567-e89b-12d3-a456-426614174000" "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "a992e82967594fa4134a0f6bcb0b4172778d6780" ] } ``` --- ## 2. 单个任务状态查询 (Query Task Status) ## 2. 任务批次聚合查询 (List Task Groups) 利用任务提交时你传入的那个 `dag_run_id`,直接获取任务的实时状态、执行结果或错误日志。数据源自数据库的精确查询($O(1)$ 复杂度)。 用于查询提交时通过 `dag_group_run` 字段绑定的批次聚合信息,系统会自动按该字段分组,并统计出每个批次下任务的各个状态总数(如成功几个、失败几个、运行中几个等)。 * **Endpoint**: `GET /api/tasks/groups` * **参数**: `?limit=50&offset=0` * **cURL 示例**: ```bash curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/groups" ``` * **测试结果 (成功响应示例)**: ```json [ { "dag_group_run": "18a2662c9ea9948a0f802e007ecd1ce6ba4e927e", "batch_id": "default", "created_at": "2026-04-02T06:21:41.675000", "total_tasks": 2, "success_tasks": 2, "failed_tasks": 0, "running_tasks": 0, "received_tasks": 0, "cancelled_tasks": 0 } ] ``` --- ## 3. 单个任务状态查询 (Query Task Status) 利用任务提交时你传入的那个 `dag_run`,直接获取任务的实时状态、执行结果或错误日志。数据源自数据库的精确查询($O(1)$ 复杂度)。 * **Endpoint**: `GET /api/tasks/{dag_run_id}` * **cURL 示例**: ```bash curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/123e4567-e89b-12d3-a456-426614174000" curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae" ``` * **测试结果 (成功响应示例)**: ```json { "task_id": "123e4567-e89b-12d3-a456-426614174000", "task_id": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "dag_id": "csst-msc-l1-mbi", "status": "success", "inputs": { "dag_id": "csst-msc-l1-mbi", "dag_run_id": "123e4567-e89b-12d3-a456-426614174000", "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100131914", "detector": "09", "batch_id": "inttest", "docker_images": { "csst-msc-l1-mbi": "latest" } "obs_id": "10100547339", "detector": "24", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "priority": 1, "data_list": [ "69c3d451aed3b579a8c0b58a" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" }, "outputs": { "finished_at": "2026-04-01T15:30:00.123456" "finished_at": "2026-04-02T06:30:00.123456" }, "error_message": null, "execution_logs": null, "created_at": "2026-04-01T15:00:00.000000Z", "updated_at": "2026-04-01T15:30:00.123456Z" "created_at": "2026-04-02T06:22:33.846000Z", "updated_at": "2026-04-02T06:30:00.123456Z" } ``` * **状态枚举 (`status`)**: Loading @@ -99,41 +197,58 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 --- ## 3. 基于异构 JSON 字段的批量检索 (Search by JSON Inputs) ## 4. 基于异构 JSON 字段的模糊批量检索 (Search by JSON Inputs) 得益于底层 `PostgreSQL JSONB + GIN 索引` 的设计,您可以随意根据您当初提交的 `inputs` 中的任意嵌套字段来进行高性能的反向检索,这对于批量查询非常有用。 得益于底层 `PostgreSQL JSONB` 字段的支持与最新的 API Gateway 升级,您现在可以随意根据当初提交的 `inputs` 中的任意嵌套字段来进行高性能的反向**模糊检索 (ILIKE)**。 * 如果提供的是字符串类型,系统会自动进行类似 `%value%` 的部分匹配。 * 如果提供的是非字符串类型(数字等),系统会进行精确匹配。 * **Endpoint**: `POST /api/tasks/search` * **参数**: `?limit=100&offset=0` * **cURL 示例** (查询所有 inputs 中 `obs_id` 为 `10100131914` 且 `batch_id` 为 `inttest` 的任务): * **cURL 示例** (查询所有 inputs 中 `obs_id` 包含 `547339` 且 `batch_id` 为 `default` 的任务): ```bash curl -s -X POST "$AIRFLOW_API_GATEWAY/api/tasks/search?limit=10" \ -H "Content-Type: application/json" \ -d '{ "obs_id": "10100131914", "batch_id": "inttest" "obs_id": "547339", "batch_id": "default" }' ``` * **测试结果 (成功响应示例)**: ```json [ { "task_id": "123e4567-e89b-12d3-a456-426614174000", "task_id": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "dag_id": "csst-msc-l1-mbi", "status": "success", "inputs": { "dag_id": "csst-msc-l1-mbi", "dag_run_id": "123e4567-e89b-12d3-a456-426614174000", "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100131914", "detector": "09", "batch_id": "inttest", "docker_images": { "csst-msc-l1-mbi": "latest" } "obs_id": "10100547339", "detector": "24", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "priority": 1, "data_list": [ "69c3d451aed3b579a8c0b58a" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" } } ] Loading @@ -141,9 +256,35 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 --- ## 4. 关于深度日志查询 ## 5. 获取任务动态字段与字典 (Metadata Fields) 用于前端渲染下拉框和级联选择。接口会扫描最近 1000 个任务,返回 `inputs` JSONB 中曾经出现过的所有 key,以及它们对应的所有不重复的 value 集合。 * **Endpoint**: `GET /api/tasks/metadata/fields` * **cURL 示例**: ```bash curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/metadata/fields" ``` * **测试结果 (成功响应示例)**: ```json { "fields": { "batch_id": ["default"], "dataset": ["test-msc-c9-25sqdeg-v3"], "detector": ["24", "25"], "instrument": ["MSC"], "obs_group": ["W5"], "obs_id": ["10100547339"], "obs_type": ["WIDE"] } } ``` --- ## 6. 关于深度日志查询 目前 API 网关支持返回关键错误信息 (`error_message`) 和简要摘要 (`outputs`)。如果您需要对失败的任务进行排错,可以通过上面返回的 `task_id`(它同时也是 Airflow 的 `dag_run_id`)前往 Elasticsearch 或 Kibana 进行深入的全量全文日志检索: 目前 API 网关的 `GET /api/tasks/{dag_run_id}` 接口已集成了从 Elasticsearch 拉取全量执行日志 (`execution_logs`) 的能力。如果您还需要进入 Kibana 进行更复杂的汇聚分析或全栈链路追踪,可以: * **Kibana URL**: `http://10.73.0.27:35601` * **检索条件示例**: 在 Kibana 的 `filebeat-*` 索引中搜索 `"123e4567-e89b-12d3-a456-426614174000"`。 No newline at end of file * **检索条件示例**: 在 Kibana 的 `filebeat-*` 索引中搜索 `"07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae"`。 Loading
.gitignore +1 −0 Original line number Diff line number Diff line Loading @@ -52,3 +52,4 @@ cookies.txt # Tests tests/ tmp/
README.md +79 −38 Original line number Diff line number Diff line Loading @@ -5,17 +5,11 @@ ## 前提条件 - 已安装 Docker 和 Docker Compose - **目标机器账号权限要求**:用于部署的执行账号(无论是手动执行 `make` 的当前用户,还是 Ansible Inventory 中的 `ansible_user`)**必须具有免 `sudo` 执行 Docker 命令的权限**(即该用户已加入 `docker` 用户组),或者拥有完整的 `root/sudo` 权限。 - Git - Make - Linux 环境(因涉及权限与路径映射) ## 架构升级说明 1. **DAG Factory 模式**: `dags/common/dag_factory.py` 提取了通用的 `DockerOperator` 逻辑,支持数十个流程快速接入。 2. **统一镜像仓库变量**: 引入了 `HARBOR` 环境变量控制镜像前缀,可根据部署环境自动切换仓库。 3. **远程日志收集**: 增加了 Elasticsearch (39200)、Kibana (35601) 和 Filebeat 容器,实现了任务日志的统一索引、可视化和快速检索。 4. **统一调度网关 (API Gateway)**: 针对 10 万级并发与异构 JSON 检索需求,新增了基于 FastAPI 和 Postgres JSONB 的轻量级高性能网关,屏蔽了 Airflow 的复杂 API 和鉴权逻辑。 ## 部署步骤 本系统支持两种部署方式:**基于 Make 的单机/手动部署** 和 **基于 Ansible 的多节点自动化部署**。 Loading @@ -25,42 +19,87 @@ 如果您需要在真实的集群(1个 Master 节点 + 多个 Worker 节点)上部署该系统,我们提供了开箱即用的 Ansible Playbooks。由于我们有多个不同的部署环境(如 `p368`, `csu`, `zjlab`),我们为每个环境准备了专属的 Ansible Inventory 文件。 #### 1. 准备工作 - 确保在执行机上已安装 `ansible`。 - 确保执行机配置了到所有目标节点(Master & Workers)的 SSH 免密登录(具有 root 或 sudo 权限)。 - **安装 Ansible**:如果您的执行机(通常是您的本地电脑或 Master 节点)尚未安装 Ansible,请先通过以下命令安装: **Ubuntu / Debian:** ```bash sudo apt update sudo apt install -y software-properties-common sudo apt-add-repository --yes --update ppa:ansible/ansible sudo apt install -y ansible ``` **Python pip (跨平台推荐):** ```bash pip3 install ansible ``` - **配置 SSH 免密与权限**:确保执行机配置了到所有目标节点(Master & Workers)的 SSH 免密登录。此外,清单中配置的 `ansible_user` 用户必须在目标机器上具有执行 `docker` 和 `docker compose` 命令的权限(将其加入 docker 组或具有 sudo 权限)。 #### 2. 配置环境对应的 Inventory 文件 根据您要部署的环境,编辑 `ansible/` 目录下对应的清单文件(例如 `inventory.csu.ini`): 根据您要部署的环境,编辑 `ansible/inventories/` 目录下对应的清单文件(例如 `inventory.csu.ini`)。 **关于部署路径 (`deploy_dir`) 的重要说明:** 为了确保在各种环境(尤其是包含 NFS/NAS 共享存储,或者单机混合部署)下不会产生文件覆盖与权限冲突,**我们强烈推荐在 Inventory 中为每个节点(或组)显式指定独立的 `deploy_dir`**。 即使您的 Master 和 Worker 在不同的物理机上,如果在 `/opt/` 下挂载了同一块共享存储,不区分路径也会导致容器配置被互相覆盖。 推荐部署示例(为不同角色或节点指定不同路径与别名): ```ini [master] # 替换为您的 Master 节点 IP 192.168.25.18 ansible_user=root # node-master 为节点的别名,后续可通过 ansible node-master 进行针对性操作 node-master ansible_host=192.168.25.18 ansible_user=root deploy_dir=/opt/csst-airflow-master [worker] # 添加您所有的 Worker 节点 IP 192.168.25.19 ansible_user=root 192.168.25.20 ansible_user=root node-worker1 ansible_host=192.168.25.19 ansible_user=root deploy_dir=/opt/csst-airflow-worker1 node-worker2 ansible_host=192.168.25.20 ansible_user=root deploy_dir=/opt/csst-airflow-worker2 [all:vars] # 这里已经预设好了对应的环境名称,无需修改 env_name=csu # 远程服务器上的部署路径 # 全局默认部署路径(仅在未明确指定的节点上生效) deploy_dir=/opt/csst-airflow ``` #### 3. 配置环境变量 确保对应环境的 `.env` 文件(例如 `envs/.env.csu`)中的 `MASTER_IP` 配置正确,以便 Worker 节点能找到 Master 节点的服务。 为了让集群中的所有组件(特别是 Worker 节点)能够正确找到主节点和其他依赖,您需要为该环境配置专用的 `.env` 文件。 在 `envs/` 目录下,复制或编辑对应环境的 `.env` 文件(例如 `envs/csu.env`)。您**必须**修改 `MASTER_IP` 和 `HARBOR` 这两个关键变量: ```env # --- Airflow Core Config --- AIRFLOW_IMAGE_NAME=apache/airflow:3.0.5 AIRFLOW_UID=50000 # --- Cluster & Network Config --- # 必须修改:将其替换为 Master 节点的实际 IP 地址 MASTER_IP=192.168.25.18 # 必须修改:将其替换为您当前环境使用的镜像仓库地址 HARBOR=csu-harbor.csst.nao:10443 # --- API Gateway Config (其余端口和内部 IP 默认会自动引用 MASTER_IP,通常无需修改) --- POSTGRES_HOST=${MASTER_IP} AIRFLOW_HOST=${MASTER_IP} REDIS_HOST=${MASTER_IP} ELASTICSEARCH_HOST=${MASTER_IP} ``` #### 4. 一键部署集群 在 `docker-celery-3.0.5` 目录下执行(以 `csu` 环境为例): ```bash # 1. 部署所有节点的基础组件与服务容器 ansible-playbook -i ansible/inventory.csu.ini ansible/deploy.yml ansible-playbook -i ansible/inventories/inventory.csu.ini ansible/deploy.yml # 2. 同步 DAGs 代码到所有节点 (以后每次更新 DAGs 都可以单独执行此命令) ansible-playbook -i ansible/inventory.csu.ini ansible/sync_dags.yml # 2. 同步 DAGs 代码到所有节点 (仅更新 DAGs 业务代码时执行此命令即可) ansible-playbook -i ansible/inventories/inventory.csu.ini ansible/sync_dags.yml ``` #### 5. 软件升级与更新 当您在 Master 节点上通过 `git pull` 拉取了最新的框架代码、修改了 `docker-compose.yaml` 或者是更新了自定义的镜像版本后,您只需要执行以下一键更新命令: ```bash ansible-playbook -i ansible/inventories/inventory.csu.ini ansible/update.yml ``` *(该脚本会自动将最新的代码同步给所有集群节点,并让各节点拉取最新镜像、重启发生了变更的服务容器,而不会中断未发生变更的服务。)* --- Loading @@ -81,7 +120,7 @@ make fix-permission ```bash make list-envs ``` *(系统将自动读取 `envs/` 目录下的 `.env.*` 文件后缀)* *(系统将自动读取 `envs/` 目录下的 `*.env` 文件)* #### 3. 初始化 Airflow 数据库 Loading Loading @@ -117,17 +156,15 @@ make up-worker ENV=csu | 命令 | 描述 | |------|------| | `make all` | **一键完整部署**:停止、拉取、建目录、修权限、初始化、启动 | | `make list-envs` | 列出 `envs/` 目录下所有可用的环境名称 | | `make mkdir` | 创建所有必需的挂载目录 | | `make fix-permission` | 修复 `es_data` 目录的属主权限 (UID 1000) | | `make copy-env-*` | 复制对应的预设环境变量文件 | | `make init` | 初始化 Airflow 数据库 | | `make up` | 启动核心调度器、API 网关与基础设施 (不含 Worker) | | `make up-master` | 启动核心组件并附带开启 Flower 监控 | | `make up-worker` | 启动 Celery Worker 以执行具体的 Task | | `make down` | 停止所有服务 | | `make clean` | 停止服务并清理容器 | | `make ps` | 列出运行中的服务 | | `make fix-permission` | 修复 `es_data` 目录的属主权限 (UID 50000) | | `make init ENV=<env>` | 初始化 Airflow 数据库 | | `make up-master ENV=<env>` | 启动主控节点 (包含网关与所有监控大盘) | | `make up-worker ENV=<env>` | 启动计算节点 (仅包含 Worker 和监控探针) | | `make down ENV=<env>` | 停止并移除指定环境的容器 | | `make clean ENV=<env>` | 停止服务并清理容器 | | `make ps ENV=<env>` | 列出当前运行中的服务 | ## 访问服务与系统端口映射 Loading @@ -138,8 +175,11 @@ make up-worker ENV=csu | 服务组件 | 宿主机端口 | 用途说明 | 访问地址示例 | | :--- | :--- | :--- | :--- | | **API Gateway** | `38000` | **业务系统唯一对接入口**。提供高并发任务提交与基于 JSONB 的状态检索 | `http://<MASTER_IP>:38000` | | **Task Portal** | `38501` | **统一任务与监控控制台**。提供带进度预估的 JSON 筛选及 Grafana 资源大盘内嵌页 | `http://<MASTER_IP>:38501` | | **API Gateway** | `38000` | **业务系统唯一对接入口**。提供高并发任务提交、任务控制(取消/重试)与状态检索 | `http://<MASTER_IP>:38000` | | **Airflow Web UI** | `38080` | Airflow 原生控制台与官方 API (账号密码默认: `airflow`/`airflow`) | `http://<MASTER_IP>:38080` | | **Grafana** | `33000` | 时间序列数据监控与仪表盘可视化 (账号密码默认: `admin`/`admin`) | `http://<MASTER_IP>:33000` | | **Prometheus** | `39090` | 集群监控指标拉取服务器 | `http://<MASTER_IP>:39090` | | **Kibana** | `35601` | 集中式日志可视化中心 (直接在 Discover 页面选择 `airflow-*` 视图检索日志) | `http://<MASTER_IP>:35601` | | **Flower** | `35555` | Celery 集群状态监控面板 (**仅在执行 `make up-master` 时启动**) | `http://<MASTER_IP>:35555` | | **Elasticsearch** | `39200` | 存储运行日志的底层搜索引擎 API | `http://<MASTER_IP>:39200` | Loading @@ -147,9 +187,9 @@ make up-worker ENV=csu | **Redis** | `36379` | Celery 消息队列中间件 | `redis://<MASTER_IP>:36379` | ### 计算节点 (Worker Node) 执行 `make up-worker` 的机器: * **不暴露任何宿主机端口**。 * Worker 节点只需通过 `.env` 中的 `MASTER_IP` 主动连接到主控节点的 Redis 和 Postgres 即可静默消费任务。 执行 `make up-worker ENV=<env>` 的机器: * **不暴露核心业务端口**,仅暴露 `cAdvisor` (38081) 和 `Node Exporter` (通过 Host 网络) 供主节点拉取监控指标。 * Worker 节点只需通过 `*.env` 中的 `MASTER_IP` 环境变量,主动连接到主控节点的 Redis 和 Postgres 即可静默消费任务。 --- Loading @@ -164,8 +204,9 @@ make up-worker ENV=csu - 确保 Docker 正在运行 - 确保宿主机端口未被占用:`38080` (Airflow), `39200` (ES), `35601` (Kibana) - 若 Elasticsearch 启动失败 (Unhealthy),请查看权限问题,再次执行 `make fix-permission` 并 `docker compose restart elasticsearch` - 若 Elasticsearch 启动失败 (Unhealthy),请查看权限问题,再次执行 `make fix-permission` 并 `docker compose restart elasticsearch` (或 `docker-compose restart elasticsearch`) - 检查容器日志以了解错误: ```bash docker compose logs -f # 或者使用老版本:docker-compose logs -f ``` No newline at end of file
api.md +190 −49 Original line number Diff line number Diff line Loading @@ -16,7 +16,7 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 ## 1. 批量任务触发 (Batch Submit Tasks) 向系统提交任务。您可以直接发送您业务中定义的 DAG 执行消息(**必须包含 `dag_id` 和 `dag_run_id`**),网关会直接复用该 `dag_run_id` 并在后台异步推送。整个 JSON 消息体会被完整存入底层数据库以供后续检索。如果重复提交相同 `dag_run_id`,网关会忽略重复记录,实现幂等。 向系统提交任务。为了支持批次聚合查看(如按观测批次),请求体采用嵌套 JSON 结构,包含 `dag_group_run` (该批次的公用信息) 和 `dag_run_list` (每个子任务的具体信息)。**每个子任务必须包含 `dag`(或 `dag_id`)和 `dag_run`(或 `dag_run_id`)**,网关会直接复用该 `dag_run` 作为底层任务 ID 并在后台异步推送。整个 JSON 消息体会被完整存入底层数据库以供后续检索。如果重复提交相同 ID,网关会忽略重复记录,实现幂等。 * **Endpoint**: `POST /api/tasks/batch-submit` * **请求头**: `Content-Type: application/json` Loading @@ -24,71 +24,169 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 ```bash curl -X POST "$AIRFLOW_API_GATEWAY/api/tasks/batch-submit" \ -H "Content-Type: application/json" \ -d '[ -d '{ "dag_group_run": { "dag_group": "default", "dag_group_run": "18a2662c9ea9948a0f802e007ecd1ce6ba4e927e", "batch_id": "default", "priority": 1, "created_time": "2026-04-02T06:21:41.675" }, "dag_run_list": [ { "dag_id": "csst-msc-l1-mbi", "dag_run_id": "123e4567-e89b-12d3-a456-426614174000", "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100131914", "detector": "09", "batch_id": "inttest", "docker_images": { "csst-msc-l1-mbi": "latest" } "obs_id": "10100547339", "detector": "24", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "priority": 1, "data_list": [ "69c3d451aed3b579a8c0b58a" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" }, { "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100547339", "detector": "25", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "a992e82967594fa4134a0f6bcb0b4172778d6780", "priority": 1, "data_list": [ "69c3d45519ed78b54758f180" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" } ]' ] }' ``` * **测试结果 (成功响应示例)**: ```json { "status": "accepted", "task_ids": [ "123e4567-e89b-12d3-a456-426614174000" "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "a992e82967594fa4134a0f6bcb0b4172778d6780" ] } ``` --- ## 2. 单个任务状态查询 (Query Task Status) ## 2. 任务批次聚合查询 (List Task Groups) 利用任务提交时你传入的那个 `dag_run_id`,直接获取任务的实时状态、执行结果或错误日志。数据源自数据库的精确查询($O(1)$ 复杂度)。 用于查询提交时通过 `dag_group_run` 字段绑定的批次聚合信息,系统会自动按该字段分组,并统计出每个批次下任务的各个状态总数(如成功几个、失败几个、运行中几个等)。 * **Endpoint**: `GET /api/tasks/groups` * **参数**: `?limit=50&offset=0` * **cURL 示例**: ```bash curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/groups" ``` * **测试结果 (成功响应示例)**: ```json [ { "dag_group_run": "18a2662c9ea9948a0f802e007ecd1ce6ba4e927e", "batch_id": "default", "created_at": "2026-04-02T06:21:41.675000", "total_tasks": 2, "success_tasks": 2, "failed_tasks": 0, "running_tasks": 0, "received_tasks": 0, "cancelled_tasks": 0 } ] ``` --- ## 3. 单个任务状态查询 (Query Task Status) 利用任务提交时你传入的那个 `dag_run`,直接获取任务的实时状态、执行结果或错误日志。数据源自数据库的精确查询($O(1)$ 复杂度)。 * **Endpoint**: `GET /api/tasks/{dag_run_id}` * **cURL 示例**: ```bash curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/123e4567-e89b-12d3-a456-426614174000" curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae" ``` * **测试结果 (成功响应示例)**: ```json { "task_id": "123e4567-e89b-12d3-a456-426614174000", "task_id": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "dag_id": "csst-msc-l1-mbi", "status": "success", "inputs": { "dag_id": "csst-msc-l1-mbi", "dag_run_id": "123e4567-e89b-12d3-a456-426614174000", "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100131914", "detector": "09", "batch_id": "inttest", "docker_images": { "csst-msc-l1-mbi": "latest" } "obs_id": "10100547339", "detector": "24", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "priority": 1, "data_list": [ "69c3d451aed3b579a8c0b58a" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" }, "outputs": { "finished_at": "2026-04-01T15:30:00.123456" "finished_at": "2026-04-02T06:30:00.123456" }, "error_message": null, "execution_logs": null, "created_at": "2026-04-01T15:00:00.000000Z", "updated_at": "2026-04-01T15:30:00.123456Z" "created_at": "2026-04-02T06:22:33.846000Z", "updated_at": "2026-04-02T06:30:00.123456Z" } ``` * **状态枚举 (`status`)**: Loading @@ -99,41 +197,58 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 --- ## 3. 基于异构 JSON 字段的批量检索 (Search by JSON Inputs) ## 4. 基于异构 JSON 字段的模糊批量检索 (Search by JSON Inputs) 得益于底层 `PostgreSQL JSONB + GIN 索引` 的设计,您可以随意根据您当初提交的 `inputs` 中的任意嵌套字段来进行高性能的反向检索,这对于批量查询非常有用。 得益于底层 `PostgreSQL JSONB` 字段的支持与最新的 API Gateway 升级,您现在可以随意根据当初提交的 `inputs` 中的任意嵌套字段来进行高性能的反向**模糊检索 (ILIKE)**。 * 如果提供的是字符串类型,系统会自动进行类似 `%value%` 的部分匹配。 * 如果提供的是非字符串类型(数字等),系统会进行精确匹配。 * **Endpoint**: `POST /api/tasks/search` * **参数**: `?limit=100&offset=0` * **cURL 示例** (查询所有 inputs 中 `obs_id` 为 `10100131914` 且 `batch_id` 为 `inttest` 的任务): * **cURL 示例** (查询所有 inputs 中 `obs_id` 包含 `547339` 且 `batch_id` 为 `default` 的任务): ```bash curl -s -X POST "$AIRFLOW_API_GATEWAY/api/tasks/search?limit=10" \ -H "Content-Type: application/json" \ -d '{ "obs_id": "10100131914", "batch_id": "inttest" "obs_id": "547339", "batch_id": "default" }' ``` * **测试结果 (成功响应示例)**: ```json [ { "task_id": "123e4567-e89b-12d3-a456-426614174000", "task_id": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "dag_id": "csst-msc-l1-mbi", "status": "success", "inputs": { "dag_id": "csst-msc-l1-mbi", "dag_run_id": "123e4567-e89b-12d3-a456-426614174000", "dataset": "test-msc-c9-25sqdeg-v3", "instrument": "MSC", "obs_type": "WIDE", "obs_group": "W5", "obs_id": "10100131914", "detector": "09", "batch_id": "inttest", "docker_images": { "csst-msc-l1-mbi": "latest" } "obs_id": "10100547339", "detector": "24", "filter": "", "custom_id": "", "batch_id": "default", "pmapname": "", "ref_cat": "", "dag_group": "default", "dag": "csst-msc-l1-mbi", "dag_group_run": "3c4033ea77904a27b23a50be39c5a0ad1e98869b", "dag_run": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae", "priority": 1, "data_list": [ "69c3d451aed3b579a8c0b58a" ], "extra_kwargs": {}, "created_time": "2026-04-02T06:22:33.846", "rerun": -1, "status_code": -1024, "n_file_expected": 1, "n_file_found": 1, "object": "", "proposal_id": "" } } ] Loading @@ -141,9 +256,35 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000 --- ## 4. 关于深度日志查询 ## 5. 获取任务动态字段与字典 (Metadata Fields) 用于前端渲染下拉框和级联选择。接口会扫描最近 1000 个任务,返回 `inputs` JSONB 中曾经出现过的所有 key,以及它们对应的所有不重复的 value 集合。 * **Endpoint**: `GET /api/tasks/metadata/fields` * **cURL 示例**: ```bash curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/metadata/fields" ``` * **测试结果 (成功响应示例)**: ```json { "fields": { "batch_id": ["default"], "dataset": ["test-msc-c9-25sqdeg-v3"], "detector": ["24", "25"], "instrument": ["MSC"], "obs_group": ["W5"], "obs_id": ["10100547339"], "obs_type": ["WIDE"] } } ``` --- ## 6. 关于深度日志查询 目前 API 网关支持返回关键错误信息 (`error_message`) 和简要摘要 (`outputs`)。如果您需要对失败的任务进行排错,可以通过上面返回的 `task_id`(它同时也是 Airflow 的 `dag_run_id`)前往 Elasticsearch 或 Kibana 进行深入的全量全文日志检索: 目前 API 网关的 `GET /api/tasks/{dag_run_id}` 接口已集成了从 Elasticsearch 拉取全量执行日志 (`execution_logs`) 的能力。如果您还需要进入 Kibana 进行更复杂的汇聚分析或全栈链路追踪,可以: * **Kibana URL**: `http://10.73.0.27:35601` * **检索条件示例**: 在 Kibana 的 `filebeat-*` 索引中搜索 `"123e4567-e89b-12d3-a456-426614174000"`。 No newline at end of file * **检索条件示例**: 在 Kibana 的 `filebeat-*` 索引中搜索 `"07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae"`。