Replies: 1 comment
|
This is expected with the current behavior: asset events are only registered when the producing task finishes in You can get the behavior you want today without changing Airflow by moving the from airflow.sdk import Asset, dag, task
from airflow.sdk.exceptions import AirflowSkipException
from airflow.providers.standard.operators.empty import EmptyOperator
source_b = Asset("s3://bucket/source_b")
@dag(schedule="@daily")
def dag_b():
@task
def ingest():
if not source_changed():
raise AirflowSkipException("source unchanged") # skip is fine now
load()
publish = EmptyOperator(
task_id="publish_source_b",
outlets=[source_b],
trigger_rule="none_failed", # runs if upstream succeeded OR skipped, not if it failed
)
ingest() >> publish
dag_b()
That covers both of your skip paths, because the task carrying the outlet is never the one being skipped. One gotcha: if the skip comes from a If you want consumers to be able to tell "unchanged" from "new data", you can also attach extra info to the event from a @task(outlets=[source_b], trigger_rule="none_failed")
def publish(*, outlet_events, ti):
changed = ti.xcom_pull(task_ids="ingest") is not None
outlet_events[source_b].extra = {"changed": changed}Changing core so skipped producers emit events would be a behavior change for everyone relying on the current semantics, so if you want that as a built-in option it's probably worth opening a feature request / dev-list discussion — but the pattern above works on current Airflow 3. |
Uh oh!
There was an error while loading. Please reload this page.
Hi comrades!
Problem
A task that declares outlets=[Dataset(...)] only emits a dataset event when it ends in success. If the task ends in skipped, no event is produced, and every consumer DAG scheduled on that dataset stays blocked waiting for an update that will never arrive.
This blocks us in a specific case: our ingestion DAGs skip work for sources that haven't changed. A skip there is a normal, expected outcome — the data is current, there is simply nothing new to load. But the downstream datamart DAG can't tell "nothing to ingest" apart from "ingestion hasn't happened yet", so it never triggers, and rarely-updated sources hold the whole datamart hostage.
Two skip paths, both affected
Worker-side skip — the producer task runs and raises AirflowSkipException.
Scheduler cascade skip — the producer task never runs at all, because an upstream branch wasn't taken (BranchPythonOperator / ShortCircuitOperator) or its trigger rule wasn't satisfied, and the scheduler marks it skipped directly.
Expected behaviour
When a producer DAG run reaches a terminal state, its declared datasets should be marked updated whether the producing task ended success or skipped. Failures should keep the current behaviour — no event.
Example
dag_a is scheduled on datasets B, C, D, produced by task_b in dag_b, task_c in dag_c, task_d in dag_d respectively.
Today: if dag_b runs and task_b skips, B gets no event and dag_a never triggers, even though C and D updated normally.
Wanted: once dag_b, dag_c and dag_d have all run, B, C and D are all marked updated and dag_a triggers.
All reactions