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

feat(api-gateway): 新增统一API网关以支持海量任务提交与异构JSON检索

parent f1f1e3df
Loading
Loading
Loading
Loading
+1 −0
Original line number Diff line number Diff line
@@ -48,3 +48,4 @@ dist/

# Database
db.sqlite3
cookies.txt
+13 −12
Original line number Diff line number Diff line
@@ -14,6 +14,7 @@
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 和鉴权逻辑。

## 部署步骤

@@ -54,7 +55,7 @@ make fix-permission
make init
```

### 5. 启动所有服务
### 5. 启动所有核心服务 (含 API Gateway)

> **⚠️ 部署强烈建议**:
> 如果您是在单台机器上进行完整的测试或生产部署,**请不要只使用 `make up`**,因为这只会启动调度器而不会启动任何计算节点(任务将永远卡在排队中)。
@@ -66,11 +67,13 @@ make init
> ```

#### 命令区别与集群部署指南:
*   **`make up`**:启动核心调度服务(Scheduler、API Server、DAG Processor)以及数据库组件(Postgres、Redis、Elasticsearch、Kibana 等)。**注意:它不会启动实际执行任务的 Celery Worker。**
*   **`make up`**:启动核心调度服务(Scheduler、API Server、DAG Processor)、API Gateway (网关 38000 端口) 以及数据库组件(Postgres、Redis、Elasticsearch、Kibana 等)。**注意:它不会启动实际执行任务的 Celery Worker。**
*   **`make up-master`**:不仅执行 `make up` 的所有内容,还会**额外启动 Flower 服务**(一个用于监控 Celery 队列状态的 Web UI)。
    *   *为什么 Flower 不作为默认启动?* 因为它是可选的可视化组件。Airflow 核心调度并不依赖它。官方出于节省服务器资源和安全端口暴露的考量,将其设置为按需启动(Profile隔离)
    *   *为什么 Flower 不作为默认启动?* 因为它是可选的可视化组件。官方出于节省服务器资源和安全考量,将其设置为按需启动。
*   **`make up-worker`****专门用于启动 Celery Worker 容器**
    *   *分布式部署场景*:您可以在多台不同的机器上克隆本代码,并**仅运行** `make up-worker` 来横向扩展计算能力(注意需修改 `.env` 中的 `MASTER_IP` 指向主节点)。
    *   *分布式部署场景*:您可以在多台不同的机器上克隆本代码,并**仅运行** `make up-worker` 来横向扩展计算能力。

---

## Makefile 命令速查

@@ -81,7 +84,7 @@ make init
| `make fix-permission` | 修复 `es_data` 目录的属主权限 (UID 1000) |
| `make copy-env-*` | 复制对应的预设环境变量文件 |
| `make init` | 初始化 Airflow 数据库 |
| `make up` | 启动核心调度器与基础设施 (不含 Worker) |
| `make up` | 启动核心调度器、API 网关与基础设施 (不含 Worker) |
| `make up-master` | 启动核心组件并附带开启 Flower 监控 |
| `make up-worker` | 启动 Celery Worker 以执行具体的 Task |
| `make down` | 停止所有服务 |
@@ -90,6 +93,7 @@ make init

## 访问服务与系统

- **统一 API Gateway**: `http://<MASTER_IP>:38000` (业务系统唯一对接入口)
- **Airflow Web UI**: `http://<MASTER_IP>:38080` (账号/密码默认: `airflow`/`airflow`)
- **Kibana (日志中心)**: `http://<MASTER_IP>:35601` (直接在 Discover 页面选择 `airflow-*` 视图检索日志)
- **Elasticsearch API**: `http://<MASTER_IP>:39200`
@@ -97,15 +101,12 @@ make init
- **PostgreSQL**: `localhost:35432`
- **Redis**: `localhost:36379`

## DAG 状态查询示例
## DAG 任务的触发与状态查询
为了避免 Airflow 3 升级鉴权体系(强制 JWT)带来的对接复杂性,以及为了解决海量异构任务数据的检索问题,本系统**强力推荐业务方直接对接自定义的 API Gateway**

由于 Airflow 3 升级了鉴权体系(强制 JWT),如果想在外部系统(如自建看板)查询 DAG 运行状态,推荐使用 Docker CLI 获取最稳定可靠的 JSON 状态数据
外部系统仅需维护一个环境变量 `export AIRFLOW_API_GATEWAY=http://<MASTER_IP>:38000`,即可进行批量任务提交、JSONB 高性能状态检索等操作

可直接运行仓库内提供的演示脚本进行参考:
```bash
python query_status_demo.py
```
它演示了如何查询 `csst-msc-l1-mbi``csst-msc-l1-qc0` 的各步骤运行情况。
请查阅配套的 **[`api.md`](./api.md)** 获取完整的接口文档和 cURL 调用示例。

## 故障排除

+149 −0
Original line number Diff line number Diff line
# CSST Airflow 统一调度网关 API 接口文档

基于最新架构升级方案,目前系统对外统一通过 **FastAPI 网关 (API Gateway)** 提供服务。所有任务的触发、查询与追踪均通过网关进行,不再需要业务端去直接连接 Airflow REST API 或 Elasticsearch。所有的任务数据(异构输入、状态、结果)均使用 `PostgreSQL JSONB` 结构聚合存储。

> **提示**:API 网关默认运行在 38000 端口。您不再需要通过 `airflow/airflow` 获取 JWT Token。

---

