Commit d844b898 authored by BO ZHANG's avatar BO ZHANG 🏀
Browse files

feat(ui): 新增DAGs与Workers监控页面并支持自定义worker并发度

parent e8a04b82
Loading
Loading
Loading
Loading
+5 −0
Original line number Diff line number Diff line
@@ -18,6 +18,11 @@

如果您需要在真实的集群(1个 Master 节点 + 多个 Worker 节点)上部署该系统,我们提供了开箱即用的 Ansible Playbooks。由于我们有多个不同的部署环境(如 `p368`, `csu`, `zjlab`),我们为每个环境准备了专属的 Ansible Inventory 文件。

**💡 进阶:自定义 Worker 节点的并发度**
您可以在 Inventory 文件(例如 `inventory.p368.ini`)中,为每台 Worker 机器灵活配置变量:
- `worker_concurrency`: 设置该台机器上 Worker 容器的 Slots 数量(例如 16 或 64)。
*(通过这种方式,您可以让高配物理机加大并发,低配物理机少跑,充分压榨集群性能!)*

#### 1. 准备工作
- **安装 Ansible**:如果您的执行机(通常是您的本地电脑或 Master 节点)尚未安装 Ansible,请先通过以下命令安装:

+43 −127
Original line number Diff line number Diff line
@@ -26,70 +26,19 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000
      -H "Content-Type: application/json" \
      -d '{
        "dag_group_run": {
          "dag_group": "default",
          "dag_group_run": "18a2662c9ea9948a0f802e007ecd1ce6ba4e927e",
          "batch_id": "default",
          "priority": 1,
          "created_time": "2026-04-02T06:21:41.675"
          "batch_id": "inttest",
          "obs_group": "W5",
          "dataset": "test-msc-c9-25sqdeg-v3"
        },
        "dag_run_list": [
          {
            "dataset": "test-msc-c9-25sqdeg-v3",
            "instrument": "MSC",
            "obs_type": "WIDE",
            "obs_group": "W5",
            "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",
            "dag_run": "123e4567-e89b-12d3-a456-426614174000",
            "obs_id": "10100131914",
            "detector": "09",
            "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": ""
            "dag_group_run": "group_run_123e4567"
          }
        ]
      }'
@@ -99,8 +48,7 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000
    {
      "status": "accepted",
      "task_ids": [
        "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae",
        "a992e82967594fa4134a0f6bcb0b4172778d6780"
        "123e4567-e89b-12d3-a456-426614174000"
      ]
    }
    ```
@@ -121,13 +69,13 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000
    ```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,
        "dag_group_run": "group_run_123e4567",
        "batch_id": "inttest",
        "created_at": "2026-04-01T15:00:00.000000",
        "total_tasks": 100,
        "success_tasks": 90,
        "failed_tasks": 2,
        "running_tasks": 8,
        "received_tasks": 0,
        "cancelled_tasks": 0
      }
@@ -138,55 +86,40 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000

## 3. 单个任务状态查询 (Query Task Status)

利用任务提交时你传入的那个 `dag_run`,直接获取任务的实时状态、执行结果或错误日志。数据源自数据库的精确查询($O(1)$ 复杂度)。
利用任务提交时你传入的那个 `dag_run_id`,直接获取任务的实时状态、执行结果或错误日志。数据源自数据库的精确查询($O(1)$ 复杂度)。

*   **Endpoint**: `GET /api/tasks/{dag_run_id}`
*   **cURL 示例**:
    ```bash
    curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae"
    curl -s -X GET "$AIRFLOW_API_GATEWAY/api/tasks/123e4567-e89b-12d3-a456-426614174000"
    ```
*   **测试结果 (成功响应示例)**:
    ```json
    {
      "task_id": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae",
      "task_id": "123e4567-e89b-12d3-a456-426614174000",
      "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": "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": ""
        "obs_id": "10100131914",
        "detector": "09",
        "batch_id": "inttest",
        "docker_images": {
          "csst-msc-l1-mbi": "latest"
        }
      },
      "outputs": {
        "finished_at": "2026-04-02T06:30:00.123456"
        "finished_at": "2026-04-01T15:30:00.123456"
      },
      "error_message": null,
      "execution_logs": null,
      "created_at": "2026-04-02T06:22:33.846000Z",
      "updated_at": "2026-04-02T06:30:00.123456Z"
      "created_at": "2026-04-01T15:00:00.000000Z",
      "updated_at": "2026-04-01T15:30:00.123456Z"
    }
    ```
*   **状态枚举 (`status`)**:
@@ -205,50 +138,33 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000

*   **Endpoint**: `POST /api/tasks/search`
*   **参数**: `?limit=100&offset=0`
*   **cURL 示例** (查询所有 inputs 中 `obs_id` 包含 `547339``batch_id``default` 的任务):
*   **cURL 示例** (查询所有 inputs 中 `obs_id` 包含 `131914``batch_id``inttest` 的任务):
    ```bash
    curl -s -X POST "$AIRFLOW_API_GATEWAY/api/tasks/search?limit=10" \
      -H "Content-Type: application/json" \
      -d '{
        "obs_id": "547339",
        "batch_id": "default"
        "obs_id": "131914",
        "batch_id": "inttest"
      }'
    ```
