Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -29,15 +29,15 @@


@task.branch()
def should_run(**kwargs) -> str:
def should_run(logical_date=None) -> str:
Comment thread
amoghrajesh marked this conversation as resolved.
"""
Determine which empty_task should be run based on if the logical date minute is even or odd.

:param dict kwargs: Context
:param pendulum.DateTime logical_date: The logical date for the current execution
:return: Id of the task to run
"""
print(f"------------- exec dttm = {kwargs['logical_date']} and minute = {kwargs['logical_date'].minute}")
if kwargs["logical_date"].minute % 2 == 0:
print(f"------------- exec dttm = {logical_date} and minute = {logical_date.minute}")
if logical_date.minute % 2 == 0:
return "empty_task_1"
return "empty_task_2"

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,16 +52,14 @@
) as dag:

@task(task_id="get_names", task_display_name="Get names")
def get_names(**kwargs) -> list[str]:
params = kwargs["params"]
def get_names(params=None) -> list[str]:
if "names" not in params:
print("Uuups, no names given, was no UI used to trigger?")
return []
return params["names"]

@task.branch(task_id="select_languages", task_display_name="Select languages")
def select_languages(**kwargs) -> list[str]:
params = kwargs["params"]
def select_languages(params=None) -> list[str]:
selected_languages = []
for lang in ["english", "german", "french"]:
if params[lang]:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -316,8 +316,7 @@ def _get_script_interfaces_from_config() -> list[str]:
) as dag:
# [START section_3]
@task(task_display_name="Show used parameters")
def show_params(**kwargs) -> None:
params = kwargs["params"]
def show_params(params=None) -> None:
print(f"This DAG was triggered with the following parameters:\n\n{json.dumps(params, indent=4)}\n")

show_params()
Expand Down
9 changes: 3 additions & 6 deletions airflow-core/src/airflow/example_dags/tutorial_dag.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,16 +57,14 @@
# [END documentation]

# [START extract_function]
def extract(**kwargs):
ti = kwargs["ti"]
def extract(ti=None):
data_string = '{"1001": 301.27, "1002": 433.21, "1003": 502.22}'
ti.xcom_push("order_data", data_string)

# [END extract_function]

# [START transform_function]
def transform(**kwargs):
ti = kwargs["ti"]
def transform(ti=None):
extract_data_string = ti.xcom_pull(task_ids="extract", key="order_data")
order_data = json.loads(extract_data_string)

Expand All @@ -81,8 +79,7 @@ def transform(**kwargs):
# [END transform_function]

# [START load_function]
def load(**kwargs):
ti = kwargs["ti"]
def load(ti=None):
total_value_string = ti.xcom_pull(task_ids="transform", key="total_order_value")
total_order_value = json.loads(total_value_string)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,16 +64,14 @@ def tutorial_objectstorage():

# [START get_air_quality_data]
@task
def get_air_quality_data(**kwargs) -> ObjectStoragePath:
def get_air_quality_data(logical_date=None) -> ObjectStoragePath:
"""
#### Get Air Quality Data
This task gets air quality data from the Finnish Meteorological Institute's
open data API. The data is saved as parquet.
"""
import pandas as pd

logical_date = kwargs["logical_date"]

latitude = 28.6139
longitude = 77.2090

Expand Down