Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down
, '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" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down
, '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('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down
, '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('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down
, '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" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down
, '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('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down
, '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('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down
, '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); } })(); })();
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions airflow/example_dags/example_datasets.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -51,18 +51,21 @@
"""
from __future__ import annotations

import random

import pendulum

from airflow.datasets import Dataset
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.timetables.datasets import DatasetOrTimeSchedule
from airflow.timetables.trigger import CronTriggerTimetable

# [START dataset_def]
dag1_dataset = Dataset("s3://dag1/output_1.txt", extra={"hi": "bye"})
# [END dataset_def]
dag2_dataset = Dataset("s3://dag2/output_1.txt", extra={"hi": "bye"})
dag2_dataset = Dataset("s3://dag2/output_1.txt")
dag3_dataset = Dataset("s3://dag3/output_3.txt", extra={"hi": "bye"})

with DAG(
Expand All@@ -83,7 +86,11 @@
schedule=None,
tags=["produces", "dataset-scheduled"],
) as dag2:
BashOperator(outlets=[dag2_dataset], task_id="producing_task_2", bash_command="sleep 5")

def some_python_callable():
return {"some_context": "dynamic data 123", "random_number": random.randint(1, 100)}

PythonOperator(python_callable=some_python_callable, outlets=[dag2_dataset], task_id="producing_task_2")

# [START dag_dep]
with DAG(
Expand Down
8 changes: 8 additions & 0 deletions airflow/jobs/scheduler_job_runner.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -1283,9 +1283,17 @@ def _create_dag_runs_dataset_triggered(
events=dataset_events,
)

run_conf = {}
for item in dataset_events:
event: DatasetEvent = item
extra: dict | None = event.extra
if extra:
run_conf.update(extra)
Comment on lines 1286 to 1291

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.

This feels a bit too heavy-handed, but I like the idea of passing in event extras as the downstream DAG run parameters. (Should we use conf or params for this?)

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

The code was just an idea - not finally thought through. If you say "heavy handed"... what do you mean? too brute force and the user does not know what comes out? Or do you mean we need a better merging mechanism? Or some hooking to be able to inject a custom merging strategy? Or just "code is ugly" :-D

Background: I would assume that in 90% of cases a single DAG triggers a dataset. THeremight be cases where multiple events come together to trigger. In such case we need to merge extras. I'd assume most times it is "conflict free" but you never know. Might be a feature to have it "last property wins" to collect events but otherwise if users feel there are too many conflicts, individual extras can also be produced "conflict free" with individual keys.

conf vs. params:

Yes, params and dag_run.conf somehow should be merged. I believe this is a leftover in the API from the past. CONF is the dict which is used to trigger a DAG. The conf is persisted as blob with the DagRun.
During runtime the conf is available in the context as dict, representing 1:1 the conf used to trigger. No validation. Just a dict.
params in contrast have default values, conf is setting values on top and the result is JSON validated.
Both üaramsand confare available in the context and can be used. I believre mid-term we should deprecate the usage of conf in the DAG and consolidate to the (more and better functional) params. But for today params only exist during runtime.
@hussein-awala did an attempt here but it dd not make it to finish line: #29174


dag_run = dag.create_dagrun(
run_id=run_id,
run_type=DagRunType.DATASET_TRIGGERED,
conf=run_conf,
execution_date=exec_date,
data_interval=data_interval,
state=DagRunState.QUEUED,
Expand Down
13 changes: 9 additions & 4 deletions airflow/models/taskinstance.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -2374,7 +2374,7 @@ def _run_raw_task(

try:
if not mark_success:
self._execute_task_with_callbacks(context, test_mode, session=session)
result = self._execute_task_with_callbacks(context, test_mode, session=session)
if not test_mode:
self.refresh_from_db(lock_for_update=True, session=session)
self.state = TaskInstanceState.SUCCESS
Expand DownExpand Up@@ -2462,7 +2462,7 @@ def _run_raw_task(
session.add(Log(self.state, self))
session.merge(self).task = self.task
if self.state == TaskInstanceState.SUCCESS:
self._register_dataset_changes(session=session)
self._register_dataset_changes(result, session=session)

session.commit()
if self.state == TaskInstanceState.SUCCESS:
Expand All@@ -2472,18 +2472,21 @@ def _run_raw_task(

return None

def _register_dataset_changes(self, *, session: Session) -> None:
def _register_dataset_changes(self, result: Any, *, session: Session) -> None:
for obj in self.task.outlets or []:
self.log.debug("outlet obj %s", obj)
# Lineage can have other types of objects besides datasets
if isinstance(obj, Dataset):
dataset_manager.register_dataset_change(
task_instance=self,
dataset=obj,
extra=obj.extra or result if isinstance(result, dict) else {self.task_id: result},
session=session,
)

def _execute_task_with_callbacks(self, context: Context, test_mode: bool = False, *, session: Session):
def _execute_task_with_callbacks(
self, context: Context, test_mode: bool = False, *, session: Session
) -> Any:
"""Prepare Task for Execution."""
from airflow.models.renderedtifields import RenderedTaskInstanceFields

Expand DownExpand Up@@ -2571,6 +2574,8 @@ def signal_handler(signum, frame):
Stats.incr("operator_successes", tags={**self.stats_tags, "task_type": self.task.task_type})
Stats.incr("ti_successes", tags=self.stats_tags)

return result

def _execute_task(self, context, task_orig):
"""
Execute Task (optionally with a Timeout) and push Xcom results.
Expand Down