*   **测试结果 (成功响应示例)**:
    ```json
    [
      {
        "task_id": "07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae",
        "task_id": "123e4567-e89b-12d3-a456-426614174000",
        "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": "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": ""
          "obs_id": "10100131914",
          "detector": "09",
          "batch_id": "inttest",
          "dag_group_run": "group_run_123e4567"
        }
      }
    ]
@@ -269,12 +185,12 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000
    ```json
    {
      "fields": {
        "batch_id": ["default"],
        "batch_id": ["inttest", "test2"],
        "dataset": ["test-msc-c9-25sqdeg-v3"],
        "detector": ["24", "25"],
        "detector": ["09", "10"],
        "instrument": ["MSC"],
        "obs_group": ["W5"],
        "obs_id": ["10100547339"],
        "obs_id": ["10100131914", "100000123", "100000124"],
        "obs_type": ["WIDE"]
      }
    }
@@ -287,4 +203,4 @@ export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000
目前 API 网关的 `GET /api/tasks/{dag_run_id}` 接口已集成了从 Elasticsearch 拉取全量执行日志 (`execution_logs`) 的能力。如果您还需要进入 Kibana 进行更复杂的汇聚分析或全栈链路追踪,可以:

*   **Kibana URL**: `http://10.73.0.27:35601`
*   **检索条件示例**: 在 Kibana 的 `filebeat-*` 索引中搜索 `"07c4d9e5ce555ff5bdb6fd2cb785b059e5fc03ae"`
*   **检索条件示例**: 在 Kibana 的 `filebeat-*` 索引中搜索 `"123e4567-e89b-12d3-a456-426614174000"`
 No newline at end of file
+15 −14
Original line number Diff line number Diff line
@@ -4,8 +4,8 @@
  become: yes
  vars:
    # Set default if not provided in inventory
    env_name: "{{ env_name | default('p368') }}"
    deploy_dir: "{{ deploy_dir | default('/opt/csst-airflow') }}"
    _env_name: "{{ env_name | default('p368') }}"
    _deploy_dir: "{{ deploy_dir | default('/opt/csst-airflow') }}"
    local_project_dir: ".."

  tasks:
@@ -21,14 +21,14 @@

    - name: Ensure deployment directory exists
      file:
        path: "{{ deploy_dir }}"
        path: "{{ _deploy_dir }}"
        state: directory
        mode: '0755'

    - name: Synchronize project files to remote nodes
      synchronize:
        src: "{{ local_project_dir }}/"
        dest: "{{ deploy_dir }}/"
        src: "{{ playbook_dir }}/../"
        dest: "{{ _deploy_dir }}/"
        delete: no
        recursive: yes
        rsync_opts:
@@ -42,7 +42,7 @@

    - name: Ensure volume directories exist
      file:
        path: "{{ deploy_dir }}/volumes/{{ item }}"
        path: "{{ _deploy_dir }}/volumes/{{ item }}"
        state: directory
        mode: '0777'
      loop:
@@ -56,24 +56,25 @@
  hosts: master
  become: yes
  vars:
    env_name: "{{ env_name | default('p368') }}"
    deploy_dir: "{{ deploy_dir | default('/opt/csst-airflow') }}"
    _env_name: "{{ env_name | default('p368') }}"
    _deploy_dir: "{{ deploy_dir | default('/opt/csst-airflow') }}"
  tasks:
    - name: Start Master services (postgres, redis, elasticsearch, airflow core, api, etc.)
      shell: |
        {{ docker_compose_cmd }} --env-file envs/{{ env_name }}.env --profile master up -d
        {{ docker_compose_cmd }} --env-file envs/{{ _env_name }}.env --profile master up -d
      args:
        chdir: "{{ deploy_dir }}"
        chdir: "{{ _deploy_dir }}"

- name: Deploy Worker Node Services
  hosts: worker
  become: yes
  vars:
    env_name: "{{ env_name | default('p368') }}"
    deploy_dir: "{{ deploy_dir | default('/opt/csst-airflow') }}"
    _env_name: "{{ env_name | default('p368') }}"
    _deploy_dir: "{{ deploy_dir | default('/opt/csst-airflow') }}"
  tasks:
    - name: Start Worker services (celery-worker, filebeat, monitoring)
      shell: |
        {{ docker_compose_cmd }} --env-file envs/{{ env_name }}.env --profile worker up -d
        export AIRFLOW__CELERY__WORKER_CONCURRENCY={{ worker_concurrency | default(16) }}
        {{ docker_compose_cmd }} --env-file envs/{{ _env_name }}.env --profile worker up -d
      args:
        chdir: "{{ deploy_dir }}"
        chdir: "{{ _deploy_dir }}"
+4 −5
Original line number Diff line number Diff line
[master]
# Replace with your actual master node IP
p368-master ansible_host=10.73.0.27 ansible_user=root
# Because we are testing on the local p368 machine itself, use ansible_connection=local
p368-master ansible_host=127.0.0.1 ansible_connection=local

[worker]
# Add all your worker node IPs here
# e.g., p368-worker1 ansible_host=10.73.0.28 ansible_user=root
# p368-worker2 ansible_host=10.73.0.29 ansible_user=root
# Deploying 1 worker container on the same physical machine
p368-worker1 ansible_host=127.0.0.1 ansible_connection=local worker_concurrency=16

[all:vars]
# The environment name to deploy (e.g., p368, csu, zjlab)
+4 −0
Original line number Diff line number Diff line
@@ -9,3 +9,7 @@ JWT_SECRET={{ jwt_secret }}
# --- Cluster & Network Config ---
MASTER_IP={{ master_ip }}
HARBOR={{ harbor }}

# --- Celery Worker Config ---
# Defines the concurrency per worker node. Can be overridden in inventory.
WORKER_CONCURRENCY={{ worker_concurrency | default(8) }}
Loading