initial stage

This commit is contained in:
Giambattista Bloisi 2024-08-01 11:04:06 +02:00
parent 9581a86313
commit b23ddd3002
1 changed files with 3 additions and 0 deletions

View File

@ -6,6 +6,7 @@ from datetime import timedelta
import pendulum
from airflow.decorators import dag, task_group
from airflow.decorators import task
from airflow.exceptions import AirflowSkipException
from airflow.operators.empty import EmptyOperator
from airflow.operators.python import get_current_context
from airflow.utils.helpers import chain
@ -51,6 +52,8 @@ def import_s3_openaire_dump():
def compute_batches():
nonlocal entity
kwargs = get_current_context()
if entity not in kwargs["params"]["ENTITIES"]:
raise AirflowSkipException(f"Skipping {entity}")
return [[(entity, '1'), (entity, '2')], [], []]
@task(executor_config={