Commit 2b89f6a6 authored by BO ZHANG's avatar BO ZHANG 🏀
Browse files

refactor: 统一重命名source_batch_id为upstream_batch_id并重构批次处理逻辑

parent de4b2b76
Loading
Loading
Loading
Loading
+7 −7
Original line number Diff line number Diff line
@@ -92,7 +92,7 @@ DagRunGroup.trigger(
- `custom_id`
- `prc_status`
- `qc_status`
- `source_batch_id`
- `upstream_batch_id`
- `priority`
- `force_success`
- `return_data_list`
@@ -102,9 +102,9 @@ DagRunGroup.trigger(
- `data_model="raw"` 时,底层会走 level0 数据查询。
- 通用 `data.find()` 对非 `raw` 数据会按 DFS `level0/1/2.list_data_models()` 动态返回的存储层级分流到 level1 或 level2。
- `batch_id` 顶层参数表示当前新生成的 `DagRunGroup` / 输出产品所属批次。
- `data.source_batch_id` 表示检索上游 level1 数据时使用的批次条件。
- `data.source_batch_id` 为空时,框架会默认回退到顶层 `batch_id`
- 为兼容旧调用方式,若调用方仍传 `data.batch_id`,当前实现会将其视为 `source_batch_id` 的别名。
- `data.upstream_batch_id` 表示检索上游数据时使用的批次条件。
- `data.upstream_batch_id` 顶层 `batch_id` 语义分离:前者用于上游检索,后者用于当前新生成 group 的批次
- 当前实现不会再将顶层 `batch_id` 回退为 `upstream_batch_id`,也不再把 `data.batch_id` 当作它的别名。

Brick / healpix 查询补充说明:

@@ -160,7 +160,7 @@ submission:
      fixed: true
    dataset:
      value: "test-msc-c9-25sqdeg-v3"
    source_batch_id:
    upstream_batch_id:
      value: "default"
    instrument:
      value: "MSC"
@@ -185,11 +185,11 @@ submission:

约定说明:

- `submission.data.<field>` 用于配置 `data` 组字段,例如 `data_model``dataset``source_batch_id``instrument``obs_type``healpix``custom_id`
- `submission.data.<field>` 用于配置 `data` 组字段,例如 `data_model``dataset``upstream_batch_id``instrument``obs_type``healpix``custom_id`
- `submission.proc.<field>` 用于配置 `proc` 组字段,例如 `pmapname``ref_cat``extra_kwargs`
- `data_model` 虽然属于 `data` 检索条件,但通常和 DAG 定义强绑定,因此推荐在 YAML 中显式给出。
- 对 Brick DAG 来说,这里的 `data_model` 表示“生成任务时要检索的上游数据产品”,不是“决定查询访问 level1 还是 level2 的开关”。
- `source_batch_id` 建议只在非 `raw` 的 level1 / L2 检索场景中使用;raw 数据场景通常不需要该字段。
- `upstream_batch_id` 建议只在非 `raw` 的 level1 / L2 检索场景中使用;raw 数据场景通常不需要该字段。
- 大部分 L1 DAG 推荐固定为 `data_model.value: "raw"``fixed: true`
- 某些 L2 / brick DAG 需要固定到特定 level1 模型,例如:
  - `csst-msc-l2-mbi-mosaic` 推荐固定为 `csst-msc-l1-mbi`
+1 −1
Original line number Diff line number Diff line
@@ -10,7 +10,7 @@ DAG_CONFIG_DIR_V2 = os.path.join(
)
ALIGNED_SUBMISSION_DATA_FIELDS = [
    "dataset",
    "source_batch_id",
    "upstream_batch_id",
    "data_model",
    "instrument",
    "obs_type",
+4 −0
Original line number Diff line number Diff line
@@ -34,6 +34,10 @@ CSST_DAGS = {
        "csst-msc-l1-ooc",
        dispatcher=Dispatcher.dispatch_obsgroup_detector,
    ),
    "csst-msc-l1-qc0": FileDAG(
        "csst-msc-l1-qc0",
        dispatcher=Dispatcher.dispatch_file,
    ),
    "csst-mci-l1": FileDAG(
        "csst-mci-l1",
        dispatcher=Dispatcher.dispatch_file,
+2 −3
Original line number Diff line number Diff line
@@ -115,7 +115,7 @@ class BaseDAG(ABC):
                return None
            return str(v)

        source_batch_id = _as_str(data.get("source_batch_id")) or _as_str(data.get("batch_id"))
        upstream_batch_id = _as_str(data.get("upstream_batch_id"))

        plan_query = {
            "dataset": _as_str(data.get("dataset")),
@@ -140,8 +140,7 @@ class BaseDAG(ABC):
            "custom_id": _as_str(data.get("custom_id")),
            "prc_status": _normalize_optional_scalar(data.get("prc_status")),
            "qc_status": _normalize_optional_scalar(data.get("qc_status")),
            "batch_id": source_batch_id,
            "pmapname": _as_str(proc.get("pmapname")),
            "batch_id": upstream_batch_id,
            "page": 1,
            "limit": 0,
        }
+4 −4
Original line number Diff line number Diff line
@@ -156,7 +156,7 @@ class BaseBrickDAG(BaseDAG):
    ) -> DagRunGroup:
        dag_group = DagRunGroup(
            dag=self.dag_name,
            batch_id=data.get("group_batch_id", data.get("batch_id")),
            batch_id=data.get("batch_id"),
            priority=data.get("priority", "low"),
            docker_images=docker_images,
        )
@@ -188,9 +188,9 @@ class BaseBrickDAG(BaseDAG):
            "priority": dag_group.get("priority", "low"),
            "data_model": (str(data.get("data_model")) if data.get("data_model") not in (None, "") else None),
            "dataset": (str(data.get("dataset")) if data.get("dataset") not in (None, "") else None),
            "source_batch_id": (
                str(data.get("source_batch_id"))
                if data.get("source_batch_id") not in (None, "")
            "upstream_batch_id": (
                str(data.get("upstream_batch_id"))
                if data.get("upstream_batch_id") not in (None, "")
                else None
            ),
            "instrument": (str(data.get("instrument")) if data.get("instrument") not in (None, "") else None),
Loading