## 0. 快速设置环境变量 (快捷方式)
为了方便后续测试调用,你可以直接在 Bash 中设置网关的环境变量 `AIRFLOW_API_GATEWAY`

```bash
export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000
```
*(执行后,后续示例中均可直接使用 `$AIRFLOW_API_GATEWAY` 进行请求,无需手动修改 IP 和端口。)*

## 1. 批量任务触发 (Batch Submit Tasks)

向系统提交任务。您可以直接发送您业务中定义的 DAG 执行消息(**必须包含 `dag_id` 和 `dag_run_id`**),网关会直接复用该 `dag_run_id` 并在后台异步推送。整个 JSON 消息体会被完整存入底层数据库以供后续检索。如果重复提交相同 `dag_run_id`,网关会忽略重复记录,实现幂等。

*   **Endpoint**: `POST /api/tasks/batch-submit`
*   **请求头**: `Content-Type: application/json`
*   **cURL 示例**:
    ```bash
    curl -X POST "$AIRFLOW_API_GATEWAY/api/tasks/batch-submit" \
      -H "Content-Type: application/json" \
      -d '[
        {
          "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"
          }
        }
      ]'
    ```
*   **测试结果 (成功响应示例)**:
    ```json
    {
      "status": "accepted",
      "task_ids": [
        "123e4567-e89b-12d3-a456-426614174000"
      ]
    }
    ```

---

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

利用任务提交时你传入的那个 `dag_run_id`,直接获取任务的实时状态、执行结果或错误日志。数据源自数据库的精确查询($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"
    ```
*   **测试结果 (成功响应示例)**:
    ```json
    {
      "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": "10100131914",
        "detector": "09",
        "batch_id": "inttest",
        "docker_images": {
          "csst-msc-l1-mbi": "latest"
        }
      },
      "outputs": {
        "finished_at": "2026-04-01T15: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"
    }
    ```
*   **状态枚举 (`status`)**:
    *   `received`: 网关已接收,正在等待进入 Airflow 队列。
    *   `queued` / `running`: Airflow 正在排队或执行中(待扩展实时更新)。
    *   `success`: 任务执行成功。
    *   `failed`: 任务执行失败。如果是失败重试机制,此状态为**最终状态**,且 `error_message` 会包含异常堆栈。

---

## 3. 基于异构 JSON 字段的批量检索 (Search by JSON Inputs)

得益于底层 `PostgreSQL JSONB + GIN 索引` 的设计,您可以随意根据您当初提交的 `inputs` 中的任意嵌套字段来进行高性能的反向检索,这对于批量查询非常有用。

*   **Endpoint**: `POST /api/tasks/search`
*   **参数**: `?limit=100&offset=0`
*   **cURL 示例** (查询所有 inputs 中 `obs_id``10100131914``batch_id``inttest` 的任务):
    ```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"
      }'
    ```
*   **测试结果 (成功响应示例)**:
    ```json
    [
      {
        "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": "10100131914",
          "detector": "09",
          "batch_id": "inttest",
          "docker_images": {
            "csst-msc-l1-mbi": "latest"
          }
        }
      }
    ]
    ```

---

## 4. 关于深度日志查询

目前 API 网关支持返回关键错误信息 (`error_message`) 和简要摘要 (`outputs`)。如果您需要对失败的任务进行排错,可以通过上面返回的 `task_id`(它同时也是 Airflow 的 `dag_run_id`)前往 Elasticsearch 或 Kibana 进行深入的全量全文日志检索:

*   **Kibana URL**: `http://10.73.0.27:35601`
*   **检索条件示例**: 在 Kibana 的 `filebeat-*` 索引中搜索 `"123e4567-e89b-12d3-a456-426614174000"`
 No newline at end of file
+12 −0
Original line number Diff line number Diff line
FROM python:3.11-slim

WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY . .

# 暴露网关端口
EXPOSE 38000

CMD ["python", "main.py"]
+25 −0
Original line number Diff line number Diff line
# 统一 API 网关 (API Gateway)

这是一个基于 FastAPI 编写的轻量级 API 网关,作为 Airflow 的前置缓冲和统一数据聚合层。

## 功能特性
1. **扁平化任务提交**: 直接接收业务方原始 JSON,无需额外包裹结构。提取 `dag_id``dag_run_id` 后原样落库。
2. **高性能批量触发**: 支持瞬间接收 100,000+ 个任务,存入 Postgres 后通过后台异步推送给 Airflow。
3. **状态与日志聚合**: 结合 Airflow Callback,任务的状态、输入、输出和简短日志全部聚合在同一张 Postgres 表中。
4. **异构 JSON 查询**: 利用 PostgreSQL 的 JSONB 和 GIN 索引,支持对任何异构的输入参数字段进行快速反向查询。

## 部署运行

目前该网关**已经集成到了主项目的 `docker-compose.yaml` 中**
当在主项目目录中运行 `make up``make up-master``docker-compose up -d` 时,网关容器会自动构建并启动。

### 手动调试模式
如果你需要在本地脱离 Docker 进行开发调试,可以使用以下命令:
```bash
cd api_gateway
pip install -r requirements.txt
python main.py
```
> **注意**:本地调试时,请确保 `main.py` 和 `database.py` 里的环境变量(如 `POSTGRES_HOST`, `AIRFLOW_HOST`)指向了正确的测试环境 IP(比如 `10.73.0.27`)。

默认运行在 `0.0.0.0:38000`。外部系统可以将其设置为环境变量(如 `export AIRFLOW_API_GATEWAY=http://10.73.0.27:38000`)进行对接。
 No newline at end of file
Loading