Commit 6612f40d authored by Wei Shoulin's avatar Wei Shoulin
Browse files

run_id -> dag_run_id pipeline_run_id -> dag_id

#11599
merged 1 commit into
apache :  master
from
sridhar-rl :  fix-run-id
run_id -> dag_run_id, pipeline_run_id -> dag_id #11599
<here is a image 52677f8958282441-0e06207544895614>
parent 0c4e9bf0
Loading
Loading
Loading
Loading
Loading
+9 −9
Original line number Diff line number Diff line
@@ -125,21 +125,21 @@ def update_qc0_status(level0_id: str, file_type: str, qc0_status: int, dataset:
    """
    return request.put(f"/api/level0/qc0_status/{level0_id}", {'file_type': file_type, 'qc0_status': qc0_status, 'dataset': dataset})

def update_prc_status(level0_id: str, file_type: str, run_id: str, prc_status: int, dataset: str = constants.DEFAULT_DATASET) -> Result:
def update_prc_status(level0_id: str, file_type: str, dag_run_id: str, prc_status: int, dataset: str = constants.DEFAULT_DATASET) -> Result:
    """
    更新0级数据的处理状态
    
    Args:
        level0_id (str): 0级数据的ID
        file_type (str): 文件类型
        run_id (str): 运行ID
        dag_run_id (str): 运行ID
        prc_status (int): 处理状态
        dataset (str): 数据集名称
    
    Returns:
        Result: 操作结果
    """
    return request.put(f"/api/level0/prc_status/{level0_id}/{run_id}", {'file_type': file_type, 'prc_status': prc_status, 'dataset': dataset})
    return request.put(f"/api/level0/prc_status/{level0_id}/{dag_run_id}", {'file_type': file_type, 'prc_status': prc_status, 'dataset': dataset})

def write(local_file: str, 
        dataset: str = constants.DEFAULT_DATASET,
@@ -232,8 +232,8 @@ def process_list(level0_id: str) -> Result:
    return request.get(f"/api/level0/prc/{level0_id}")

def add_process(level0_id: str, 
                pipeline_id: str, 
                run_id: str, 
                dag_id: str, 
                dag_run_id: str, 
                batch_id: Optional[str] = None,
                dataset: str = constants.DEFAULT_DATASET,
                prc_status: int = -1024,                
@@ -245,8 +245,8 @@ def add_process(level0_id: str,
    
    Args:
        level0_id (str): 0级数据的ID
        pipeline_id (str): 管线ID
        run_id (str): 运行ID
        dag_id (str): 管线ID
        dag_run_id (str): 运行ID
        dataset (str): 数据集
        batch_id (str): 批次ID
        prc_time (str): 处理时间,格式为"YYYY-MM-DD HH:MM:SS"
@@ -260,8 +260,8 @@ def add_process(level0_id: str,
    """
    params = {
        'level0_id': level0_id,
        'pipeline_id': pipeline_id,
        'run_id': run_id,
        'dag_id': dag_id,
        'dag_run_id': dag_run_id,
        'dataset': dataset,
        'batch_id': batch_id,
        'prc_time': prc_time,
+10 −10
Original line number Diff line number Diff line
@@ -157,7 +157,7 @@ def write(local_file: Union[IO, str],
        level1_id: str,
        file_type: str,
        file_name: str,
        pipeline_id: str,
        dag_id: str,
        pmapname: str,
        build: int,
        level0_id: Optional[str] = None,
@@ -175,7 +175,7 @@ def write(local_file: Union[IO, str],
        level1_id (str): 1级数据的ID
        file_type (str): 文件类型
        file_name (str): 1级数据文件名
        pipeline_id (str): 管线ID
        dag_id (str): 管线ID
        pmapname (str): CCDS pmap名称
        build (int): 构建号
        dataset (str): 数据集名称
@@ -192,7 +192,7 @@ def write(local_file: Union[IO, str],
        'level1_id': level1_id, 
        'file_type': file_type, 
        'file_name': file_name, 
        'pipeline_id': pipeline_id, 
        'dag_id': dag_id, 
        'pmapname': pmapname, 
        'build': build,
        'dataset': dataset,
@@ -220,7 +220,7 @@ def generate_prc_msg(module_id: Literal['MSC', 'IFS', 'MCI', 'HSTDM', 'CPIC'],
    Args:
        module_id (str): 模块ID
        level1_id (str): 1级数据的ID
        pipeline_id (str): 流水管线ID,默认为空字符串
        dag_id (str): 流水管线ID,默认为空字符串
        dataset (str): 数据集
        batch_id (str): 批次ID

@@ -250,8 +250,8 @@ def process_list(level1_id: str) -> Result:
    return request.get(f"/api/level1/prc/{level1_id}")

def add_process(level1_id: str, 
                pipeline_id: str, 
                run_id: str, 
                dag_id: str, 
                dag_run_id: str, 
                dataset: str = constants.DEFAULT_DATASET,
                batch_id: str = constants.DEFAULT_BATCH_ID,
                prc_time: str = utils.get_current_time(),
@@ -263,8 +263,8 @@ def add_process(level1_id: str,
    
    Args:
        level1_id (str): 1级数据的ID
        pipeline_id (str): 管线ID
        run_id (str): 运行ID
        dag_id (str): 管线ID
        dag_run_id (str): 运行ID
        dataset (str): 数据集
        batch_id (str): 批次ID
        prc_time (str): 处理时间,格式为"YYYY-MM-DD HH:MM:SS"
@@ -278,8 +278,8 @@ def add_process(level1_id: str,
    """
    params = {
        'level1_id': level1_id,
        'pipeline_id': pipeline_id,
        'run_id': run_id,
        'dag_id': dag_id,
        'dag_run_id': dag_run_id,
        'dataset': dataset,
        'batch_id': batch_id,
        'prc_time': prc_time,
+3 −3
Original line number Diff line number Diff line
@@ -150,7 +150,7 @@ def write(local_file: Union[IO, str],
        level2_id: str,
        data_type: str,
        file_name: str,
        pipeline_id: str,
        dag_id: str,
        build: int,
        level0_id: Optional[str] = None,
        level1_id: Optional[str] = None,   
@@ -168,7 +168,7 @@ def write(local_file: Union[IO, str],
        level2_id (str): 2级数据的ID
        data_type (str): 数据类型,如'csst-msc-l2-mbi-cat'
        file_name (str): 2级数据文件名
        pipeline_id (str): 管线ID
        dag_id (str): 管线ID
        build (int): 构建号
        level0_id (Optional[str]): 0级数据的ID默认为 None
        level1_id (Optional[str]): 1级数据的ID默认为 None
@@ -191,7 +191,7 @@ def write(local_file: Union[IO, str],
        'brick_id': brick_id,
        'file_name': file_name,
        'data_type': data_type,
        'pipeline_id': pipeline_id,
        'dag_id': dag_id,
        'build': build,
        'dataset': dataset,
        'batch_id': batch_id,
+2 −2
Original line number Diff line number Diff line
@@ -136,8 +136,8 @@ print(result)# result.data为list类型,包含0级数据处理记录列表
from csst_dfs_client import level0
result = level0.add_process(
      level0_id='some_level0_id', 
      pipeline_id='some_pipeline_id', 
      run_id='some_run_id', 
      dag_id='some_dag_id', 
      dag_run_id='some_dag_run_id', 
      batch_id='some_batch_id',
      dataset='default_dataset',
      prc_status=1,
+4 −4
Original line number Diff line number Diff line
@@ -117,7 +117,7 @@ module_id = 'MSC'
level1_id = 'some_level1_id'
file_type = 'SCI'
file_name = 'level1_data.fits'
pipeline_id = 'some_pipeline_id'
dag_id = 'some_dag_id'
pmapname = 'some_pmapname'
build = 1
level0_id = 'some_level0_id'
@@ -125,7 +125,7 @@ dataset = 'default_dataset'
batch_id = 'default_batch_id'
qc1_status = 0

result = level1.write(local_file_path, module_id, level1_id, file_type, file_name, pipeline_id, pmapname, build, level0_id, dataset, batch_id, qc1_status)
result = level1.write(local_file_path, module_id, level1_id, file_type, file_name, dag_id, pmapname, build, level0_id, dataset, batch_id, qc1_status)
print(result)

```
@@ -152,8 +152,8 @@ print(result) # result.data为list类型,包含1级数据处理记录列表
from csst_dfs_client import level1
result = level1.add_process(
      level1_id='some_level1_id', 
      pipeline_id='some_pipeline_id', 
      run_id='some_run_id', 
      dag_id='some_dag_id', 
      dag_run_id='some_dag_run_id', 
      batch_id='some_batch_id',
      dataset='default_dataset',
      prc_status=1,
Loading