Commit 129bbfe9 authored by BO ZHANG's avatar BO ZHANG 🏀
Browse files

feat(v2): 引入新版 DAG 框架并重构项目结构

parent edd8a5a2
Loading
Loading
Loading
Loading
+137 −0
Original line number Diff line number Diff line
@@ -16,6 +16,143 @@ pip install git+https://csst-tb.bao.ac.cn/code/csst-cicd/csst-dag.git

基本用法

### v2: 通过 `DagRunGroup.trigger()` 生成 `DagRunGroup` 对象

当前 v2 推荐入口为 `DagRunGroup.trigger()`

```python
from csst_dag.v2.models import DagRunGroup

group = DagRunGroup.trigger(
    dag="csst-msc-l1-mbi",
    batch_id="manual-test",
    priority="low",
    data={
        "data_model": "raw",
        "dataset": "test-msc-c9-25sqdeg-v3",
        "instrument": "MSC",
        "obs_type": "WIDE",
        "obs_group": "W5",
    },
    proc={
        "pmapname": "csst_000155.pmap",
        "ref_cat": None,
        "extra_kwargs": {},
    },
)

# 直接打印对象,会显示 group 头部 + 前 2 个 dag run 预览
print(group)

# 如果需要提交/序列化,再转成 payload
payload = group.to_payload()
```

#### 需要多少参数

`trigger()` 的函数签名是:

```python
DagRunGroup.trigger(
    dag: str,
    batch_id: str | None = None,
    priority: str = "low",
    data: dict | None = None,
    proc: dict | None = None,
)
```

最少只需要:

- `dag`

但如果希望真正匹配出一批可用的 `dag_runs`,通常至少还需要在 `data` 中提供:

- `dataset`
- `data_model`
- `instrument`
- `obs_type`
- `obs_group`

对大多数 v2 `file` 类型 DAG,上面这 5 个 `data` 字段已经是“实际可用的最小集合”。

#### `data` 支持的参数

`data` 中当前支持这些键;未被 DFS 使用的多余键会被忽略:

- `data_model`
- `dataset`
- `instrument`
- `obs_type`
- `obs_group`
- `obs_id`
- `detector`
- `filter`
- `custom_id`
- `prc_status`
- `qc_status`
- `batch_id`
- `priority`
- `force_success`
- `return_data_list`

说明:

- `data_model="raw"` 时,底层会走 level0 数据查询。
-`raw` 时,会走 level1 数据查询。
- `batch_id` 也可以只从 `trigger(..., batch_id=...)` 顶层传入,框架会自动补到 `data` 中。

#### `proc` 支持的参数

- `pmapname`
- `ref_cat`
- `extra_kwargs`

#### 返回值

`trigger()` 返回的是一个 `DagRunGroup` 对象,而不是普通字典。

- `print(group)`:显示精简预览
- `group.dag_runs`:访问子任务对象列表
- `group.to_payload()`:转成对外提交格式

`group.to_payload()` 的结构为:

```json
{
  "group": {
    "dag": "csst-msc-l1-mbi",
    "dag_run_group": "019d...",
    "batch_id": "manual-test",
    "priority": "low",
    "created_time": "2026-04-14T05:20:55.568",
    "n_dag_runs": 720
  },
  "dag_runs": [
    {
      "job": { "...": "..." },
      "data": { "...": "..." },
      "aux": { "...": "..." },
      "proc": { "...": "..." }
    }
  ]
}
```

#### 手动测试脚本

仓库中提供了一个手动测试脚本:

```bash
python3 /nfs/airflow/csst-airflow/external/csst-dag/test_v2/run_dag_group_manual.py
```

默认会:

- 生成一个真实的 `DagRunGroup`
- `print(group)` 显示精简预览
- 如有需要,可在脚本中打开 `group.to_payload()` 的打印

```python
from csst_dag import CSST_DAGS

+13 −3
Original line number Diff line number Diff line
from .dfs import DFS
from .dag import CSST_DAGS, Dispatcher, BaseDAG, Level1DAG, Level2DAG
from ._csst import csst, CsstPlanObsid, CsstPlanObsgroup, DotDict
"""CSST DAG 包。

说明
----
该包同时提供 `v1`(legacy)与 `v2`(新 payload/models)两套并列实现。
项目代码应优先从 `csst_dag.v1.*` 或 `csst_dag.v2.*` 引用,避免依赖顶层 legacy 路径。
"""

__version__ = "0.0.0"

from . import v1, v2

__all__ = [
    "v1",
    "v2",
]
+524 B

File added.

No diff preview for this file type.

+196 B

File added.

No diff preview for this file type.

+8.78 KiB

File added.

No diff preview for this file type.

Loading