Skip to content

Fix OpenLineage jobDependencies run ID after upstream task is cleared - #73247

Closed
FrankYang0529 wants to merge 1 commit into
apache:mainfrom
FrankYang0529:airflow-openlineage-job-dependencies-emitting-try
Closed

FrankYang0529 wants to merge 1 commit into
apache:mainfrom
FrankYang0529:airflow-openlineage-job-dependencies-emitting-try

Conversation

@FrankYang0529

Copy link
Copy Markdown
Member

Why

  • Asset-triggered Dag runs emit a jobDependencies facet. For each asset event that triggered the Dag run, the facet lists the task and the OpenLineage run ID of the try that emitted the asset event. The run ID is built from the try number.
  • _extract_ol_info_from_asset_event takes the try number from AssetEvent.source_task_instance. That relationship looks up the task_instance row by Dag, run, task and map index. The row only holds the latest try, and the asset event does not record which try emitted it.
  • The facet is rebuilt for the START, COMPLETE and FAIL events of the Dag run. For example, try 1 of a task emits an asset event, and that asset event triggers a Dag run. Before that Dag run finishes, the task is cleared and scheduled again as try 2. From then on, every event of that Dag run points at try 2, even though try 1 emitted the asset event. If the Dag run is already running when the task is cleared, its START event points at try 1 and its COMPLETE event points at try 2.

How

  • DagRun.schedule_tis sets try_number and scheduled_dttm in the same update.
  • If the source task instance was scheduled at or before the event, its current try emitted the event. In that case the run ID is unchanged and no extra query runs.
  • If the source task instance was scheduled after the event, a new helper looks up the highest try in TaskInstanceHistory that was scheduled at or before the event.

Verification

  • Unit test: TZ=UTC uv run --frozen --project providers/openlineage pytest providers/openlineage/tests/unit
  • Integration test:
  1. Setup
uv sync --frozen --no-dev --package apache-airflow-providers-openlineage --package apache-airflow-providers-standard
npx -y pnpm@10.28.1 -C airflow-core/src/airflow/ui install --frozen-lockfile
npx -y pnpm@10.28.1 -C airflow-core/src/airflow/ui build
mkdir -p /tmp/ol-job-dependencies-demo/dags
cat > /tmp/ol-job-dependencies-demo/dags/orders.py <<'EOF'
import time
from pathlib import Path

import pendulum

from airflow.sdk import DAG, Asset, task

orders = Asset("demo://warehouse/orders")

with DAG(
    dag_id="orders_producer",
    schedule=None,
    start_date=pendulum.datetime(2025, 1, 1, tz="UTC"),
    is_paused_upon_creation=False,
):

    @task(outlets=[orders])
    def build_orders():
        return "orders table rebuilt"

    build_orders()

with DAG(
    dag_id="orders_report",
    schedule=[orders],
    start_date=pendulum.datetime(2025, 1, 1, tz="UTC"),
    is_paused_upon_creation=False,
):

    @task
    def make_report():
        while not Path("/tmp/ol-job-dependencies-demo/finish").exists():
            time.sleep(1)

    make_report()
EOF
cat > /tmp/ol-job-dependencies-demo/airflow.env <<'EOF'
AIRFLOW_HOME=/tmp/ol-job-dependencies-demo/airflow-home
AIRFLOW__CORE__DAGS_FOLDER=/tmp/ol-job-dependencies-demo/dags
AIRFLOW__CORE__LOAD_EXAMPLES=False
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_ALL_ADMINS=True
AIRFLOW__OPENLINEAGE__TRANSPORT='{"type": "file", "log_file_path": "/tmp/ol-job-dependencies-demo/airflow-home/ol_events.jsonl", "append": true}'
EOF
  1. Start Airflow
set -a; . /tmp/ol-job-dependencies-demo/airflow.env; PATH="$PWD/.venv/bin:$PATH"; exec airflow standalone
  1. Run the Dag

Trigger the Dag:

curl -s -X POST "http://localhost:8080/api/v2/dags/orders_producer/dagRuns" -H 'Content-Type: application/json' -d '{"logical_date": null}' | jq -c '{dag_run_id, state}'

