Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
AIP-72: Allow retrieving Variable from Task Context by amoghrajesh · Pull Request #45431 · apache/airflow · GitHub
Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' AIP-72: Allow retrieving Variable from Task Context by amoghrajesh · Pull Request #45431 · apache/airflow · GitHub
Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' AIP-72: Allow retrieving Variable from Task Context by amoghrajesh · Pull Request #45431 · apache/airflow · GitHub
Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' AIP-72: Allow retrieving Variable from Task Context by amoghrajesh · Pull Request #45431 · apache/airflow · GitHub
Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' AIP-72: Allow retrieving Variable from Task Context by amoghrajesh · Pull Request #45431 · apache/airflow · GitHub
Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' AIP-72: Allow retrieving Variable from Task Context by amoghrajesh · Pull Request #45431 · apache/airflow · GitHub
Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); AIP-72: Allow retrieving Variable from Task Context by amoghrajesh · Pull Request #45431 · apache/airflow · GitHub
Skip to content

AIP-72: Allow retrieving Variable from Task Context - #45431

Merged
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context
Jan 7, 2025
Merged

AIP-72: Allow retrieving Variable from Task Context#45431
amoghrajesh merged 6 commits into
apache:mainfrom
astronomer:AIP72-variables-from-context

Conversation

@amoghrajesh

@amoghrajeshamoghrajesh commented Jan 6, 2025

Copy link
Copy Markdown
Contributor

closes: #45421

Summary of changes

  1. Added a minimal Variable user-facing definition which will be used in DAG files by DAG authors
  2. Added logic to get Variables in the context - both in "value" and "json" format
  • "value" is the raw form
  • "json" is the deserialised json form, we are trying to keep the contract between SDK and API server simple, they interact only in strings and the responsibility of serialising + deserialising lies on the client before sending it to the task sdk, not on the API server. This will enable multi language support too.

Object Glossary

-VariableResponse is auto-generated and tightly coupled with the API schema.
-VariableResult is runtime-specific and meant for internal communication between Supervisor & Task Runner.
-Variable class here is where the public-facing, user-relevant aspects are exposed, hiding internal details.

Testing

DAG:

from __future__ import annotations
from airflow.models.baseoperator import BaseOperator
from airflow.models.dag import dag
class CustomOperator(BaseOperator):
def execute(self, context):
import os
os.environ["AIRFLOW_VAR_HI_MESSAGE"] = "hello_world"
os.environ["AIRFLOW_VAR_JSON_VAR"] = "{\r\n \"key1\": \"value1\",\r\n \"key2\": \"value2\",\r\n \"enabled\": true,\r\n \"threshold\": 42\r\n}"
task_id = context["task_instance"].task_id
print(f"Hello World {task_id}!")
print(context)
print(context["var"]["value"].hi_message)
print(context["var"]["json"].json_var)
@dag()
def var_from_context():
CustomOperator(task_id="hello")
var_from_context()

This dag tests both the scenarios of a regular value context as well as json.

Case1: Variables found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:37:34.612546 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.613168 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.616934 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:37:34.617131 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:37:34.641905 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:37:34.642284 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:37:34.642393 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:37:34.037857+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9e-dadf-72eb-a47b-6c36e0fea3b7'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:37:34.037857+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 37, 34, 492777, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 37, 34, 37857, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:37:34.037857+00:00', 'ts_nodash': '20250106T123734', 'ts_nodash_with_tz': '20250106T123734.037857+0000'} [task] chan=stdout
2025-01-06 12:37:34.642268 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:37:34.648998 [info ] Variable(key='hi_message', value='hello_world', description=None) [task] chan=stdout
2025-01-06 12:37:34.649000 [debug ] Sending request [task] json={"key":"json_var","type":"GetVariable"}
2025-01-06 12:37:34.652506 [warning ] Pydantic serializer warnings:
Expected `str` but got `dict` with value `{'api_key': '12345', 'region': 'us-east-1'}` - serialized value may not be as expected [py.warnings] category=UserWarning filename=/usr/local/lib/python3.9/site-packages/pydantic/main.py lineno=426
2025-01-06 12:37:34.652586 [debug ] Sending request [task] json={"state":"success","end_date":"2025-01-06T12:37:34.652550Z","type":"TaskState"}
2025-01-06 12:37:34.652845 [info ] Variable(key='json_var', value={'api_key': '12345', 'region': 'us-east-1'}, description=None) [task] chan=stdout

Case 2: Variable not found

image

Logs:

2c27ff9a5949
▶ Log message source details
2025-01-06 12:34:57.712062 [info ] Filling up the DagBag from /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.712579 [debug ] Importing /files/dags/var_from_context.py [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717120 [debug ] Loaded DAG <DAG: var_from_context> [airflow.models.dagbag.DagBag]
2025-01-06 12:34:57.717434 [debug ] DAG file parsed [task] file=/files/dags/var_from_context.py
2025-01-06 12:34:57.743903 [warning ] CustomOperator.execute cannot be called outside TaskInstance! [airflow.task.operators.unusual_prefix_c6632fd34e048ff55a9057c21ee5a54c16b99828_var_from_context.CustomOperator]
2025-01-06 12:34:57.744163 [info ] Hello World hello! [task] chan=stdout
2025-01-06 12:34:57.744270 [info ] {'dag': <DAG: var_from_context>, 'inlets': [], 'map_index_template': None, 'outlets': [], 'run_id': 'manual__2025-01-06T12:34:56.856856+00:00', 'task': <Task(CustomOperator): hello>, 'task_instance': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'ti': RuntimeTaskInstance(id=UUID('01943b9c-74e4-7b93-ba7e-7401dde98944'), task_id='hello', dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', try_number=1, map_index=-1, task=<Task(CustomOperator): hello>), 'var': {'json': <VariableAccessor (dynamic access)>, 'value': <VariableAccessor (dynamic access)>}, 'conn': <ConnectionAccessor (dynamic access)>, 'dag_run': DagRun(dag_id='var_from_context', run_id='manual__2025-01-06T12:34:56.856856+00:00', logical_date=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_start=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), data_interval_end=datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), start_date=datetime.datetime(2025, 1, 6, 12, 34, 57, 589490, tzinfo=TzInfo(UTC)), end_date=None, run_type=<DagRunType.MANUAL: 'manual'>, conf={}), 'data_interval_end': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'data_interval_start': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'logical_date': datetime.datetime(2025, 1, 6, 12, 34, 56, 856856, tzinfo=TzInfo(UTC)), 'ds': '2025-01-06', 'ds_nodash': '20250106', 'task_instance_key_str': 'var_from_context__hello__20250106', 'ts': '2025-01-06T12:34:56.856856+00:00', 'ts_nodash': '20250106T123456', 'ts_nodash_with_tz': '20250106T123456.856856+0000'} [task] chan=stdout
2025-01-06 12:34:57.744177 [debug ] Sending request [task] json={"key":"hi_message","type":"GetVariable"}
2025-01-06 12:34:57.749149 [debug ] Sending request [task] json={"state":"failed","end_date":"2025-01-06T12:34:57.749108Z","type":"TaskState"}

TODO:

  • Writing variables from task SDK

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/context.py Outdated
Comment threadtask_sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment threadtask_sdk/tests/api/test_client.py
Comment threadtask_sdk/tests/api/test_client.py

@kaxilkaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few nits but lgtm

@amoghrajesh

Copy link
Copy Markdown
ContributorAuthor

Unrelated failure. Merging.

@amoghrajesh
amoghrajesh merged commit a6da8df into apache:mainJan 7, 2025
@amoghrajesh
amoghrajesh deleted the AIP72-variables-from-context branch January 7, 2025 06:05
HariGS-DB pushed a commit to HariGS-DB/airflow that referenced this pull request Jan 16, 2025
got686-yandex pushed a commit to got686-yandex/airflow that referenced this pull request Jan 30, 2025
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Allow retrieving Variable from Task Context

2 participants

@amoghrajesh@kaxil