Check the state of build_orders:

curl -s "http://localhost:8080/api/v2/dags/orders_producer/dagRuns/~/taskInstances" | jq -c '.task_instances[] | {task_id, try_number, state}'

Clear build_orders once it is success:

curl -s -X POST "http://localhost:8080/api/v2/dags/orders_producer/clearTaskInstances" -H 'Content-Type: application/json' -d '{"dry_run": false, "only_failed": false, "task_ids": ["build_orders"]}' | jq -c '.task_instances[] | {task_id, state}'

Let orders_report finish:

touch /tmp/ol-job-dependencies-demo/finish

Check both orders_report runs are success:

curl -s "http://localhost:8080/api/v2/dags/orders_report/dagRuns" | jq -c '.dag_runs[] | {dag_run_id, state}'
  1. Check result
jq -r 'select(.job.name == "orders_producer.build_orders" and .eventType == "COMPLETE") | "build_orders try \(.run.facets.airflow.taskInstance.try_number): run \(.run.runId)"' /tmp/ol-job-dependencies-demo/airflow-home/ol_events.jsonl
jq -r 'select(.job.name == "orders_report" and .eventType == "COMPLETE") | .run.facets.jobDependencies.upstream[] | "asset event \(.airflow.asset_events[0].asset_event_id): upstream run \(.run.runId)"' /tmp/ol-job-dependencies-demo/airflow-home/ol_events.jsonl | sort

On main, asset event 1 shows the run ID of try 2:

build_orders try 1: run 01a0aa0b-f137-7d33-b737-24e810559276
build_orders try 2: run 01a0aa0b-f137-7f1d-975b-87e1216ea322
asset event 1: upstream run 01a0aa0b-f137-7f1d-975b-87e1216ea322
asset event 2: upstream run 01a0aa0b-f137-7f1d-975b-87e1216ea322

On this branch, asset event 1 shows the run ID of try 1:

build_orders try 1: run 01a0a9d4-61ca-7dc3-b577-be0c3728e983
build_orders try 2: run 01a0a9d4-61ca-7716-92a2-1a0f5372c4d4
asset event 1: upstream run 01a0a9d4-61ca-7dc3-b577-be0c3728e983
asset event 2: upstream run 01a0a9d4-61ca-7716-92a2-1a0f5372c4d4

Was generative AI tooling used to co-author this PR?
  • Yes - Claude Code

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@FrankYang0529
FrankYang0529 marked this pull request as ready for review September 17, 2026 00:24
This was referenced Sep 25, 2026
@potiuk potiuk added the closed because of open PR limit Closed as a one-time step of introducing the open pull request limit label Sep 25, 2026
@potiuk

potiuk commented Sep 25, 2026

Copy link
Copy Markdown
Member

Hello @FrankYang0529 - thank you for your contributions to Apache Airflow!

The Airflow community has introduced a limit of 5 open pull requests at a time for contributors without write access to the repository. You currently have 30 open pull requests, so - as a one-time step of introducing the limit - we closed the ones where maintainers have not engaged yet:

These pull requests stay open because maintainers are already engaged in them - they count towards your limit:

This is not a judgement of you or of your changes. We never told contributors before that opening many pull requests at once was a problem, so there is nothing to feel bad about - and nothing is lost: your branches, commits and the review history stay where they are.

What we ask you to do is to make your first prioritization decision: choose which of the pull requests above matter most to you, and reopen them (up to 5 open at a time, including the ones still open) with the "Reopen pull request" button or gh pr reopen <PR_NUMBER> --repo apache/airflow. Reopen the ones you are ready to follow through - keep them rebased, respond to review comments and fix failing checks.

While your pull requests are waiting for review, the most valuable thing you can do is help in other ways - reviewing other contributors' pull requests, helping with issues, and taking part in the discussions on the devlist and Slack.

Why we introduced the limit, what it means for you and how to reopen or restore a pull request is explained in https://github.com/apache/airflow/blob/main/contributing-docs/32_open_pull_request_limit.rst.


Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

closed because of open PR limit Closed as a one-time step of introducing the open pull request limit

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants