This is an automated email from the ASF dual-hosted git repository.
kaxil pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new af3ba4c3a72 Add use-case pages to the common.ai provider docs (#73542)
af3ba4c3a72 is described below
commit af3ba4c3a72ceba130d4db0b0b99b893923f515f
Author: Kaxil Naik <[email protected]>
AuthorDate: Tue Sep 22 13:58:44 2026 +0100
Add use-case pages to the common.ai provider docs (#73542)
The provider docs were organized by operator, so a reader who wanted to
know what to build with common.ai found a parameter reference and a list of
example Dags grouped by mechanism. This adds a "What you can build" section
after Getting started with an overview and ten pages, each starting from a job
a data team already has, embedding the Dag that does it, and saying what
Airflow adds over a script. The end-to-end pipelines page is folded into the
new pages with a redirect, and the e [...]
Also fixes the triage example's include marker so the embedded snippet
shows its output class, a comment on the survey example that claimed the schema
check raises when it only reports, and stale task ids in the 10-K module
docstrings.
---
providers/common/ai/docs/end_to_end_pipelines.rst | 176 ---------------------
providers/common/ai/docs/examples.rst | 87 +++++-----
providers/common/ai/docs/index.rst | 27 +++-
providers/common/ai/docs/quickstart.rst | 1 +
providers/common/ai/docs/redirects.txt | 1 +
.../ai/docs/use_cases/ask_questions_over_pdfs.rst | 104 ++++++++++++
.../ai/docs/use_cases/classify_reviews_in_bulk.rst | 89 +++++++++++
.../ai/docs/use_cases/compare_10k_filings.rst | 118 ++++++++++++++
.../ai/docs/use_cases/explain_revenue_anomaly.rst | 78 +++++++++
.../docs/use_cases/gate_loads_on_schema_drift.rst | 76 +++++++++
providers/common/ai/docs/use_cases/index.rst | 108 +++++++++++++
.../docs/use_cases/monthly_report_from_a_csv.rst | 103 ++++++++++++
.../docs/use_cases/research_agent_with_review.rst | 88 +++++++++++
.../ai/docs/use_cases/route_pipeline_failures.rst | 88 +++++++++++
.../ai/docs/use_cases/triage_support_tickets.rst | 72 +++++++++
.../ai/docs/use_cases/weekly_status_report.rst | 120 ++++++++++++++
.../ai/example_dags/example_langchain_10k.py | 2 +-
.../ai/example_dags/example_llamaindex_10k.py | 2 +-
.../example_dags/example_llm_analysis_pipeline.py | 4 +-
.../ai/example_dags/example_llm_survey_analysis.py | 4 +-
20 files changed, 1120 insertions(+), 228 deletions(-)
diff --git a/providers/common/ai/docs/end_to_end_pipelines.rst
b/providers/common/ai/docs/end_to_end_pipelines.rst
deleted file mode 100644
index 3c38d07891a..00000000000
--- a/providers/common/ai/docs/end_to_end_pipelines.rst
+++ /dev/null
@@ -1,176 +0,0 @@
- .. Licensed to the Apache Software Foundation (ASF) under one
- or more contributor license agreements. See the NOTICE file
- distributed with this work for additional information
- regarding copyright ownership. The ASF licenses this file
- to you under the Apache License, Version 2.0 (the
- "License"); you may not use this file except in compliance
- with the License. You may obtain a copy of the License at
-
- .. http://www.apache.org/licenses/LICENSE-2.0
-
- .. Unless required by applicable law or agreed to in writing,
- software distributed under the License is distributed on an
- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- KIND, either express or implied. See the License for the
- specific language governing permissions and limitations
- under the License.
-
-End-to-end pipelines
-====================
-
-The Dags in this guide combine several patterns from the operator and hook
guides into one
-production-shaped pipeline. Each section explains the architecture -- how the
Dags are split,
-why, and how the pieces are wired together -- rather than the mechanics of a
single operator,
-which the linked guides already cover. Read the full source for the runnable
Dag.
-
-LlamaIndex RAG shapes
------------------------
-
-`example_llamaindex_rag.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_rag.py>`__
-walks the same load -> embed -> retrieve -> answer pattern through three
shapes, from simplest to
-production-shaped:
-
-- ``example_llamaindex_rag_pipeline`` (``schedule=None``) -- everything in one
Dag:
- ``DocumentLoaderOperator`` parses local files,
``LlamaIndexEmbeddingOperator`` chunks and
- persists the index, ``LlamaIndexRetrievalOperator`` retrieves for a fixed
question, and an
- ``LLMOperator`` synthesizes the answer.
-- ``example_llamaindex_index_pdf`` / ``example_llamaindex_query`` -- the
load/embed step moved
- into its own weekly Dag (``schedule="@weekly"``) so the index stays fresh as
PDFs arrive, while
- a second Dag (``schedule=None``, triggered with a ``question`` param)
retrieves and answers on
- demand against the persisted index -- the same index/query split the SEC
10-K pipelines below
- use, without their multi-company retrieval fan-out.
-- ``example_llamaindex_multi_source`` -- two ``DocumentLoaderOperator`` calls
tag documents from
- different sources via ``metadata_fields`` before merging and embedding them
into one index, for
- filtered retrieval downstream.
-
-See :doc:`operators/document_loader`, :doc:`operators/llamaindex_embedding`,
and
-:doc:`operators/llamaindex_retrieval` for how the individual operators work.
-
-SEC 10-K financial analysis
-----------------------------
-
-`example_llamaindex_10k.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py>`__
and
-`example_langchain_10k.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_10k.py>`__
-compare companies' SEC 10-K filings, fetched live from EDGAR. Each file splits
the work into two
-Dags:
-
-- An **indexing Dag** (``schedule="@weekly"``) that fetches each company's
latest 10-K and
- persists an index per ticker to disk.
-- An **analysis Dag** (``schedule=None``, triggered manually or on demand)
that decomposes a
- comparison question, retrieves against the indexes, and produces a reviewed
report.
-
-The two Dags are not linked by an Asset or XCom -- they agree only on a shared
on-disk index
-path (``INDEX_BASE_DIR/{ticker}``). Run the indexing Dag at least once before
triggering
-analysis for a given company.
-
-Both files build the same analysis shape; the only structural difference is
that LlamaIndex has
-dedicated ``LlamaIndexEmbeddingOperator``/``LlamaIndexRetrievalOperator``
classes, while
-LangChain does not yet, so the LangChain file's indexing and retrieval steps
are plain ``@task``
-functions calling
:class:`~airflow.providers.common.ai.hooks.langchain.LangChainHook` and FAISS
-directly.
-
-Analysis Dag flow:
-
-1. A human confirms or edits the comparison question and the tickers to compare
- (``HITLEntryOperator``).
-2. ``@task.llm`` decomposes the question into one sub-question per company --
the LLM decides how
- many sub-questions are needed, not the Dag author:
-
- .. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
- :language: python
- :start-after: [START 10k_decompose]
- :end-before: [END 10k_decompose]
-
-3. Dynamic Task Mapping fans retrieval out, one mapped task per sub-question,
each querying that
- company's index:
-
- .. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
- :language: python
- :start-after: [START 10k_dtm_retrieval]
- :end-before: [END 10k_dtm_retrieval]
-
-4. The mapped results are zipped back to their sub-questions (mapped task
outputs preserve input
- order) and synthesized into a structured ``AnalysisReport`` by an
``LLMOperator`` bounded by
- ``UsageLimits``.
-5. ``ApprovalOperator`` gates the report before hand-off.
-
-See :doc:`operators/llamaindex_embedding`,
:doc:`operators/llamaindex_retrieval`, and
-:doc:`hooks/langchain` for how the individual operators/hooks behave.
-
-Natural-language survey analysis
-----------------------------------
-
-`example_llm_survey_analysis.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py>`__
-answers natural-language questions over a CSV with the same core three tasks --
-``LLMSQLQueryOperator`` generates SQL, ``AnalyticsOperator`` runs it, a
``@task`` extracts the
-rows (see :doc:`operators/llm_sql` for how those two operators work together).
The file defines
-two Dags around that core, shaped for different operating modes:
-
-- ``example_llm_survey_interactive`` (``schedule=None``) -- a human confirms
or edits the
- question before SQL generation (``HITLEntryOperator``) and approves the
extracted result
- afterward (``ApprovalOperator``). It assumes the CSV is already in place.
-- ``example_llm_survey_scheduled`` (``schedule="@monthly"``) -- downloads the
CSV itself
- (``HttpOperator``), validates its schema against a reference CSV
- (``LLMSchemaCompareOperator``) before generating SQL, and ends by emailing
or logging the
- result. No human gate -- suited to recurring reporting.
-
-Agentic multi-dimensional survey synthesis
----------------------------------------------
-
-`example_llm_survey_agentic.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_agentic.py>`__
-answers a research question that a single SQL query cannot -- one that spans
several
-dimensions (executor, deployment, cloud, Airflow version) -- by fanning both
SQL generation and
-execution out with Dynamic Task Mapping, one mapped pair per dimension, then
synthesizing:
-
-1. ``decompose_question`` returns one sub-question per dimension.
-2. ``LLMSQLQueryOperator.partial(...).expand(prompt=sub_questions)`` generates
SQL for each
- dimension as an independent mapped task instance.
-3. ``AnalyticsOperator.partial(...).expand(...)`` runs each query. If one
dimension's query
- fails, only that mapped instance retries -- the other three keep their
results.
-4. ``collect_results`` zips the dimension labels back onto the
(order-preserved) results.
-5. An ``LLMOperator`` synthesizes the four labeled result sets into one
narrative, then
- ``ApprovalOperator`` gates it.
-
-The Dag's own docstring explains the reasoning for fanning out rather than
looping inside a
-single LLM call: each sub-query becomes a named, logged task instance instead
of a hidden tool
-call, so a failure is isolated and retryable, and every intermediate result is
visible in XCom
-instead of being an opaque step inside an LLM reasoning loop.
-
-AIP progress tracking: pipeline vs. agent
---------------------------------------------
-
-`example_aip_progress_tracker.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py>`__
-tracks the same thing -- progress on a set of Airflow Improvement Proposals,
checked against
-Confluence specs and GitHub activity -- with two different architectures, to
compare the tradeoff directly:
-
-- ``example_aip_progress_tracker`` -- a **deterministic pipeline**. Evidence
is gathered by
- fixed tasks; dynamically-mapped ``LLMOperator`` calls analyze each AIP with
structured output;
- the per-AIP analyses are synthesized into one report; a second
``LLMOperator`` validates that
- synthesis against the raw evidence and flags unsupported claims; a
plain-Python task applies
- only the flagged corrections mechanically (string/regex replacement, no LLM
involved); a human
- reviews the corrected report.
-- ``example_aip_progress_tracker_skills`` -- an **autonomous agent**. A single
``AgentOperator``
- loaded with an Agent Skill (``AgentSkillsToolset``) and a handful of custom
tool functions
- decides its own evidence-gathering strategy and tool-call order, then the
same human review
- gate.
-
-The pipeline's three-layer defense against hallucination is the pattern worth
taking away:
-structured LLM output, then a second LLM call whose only job is to judge the
first against the
-original evidence, then a deterministic (non-LLM) step that applies exactly
what was flagged:
-
-.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
- :language: python
- :start-after: [START aip_tracker_validation]
- :end-before: [END aip_tracker_validation]
-
-The agent variant collapses that whole graph into one operator call, trading
auditability of
-each step for simplicity and letting the model decide how to use its tools:
-
-.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
- :language: python
- :start-after: [START aip_tracker_skills_operator]
- :end-before: [END aip_tracker_skills_operator]
-
-Reach for the pipeline when accuracy matters and every step needs to be
auditable; reach for the
-agent when the problem is open-ended enough that a fixed task graph would just
be re-encoding
-what the model can already work out from the skill instructions.
diff --git a/providers/common/ai/docs/examples.rst
b/providers/common/ai/docs/examples.rst
index 0ac806e16a2..89c86a2d854 100644
--- a/providers/common/ai/docs/examples.rst
+++ b/providers/common/ai/docs/examples.rst
@@ -17,16 +17,14 @@
.. _howto/examples:
-Examples
-========
+All example Dags
+================
-Every operator, decorator, and integration in this provider has a runnable Dag
under
-`example_dags
<https://github.com/apache/airflow/tree/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags>`__.
-This page groups them by scenario. Guides embed the same Dags inline where a
step-by-step
-walkthrough exists; the rest are linked directly to source below.
+Every operator, decorator and integration has a runnable Dag under
+`example_dags
<https://github.com/apache/airflow/tree/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags>`__,
+listed here by operator. Guides embed the Dags they walk through; the rest
link to source.
-New to the provider? Start with :doc:`operators/index` to pick an operator,
then come back here
-for a worked example of the pattern you need.
+For what to build rather than how, start at :doc:`use_cases/index`.
Single-prompt tasks
--------------------
@@ -38,17 +36,21 @@ Single-prompt tasks
* - Guide
- What it shows
* - :doc:`operators/llm`
- - ``@task.llm`` for text, structured output, classification, and dynamic
task mapping over
- LLM results
+ - Summarize and extract entities, grade incident severity, and
+ :doc:`triage a queue of support tickets
<use_cases/triage_support_tickets>` one mapped
+ task at a time
(`example_llm.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm.py>`__,
`example_llm_classification.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_classification.py>`__,
`example_llm_analysis_pipeline.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_analysis_pipeline.py>`__).
* - :doc:`operators/llm_branch`
- - ``@task.llm_branch`` picking which downstream task runs.
+ - Let the model pick which downstream task runs, and
+ :doc:`route a failed task to rerun, page or ignore
<use_cases/route_pipeline_failures>`
+ with a confidence bar.
* - :doc:`operators/llm_file_analysis`
- ``@task.llm_file_analysis`` reasoning over files, images, and PDFs.
* - :doc:`operators/llm_schema_compare`
- - ``@task.llm_schema_compare`` comparing two schemas with an LLM.
+ - Compare two schemas and
+ :doc:`block a load when they drifted
<use_cases/gate_loads_on_schema_drift>`.
* - :doc:`operators/llm_sql`
- ``@task.llm_sql`` generating SQL from a natural-language question.
@@ -62,8 +64,9 @@ Batch processing
* - Guide
- What it shows
* - :doc:`operators/llm_batch`
- - ``@task.llm_batch`` submitting many prompts as one OpenAI/Anthropic
batch job, with
- structured output and manifest-based results
+ - :doc:`Classify thousands of reviews at half the price
<use_cases/classify_reviews_in_bulk>`
+ through one OpenAI or Anthropic batch job, with structured output and
results landed on
+ object storage
(`example_llm_batch.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_batch.py>`__).
Agents & tools
@@ -90,9 +93,10 @@ Agents & tools
- Connecting an agent to an MCP server through an Airflow connection.
* - :doc:`hitl_review`
- Adding a human-in-the-loop review gate to agent output.
- * - `example_langchain_tool_agent.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py>`__
+ * - :doc:`use_cases/research_agent_with_review`
- A LangChain ReAct agent that decides its own tool calls, composed with
``LLMOperator`` for
- report formatting and AIP-90 HITL review.
+ report formatting and AIP-90 HITL review
+ (`example_langchain_tool_agent.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py>`__).
Retrieval & document processing
--------------------------------
@@ -112,35 +116,40 @@ Retrieval & document processing
* - :doc:`hooks/llamaindex`
- ``LlamaIndexHook`` plus the embedding and retrieval operators.
-End-to-end scenarios
----------------------
+By use case
+-----------
-Production-shaped Dags that combine several of the patterns above into one
pipeline. See
-:doc:`end_to_end_pipelines` for the architecture behind each one.
+Dags written around a job. Each has a page under :doc:`use_cases/index` with
the Dag
+embedded and steps to run it.
.. list-table::
:header-rows: 1
:widths: 30 70
- * - Example
- - What it shows
- * - `example_llamaindex_rag.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_rag.py>`__
- - Three RAG shapes with LlamaIndex: a single load-embed-retrieve-answer
Dag, a production-shaped
- split into scheduled indexing plus on-demand query Dags, and
multi-source RAG.
- * - `example_llamaindex_10k.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py>`__,
- `example_langchain_10k.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_10k.py>`__
- - SEC 10-K financial analysis against live EDGAR filings: one Dag indexes
filings on a
- schedule, the other decomposes a comparison question at runtime and
fans out retrieval with
- Dynamic Task Mapping. One variant per RAG library.
- * - `example_llm_survey_analysis.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py>`__
- - Natural-language querying of a CSV via ``LLMSQLQueryOperator`` and
``AnalyticsOperator``,
- interactive and scheduled variants.
- * - `example_llm_survey_agentic.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_agentic.py>`__
- - The same survey data, analyzed by fanning a multi-dimensional research
question into an
- agentic multi-query synthesis pipeline.
- * - `example_aip_progress_tracker.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py>`__
- - The same use case -- tracking Airflow Improvement Proposal progress --
solved two ways, a
- deterministic pipeline and an autonomous agent, to compare the tradeoff.
+ * - Use case
+ - Source
+ * - :doc:`use_cases/ask_questions_over_pdfs`
+ - `example_llamaindex_rag.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_rag.py>`__:
a weekly indexing Dag
+ plus an on-demand query Dag, with single-Dag and multi-source variants.
+ * - :doc:`use_cases/compare_10k_filings`
+ - `example_llamaindex_10k.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py>`__
and
+ `example_langchain_10k.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_10k.py>`__:
live SEC EDGAR filings,
+ per-company retrieval fan-out, human review at both ends. One variant
per RAG library.
+ * - :doc:`use_cases/monthly_report_from_a_csv`
+ - `example_llm_survey_analysis.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py>`__:
download,
+ schema check, generated SQL, email; plus an interactive variant with
HITL.
+ `example_llm_survey_agentic.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_agentic.py>`__
fans a
+ multi-dimensional question out one SQL query per dimension.
+ * - :doc:`use_cases/explain_revenue_anomaly`
+ - `example_sandbox_toolset.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_sandbox_toolset.py>`__:
an agent with a
+ read-only warehouse toolset and a sandbox for the arithmetic.
+ * - :doc:`use_cases/weekly_status_report`
+ - `example_aip_progress_tracker.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py>`__:
the same
+ report built as a deterministic pipeline with a hallucination check and
as one
+ autonomous agent.
+ * - :doc:`use_cases/research_agent_with_review`
+ - `example_langchain_tool_agent.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py>`__:
a LangChain
+ ReAct agent between a question-review gate and a report-approval gate.
Reliability
-----------
diff --git a/providers/common/ai/docs/index.rst
b/providers/common/ai/docs/index.rst
index 62be874817f..da65b8e9fa3 100644
--- a/providers/common/ai/docs/index.rst
+++ b/providers/common/ai/docs/index.rst
@@ -71,6 +71,7 @@ Getting started
* :doc:`installation` — which extra to install for your model vendor.
* :doc:`quickstart` — a connection and a first ``@task.llm`` in three steps.
* :doc:`concepts` — connections, operators, toolsets, hooks and XCom in one
page.
+* :doc:`use_cases/index` — jobs a data team already has, each with the Dag
that does it.
.. toctree::
:hidden:
@@ -91,6 +92,23 @@ Getting started
Core concepts <concepts>
Structured output <structured_output>
+.. toctree::
+ :hidden:
+ :maxdepth: 1
+ :caption: What you can build
+
+ Overview <use_cases/index>
+ Triage support tickets <use_cases/triage_support_tickets>
+ Route pipeline failures <use_cases/route_pipeline_failures>
+ Block a load on schema drift <use_cases/gate_loads_on_schema_drift>
+ Explain a revenue anomaly <use_cases/explain_revenue_anomaly>
+ Monthly report from a CSV <use_cases/monthly_report_from_a_csv>
+ Compare 10-K filings <use_cases/compare_10k_filings>
+ Ask questions over PDFs <use_cases/ask_questions_over_pdfs>
+ Weekly status report <use_cases/weekly_status_report>
+ Classify reviews in bulk <use_cases/classify_reviews_in_bulk>
+ Research agent with review <use_cases/research_agent_with_review>
+
.. toctree::
:hidden:
:maxdepth: 1
@@ -169,19 +187,12 @@ Getting started
Securing agent tools <agent_security>
Troubleshooting <troubleshooting>
-.. toctree::
- :hidden:
- :maxdepth: 1
- :caption: Examples
-
- Examples by scenario <examples>
- End-to-end pipelines <end_to_end_pipelines>
-
.. toctree::
:hidden:
:maxdepth: 1
:caption: References
+ All example Dags <examples>
Configuration <configurations-ref>
Python API <_api/airflow/providers/common/ai/index>
diff --git a/providers/common/ai/docs/quickstart.rst
b/providers/common/ai/docs/quickstart.rst
index 3c23bd323d5..38f0685bad9 100644
--- a/providers/common/ai/docs/quickstart.rst
+++ b/providers/common/ai/docs/quickstart.rst
@@ -83,6 +83,7 @@ requirements.
Where to go next
-----------------
+- :doc:`use_cases/index` — start here for what to build.
- :doc:`operators/index` — the full set of operators and ``@task`` decorators
(file analysis, SQL, branching, schema comparison).
- :doc:`toolsets/index` — give an agent tools built from Airflow hooks, SQL
diff --git a/providers/common/ai/docs/redirects.txt
b/providers/common/ai/docs/redirects.txt
index b7f775cf404..68309579391 100644
--- a/providers/common/ai/docs/redirects.txt
+++ b/providers/common/ai/docs/redirects.txt
@@ -20,3 +20,4 @@ toolsets.rst toolsets/index.rst
choosing_a_toolset.rst toolsets/index.rst
sandbox.rst sandbox/index.rst
hooks/index.rst concepts.rst
+end_to_end_pipelines.rst use_cases/index.rst
diff --git a/providers/common/ai/docs/use_cases/ask_questions_over_pdfs.rst
b/providers/common/ai/docs/use_cases/ask_questions_over_pdfs.rst
new file mode 100644
index 00000000000..a97f5da84cf
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/ask_questions_over_pdfs.rst
@@ -0,0 +1,104 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Ask questions over a growing PDF corpus
+=======================================
+
+A folder of quarterly reports keeps growing and people keep asking questions
answered
+somewhere inside it. One Dag keeps a vector index fresh weekly; the other
answers a question
+on demand from retrieved excerpts, citing them by number. Indexing runs on a
schedule so no
+query pays for embedding, and ``question`` is a param, so anyone can trigger
it from the UI,
+CLI or REST API.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/document_loader` -- ``DocumentLoaderOperator`` parses a
glob of PDFs
+ into documents.
+* :doc:`../operators/llamaindex_embedding` -- ``LlamaIndexEmbeddingOperator``
chunks,
+ embeds and persists the index to disk.
+* :doc:`../operators/llamaindex_retrieval` -- ``LlamaIndexRetrievalOperator``
pulls the
+ top matching chunks for a question.
+* :doc:`../operators/llm` -- ``LLMOperator`` answers from the excerpts only,
with the
+ question and the context templated in from params and XCom.
+
+Run it
+------
+
+1. Install the provider with the LlamaIndex and PDF extras:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai,llamaindex,pdf]"
+
+2. Create a ``llamaindex`` connection named ``llamaindex_default`` for
embedding and
+ retrieval (see :doc:`../connections/llamaindex`); ``pydanticai_default``
writes the answer.
+
+3. Put some PDFs under ``/opt/airflow/data/reports/``, run the indexing Dag
once, then
+ ask a question:
+
+ .. code-block:: bash
+
+ airflow dags test example_llamaindex_index_pdf
+ airflow dags test example_llamaindex_query \
+ --conf '{"question": "What drove the change in operating margin?"}'
+
+The ``synthesize`` XCom holds the answer with ``[n]`` references into the
excerpts that
+``format_context`` numbered.
+
+The indexing Dag
+----------------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_rag.py
+ :language: python
+ :start-after: [START howto_llamaindex_index_dag]
+ :end-before: [END howto_llamaindex_index_dag]
+
+The query Dag
+-------------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_rag.py
+ :language: python
+ :start-after: [START howto_llamaindex_query_dag]
+ :end-before: [END howto_llamaindex_query_dag]
+
+Two more shapes
+---------------
+
+``example_llamaindex_rag_pipeline`` runs the same four operators in one Dag
with a fixed
+question. Run it once to see the chain end to end.
+
+``example_llamaindex_multi_source`` loads from two places, tags each document
with where
+it came from through ``metadata_fields``, and embeds them into one index so
retrieval can
+filter by source later:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_rag.py
+ :language: python
+ :start-after: [START howto_llamaindex_multi_source]
+ :end-before: [END howto_llamaindex_multi_source]
+
+Adapting it
+-----------
+
+* Change ``source_path`` to your folder or an object storage URI; the loader
also reads
+ DOCX, CSV and JSON.
+* Change the indexing schedule to match how often documents land, or trigger
it from an
+ Asset when the upstream Dag that drops the PDFs emits one.
+* Raise ``top_k`` for broad questions, lower it for precise ones. Keep the
excerpts-only
+ instruction in the prompt; it stops the model filling gaps from memory.
+* :doc:`compare_10k_filings` grows this into a fan-out over several indexes
with a human at
+ both ends.
diff --git a/providers/common/ai/docs/use_cases/classify_reviews_in_bulk.rst
b/providers/common/ai/docs/use_cases/classify_reviews_in_bulk.rst
new file mode 100644
index 00000000000..b51007bae38
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/classify_reviews_in_bulk.rst
@@ -0,0 +1,89 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Classify reviews in bulk at half the price
+==========================================
+
+A day of product reviews needs sentiment labels, nobody needs them in the next
minute, and
+a hundred thousand full-price calls is hard to justify. This Dag sends them
all to the
+vendor's batch API in one job at about half the per-request price, with up to
24 hours'
+turnaround. Airflow defers while the batch runs so no worker is held,
re-attaches to the
+same batch on retry instead of paying twice, and lands results as JSONL on
object storage
+with only a manifest in XCom.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/llm_batch` -- ``LLMBatchOperator`` with
``deferrable=True`` and a
+ ``result_path`` keyed by ``run_id``.
+* :doc:`../structured_output` -- every request in the batch asks for the same
+ ``Sentiment`` schema, and rows that fail validation are kept with their raw
text.
+* ``ObjectStoragePath`` -- the downstream task reads the landed rows back from
storage
+ rather than through XCom.
+
+Run it
+------
+
+1. Install the provider with the OpenAI extra. Change ``model_id`` to run the
same batch
+ on Anthropic:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai]"
+
+2. Change ``RESULT_ROOT`` in the file to a bucket you can write to, with an
object storage
+ connection for it.
+
+3. Trigger the Dag. It defers until the vendor finishes the batch, which can
take hours:
+
+ .. code-block:: bash
+
+ airflow dags test example_llm_batch_operator
+
+The ``classify_reviews`` XCom is a manifest with the result URI and counts by
status.
+``summarize_manifest`` reads the JSONL and logs a count per label.
+
+The Dag
+-------
+
+The output schema and the reader task, then the Dag body:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_batch.py
+ :language: python
+ :start-after: [START howto_operator_llm_batch_structured_output_class]
+ :end-before: [END howto_operator_llm_batch_structured_output_class]
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_batch.py
+ :language: python
+ :start-after: [START howto_operator_llm_batch_read_results]
+ :end-before: [END howto_operator_llm_batch_read_results]
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_batch.py
+ :language: python
+ :start-after: [START howto_operator_llm_batch_basic]
+ :end-before: [END howto_operator_llm_batch_basic]
+
+Adapting it
+-----------
+
+* Replace ``REVIEWS`` with a task that pulls the day's reviews from your store
and pass
+ its output as ``requests``. The operator accepts a list of prompts or of
+ per-request overrides; see :doc:`../operators/llm_batch`.
+* Keep ``run_id`` in ``result_path``. A timestamp would break re-attachment on
retry and
+ let a flaky worker submit the batch twice; :doc:`../operators/llm_batch`
explains why.
+* Load the JSONL into your warehouse from a downstream task instead of
counting labels,
+ and emit an Asset so reporting Dags can run when the day's labels land.
diff --git a/providers/common/ai/docs/use_cases/compare_10k_filings.rst
b/providers/common/ai/docs/use_cases/compare_10k_filings.rst
new file mode 100644
index 00000000000..4bb862752d0
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/compare_10k_filings.rst
@@ -0,0 +1,118 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Compare companies' 10-K filings
+===============================
+
+An analyst wants to know which of several companies carries the most
concentrated risk and
+which is growing fastest, grounded in the filings rather than the model's
memory. A weekly
+Dag fetches each latest 10-K from SEC EDGAR and indexes it. An on-demand Dag
has the model
+split the question per company, retrieves against each index in its own task
so a missing
+index retries alone, and writes a structured report. A human edits the
question on the way
+in and approves the report on the way out.
+
+Two variants build the same graph: LlamaIndex operators, or LangChain with
FAISS.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/llamaindex_embedding` and
+ :doc:`../operators/llamaindex_retrieval` -- index per ticker on a schedule,
retrieve
+ per sub-question on demand.
+* :doc:`../operators/llm` -- ``@task.llm`` with structured output for
decomposition, and
+ ``LLMOperator`` bounded by ``UsageLimits`` for synthesis.
+* :doc:`apache-airflow:authoring-and-scheduling/dynamic-task-mapping` -- the
model decides
+ how many sub-questions there are; the Dag maps over them at runtime.
+* :doc:`apache-airflow-providers-standard:operators/hitl` --
``HITLEntryOperator`` on the
+ way in, ``ApprovalOperator`` on the way out.
+* :doc:`../hooks/langchain` -- the LangChain variant does indexing and
retrieval in plain
+ ``@task`` functions over ``LangChainHook`` and FAISS.
+
+Run it
+------
+
+1. Install the provider with the LlamaIndex extra (or ``langchain`` for the
other
+ variant):
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai,llamaindex]"
+
+2. Create a ``llamaindex`` connection named ``llamaindex_default`` for
embedding and
+ retrieval (see :doc:`../connections/llamaindex`); ``pydanticai_default``
does
+ decomposition and synthesis.
+
+3. Set ``EDGAR_USER_AGENT`` in the file to your name and email. SEC requires a
contact
+ address on every EDGAR request. No API key is needed.
+
+4. Run the indexing Dag once, then trigger the analysis:
+
+ .. code-block:: bash
+
+ airflow dags test example_llamaindex_10k_index
+ airflow dags test example_llamaindex_10k_analysis
+
+The run pauses at ``analyst_input`` to confirm the question and tickers, and at
+``review_report``. Answer both from Required Actions in the UI (see the note on
+:doc:`index`). The ``synthesize_report`` XCom holds the ``AnalysisReport``.
+
+The Dags share only the index path ``INDEX_BASE_DIR/<lowercased ticker>``, so
index a
+ticker before analyzing it.
+
+The indexing Dag
+----------------
+
+One mapped ``LlamaIndexEmbeddingOperator`` per ticker, weekly:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
+ :language: python
+ :start-after: [START example_llamaindex_10k_index]
+ :end-before: [END example_llamaindex_10k_index]
+
+The analysis Dag
+----------------
+
+The output types come first. ``DecomposedQuestion`` is what the model returns
from the
+decomposition step, and ``AnalysisReport`` is what the reviewer approves:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
+ :language: python
+ :start-after: [START 10k_structured_output]
+ :end-before: [END 10k_structured_output]
+
+The Dag itself:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
+ :language: python
+ :start-after: [START example_llamaindex_10k_analysis]
+ :end-before: [END example_llamaindex_10k_analysis]
+
+Decomposition is the step to notice: the model returns a list and the Dag maps
retrieval
+over it, so the model decides at runtime how many tasks run, each with its own
log and
+retries.
+
+Adapting it
+-----------
+
+* Change ``DEFAULT_TICKERS`` or pass ``tickers`` as a Dag param. Any US-listed
company
+ works; EDGAR resolves the ticker.
+* Replace ``fetch_filings`` with your own document source; nothing downstream
cares where
+ the text came from.
+* Tighten ``UsageLimits`` on ``synthesize_report`` to cap spend per run.
+* The LangChain build is ``example_langchain_10k.py`` (Dag ids
``example_langchain_10k_index``
+ and ``example_langchain_10k_analysis``); it adds a ``langchain_default``
connection for
+ embeddings.
diff --git a/providers/common/ai/docs/use_cases/explain_revenue_anomaly.rst
b/providers/common/ai/docs/use_cases/explain_revenue_anomaly.rst
new file mode 100644
index 00000000000..29b96acb7cd
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/explain_revenue_anomaly.rst
@@ -0,0 +1,78 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Explain a revenue anomaly
+=========================
+
+Daily revenue moved more than ten percent against its trailing average and
someone has
+to find out which region and channel explain it before the morning stand-up.
This Dag
+gives an agent two tools: read-only queries against two warehouse tables,
capped at 500
+rows, and a pandas sandbox for the pivots. It returns a typed finding with the
suspected
+cause and a confidence. Airflow injects the run date, keeps the warehouse
credential where
+the model never sees it, and destroys the sandbox when the run ends.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/agent` -- ``AgentOperator`` runs a multi-turn agent as
one task,
+ with ``output_type`` for a typed result.
+* :doc:`../toolsets/sql` -- ``SQLToolset`` exposes named query operations; the
model sees
+ rows, never the connection.
+* :doc:`../sandbox/index` -- ``SandboxToolset`` runs everything the model
writes off the
+ worker.
+* Templating -- ``{{ ds }}`` in the prompt ties the question to the Dag run's
date.
+
+Run it
+------
+
+1. Install the provider with the SQL and Modal extras and authenticate with
Modal:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai,sql,modal]"
+ modal setup
+
+2. Create a database connection named ``warehouse`` that has
+ ``analytics.daily_revenue`` and ``analytics.orders`` tables.
+
+3. Trigger the Dag for a date:
+
+ .. code-block:: bash
+
+ airflow dags test example_sandbox_agent_investigation 2026-03-01
+
+The ``investigate`` task log shows every SQL call and every script the agent
ran in the
+sandbox, and its XCom holds the ``Findings`` record.
+
+The Dag
+-------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_sandbox_toolset.py
+ :language: python
+ :start-after: [START howto_sandbox_agent_investigation]
+ :end-before: [END howto_sandbox_agent_investigation]
+
+Adapting it
+-----------
+
+* Trigger it from a data-quality check that fires when the move exceeds your
threshold, so
+ the agent runs only on days that need explaining.
+* Swap ``ModalSandboxBackend`` for ``SbxSandboxBackend`` to run locally; the
same file has
+ that variant (see :doc:`../sandbox/backends`).
+* Add an :doc:`../approval_gates` step before the finding reaches the finance
channel.
+* Widen ``allowed_tables`` only as far as the question needs. The limit is
what makes the
+ agent safe to run unattended.
diff --git a/providers/common/ai/docs/use_cases/gate_loads_on_schema_drift.rst
b/providers/common/ai/docs/use_cases/gate_loads_on_schema_drift.rst
new file mode 100644
index 00000000000..b44fe95cf73
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/gate_loads_on_schema_drift.rst
@@ -0,0 +1,76 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Block a load when the schema drifts
+===================================
+
+A nightly load copies ``customers`` from Postgres into Snowflake. One morning
a column was
+renamed upstream and the load either failed halfway or, worse, succeeded with
nulls. This
+Dag compares the two schemas before the load and asks the model which
differences would
+break it. Airflow branches on the answer: compatible schemas run the load,
anything else
+notifies the team. The model only reports; no migration runs.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/llm_schema_compare` -- ``@task.llm_schema_compare`` reads
both
+ schemas through Airflow connections and returns a structured comparison with
a
+ ``compatible`` flag.
+* :ref:`Branching <apache-airflow:concepts:branching>` with ``@task.branch``.
+
+Run it
+------
+
+1. Install the provider with the SQL extra and the providers for your
databases:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai,sql]" \
+ apache-airflow-providers-postgres apache-airflow-providers-snowflake
+
+2. Create database connections ``postgres_source`` and ``snowflake_target``
that both
+ contain a ``customers`` table.
+
+3. Trigger the Dag:
+
+ .. code-block:: bash
+
+ airflow dags test example_llm_schema_compare_conditional
+
+The ``check_before_etl`` XCom holds the full comparison. Exactly one of
``run_etl`` and
+``notify_team`` runs.
+
+The Dag
+-------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_schema_compare.py
+ :language: python
+ :start-after: [START howto_operator_llm_schema_compare_conditional]
+ :end-before: [END howto_operator_llm_schema_compare_conditional]
+
+Adapting it
+-----------
+
+* Point ``db_conn_ids`` and ``table_names`` at your own source and target. The
operator
+ also compares against files on object storage, for example a Parquet landing
zone; see
+ :doc:`../operators/llm_schema_compare`.
+* Replace the ``run_etl`` placeholder with your load task and ``notify_team``
with a
+ Slack or email notifier.
+* Add an :doc:`../approval_gates` step before ``run_etl`` when a
compatible-but-changed
+ schema should still get a human's confirmation.
+* Set ``schedule`` to match the load and give the Dag the same ``start_date``
so the two
+ stay aligned.
diff --git a/providers/common/ai/docs/use_cases/index.rst
b/providers/common/ai/docs/use_cases/index.rst
new file mode 100644
index 00000000000..f6417a4581d
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/index.rst
@@ -0,0 +1,108 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+.. _howto/use_cases:
+
+What you can build
+==================
+
+Each page starts from a job a data team already has, shows the Dag that does
it, and names
+what Airflow adds over a script: a schedule, a retryable task per item, an
approval gate, or
+a fan-out sized at runtime. Every Dag ships with the provider and runs against
your
+connections.
+
+Start with :doc:`triage_support_tickets`. It needs nothing but an LLM
connection and shows
+the shape most of the others build on: structured output plus dynamic task
mapping.
+
+Every Dag needs the provider installed with the extra for your model vendor
and a
+``pydanticai`` connection named ``pydanticai_default``; :doc:`../quickstart`
covers both.
+Each page's "Run it" lists only what that Dag adds.
+
+.. list-table::
+ :header-rows: 1
+ :widths: 28 36 36
+
+ * - Scenario
+ - What the model does
+ - What Airflow does
+ * - :doc:`triage_support_tickets`
+ - Reads each ticket and returns priority, category, summary and next
action as a typed record
+ - One retryable task per ticket, results in XCom, one argument away from
a schedule
+ * - :doc:`route_pipeline_failures`
+ - Picks rerun, page or ignore from the error text, with a confidence score
+ - Runs only the chosen branch, sends low-confidence picks to a human
+ * - :doc:`gate_loads_on_schema_drift`
+ - Compares source and target schemas and reports what would break a load
+ - Branches to the load or to a notification, no migration ever runs
+ * - :doc:`explain_revenue_anomaly`
+ - Queries the warehouse and does the pivots in a sandbox to explain a
revenue move
+ - Injects the date, holds the credential, destroys the sandbox when the
run ends
+ * - :doc:`monthly_report_from_a_csv`
+ - Turns a fixed question into SQL over a CSV
+ - Downloads the file monthly, records schema changes before generating
SQL, emails the result
+ * - :doc:`compare_10k_filings`
+ - Splits a comparison question per company and writes the report
+ - Weekly indexing Dag, on-demand analysis Dag, per-company fan-out,
review at both ends
+ * - :doc:`ask_questions_over_pdfs`
+ - Answers a question from retrieved excerpts, citing them
+ - Weekly indexing Dag, query Dag triggered with a ``question`` param
+ * - :doc:`weekly_status_report`
+ - Assesses progress per proposal, then checks its own report against the
evidence
+ - Mapped evidence gathering, deterministic correction step, human review
+ * - :doc:`classify_reviews_in_bulk`
+ - Labels sentiment for every review in one batch job
+ - Waits up to 24 hours without holding a worker, lands results on object
storage
+ * - :doc:`research_agent_with_review`
+ - Decides which tools to call to answer a research question
+ - Human edits the question first, separate formatting step, approval
before delivery
+
+More ideas
+----------
+
+The same operators cover many other jobs. These do not have an example Dag
yet, but each
+can be built from the patterns shown on the pages above.
+
+* **Explain a failure in the alert.** A task with
``trigger_rule="one_failed"`` reads the
+ failed task's log, ``@task.llm`` returns a ``Literal`` root cause and a
two-sentence
+ explanation, and the notifier posts that instead of a stack trace.
+* **Data-quality triage.** Feed null rates, freshness and duplicate counts
from your checks
+ task to ``@task.llm`` with a ``Finding`` list as ``output_type``, then
+ ``LLMBranchOperator`` to page, file a ticket, or ignore.
+* **Daily incident digest.** Fetch alerts for the data interval, one
``@task.llm`` summary
+ per service with ``.expand()``, one synthesis call, ``ApprovalOperator``
before it posts.
+* **Release notes from merged pull requests.** Weekly ``HttpOperator`` fetch,
+ ``@task.llm_batch`` to classify each PR at batch prices, one call to draft
the notes, a
+ ``HITLEntryOperator`` for the editor's pass.
+* **Invoices into a table.** ``ObjectStoragePath`` lists new files,
``@task.llm_file_analysis``
+ extracts a typed row from each with ``.expand()``, a SQL operator inserts
them, and the
+ Dag emits an Asset so downstream reporting runs when the rows land.
+* **Rewrite a failing query.** On a SQL task's failure, hand the query and the
database
+ error to ``@task.llm_sql`` with the schema context, and put the rewrite in
front of a
+ reviewer before it runs.
+* **Tag and route incoming files.** ``@task.llm_file_analysis`` on each new
object decides
+ its type and sensitivity, ``@task.branch`` moves it to the right bucket.
+* **Catalog descriptions.** Nightly, for every table that changed,
``@task.llm`` writes a
+ column-level description from the schema and a sample, and a task upserts it
into the
+ catalog.
+
+:doc:`../examples` lists the same Dags by operator.
+
+.. note::
+
+ Dags with ``HITLEntryOperator`` or ``ApprovalOperator`` pause under
``airflow dags test``
+ until someone answers from Required Actions in the UI of an api-server on
the same
+ metadata database. ``airflow standalone`` gives you one.
diff --git a/providers/common/ai/docs/use_cases/monthly_report_from_a_csv.rst
b/providers/common/ai/docs/use_cases/monthly_report_from_a_csv.rst
new file mode 100644
index 00000000000..e3284500328
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/monthly_report_from_a_csv.rst
@@ -0,0 +1,103 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Monthly report from a survey CSV
+================================
+
+A published CSV, a stakeholder who wants the same question answered from it
every month,
+and nobody who wants to hand-write the SQL or learn from the report that a
column was
+renamed. This Dag downloads the file, checks its schema against a reference,
has the model
+write the SQL, runs it with Apache DataFusion, and emails the rows. Airflow
supplies the
+monthly schedule and a record of every schema change.
+
+The example uses the `Airflow community survey
<https://airflow.apache.org/survey/>`__
+CSV, which is public and needs no credentials.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/llm_sql` -- ``LLMSQLQueryOperator`` turns a question into
SQL
+ against a described schema, without executing it.
+* :doc:`../operators/llm_schema_compare` -- ``LLMSchemaCompareOperator``
records how the
+ downloaded file differs from a reference before any SQL is generated. It
reports; it does
+ not block (see Adapting it).
+* ``AnalyticsOperator`` from the ``common.sql`` provider -- runs the generated
SQL over a
+ local file with DataFusion, no database needed.
+* ``HttpOperator`` and ``SmtpHook`` from the ``http`` and ``smtp`` providers
do the
+ download and delivery.
+
+Run it
+------
+
+1. Install the provider with the SQL extra, plus the HTTP and SMTP providers:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai,sql]" \
+ "apache-airflow-providers-common-sql[datafusion]" \
+ apache-airflow-providers-http apache-airflow-providers-smtp
+
+2. Create an HTTP connection named ``airflow_website`` with host
+ ``https://airflow.apache.org`` and no auth.
+
+3. Optionally set ``SMTP_CONN_ID`` and ``NOTIFY_EMAIL`` in the environment.
Without
+ them the result goes to the task log.
+
+4. Trigger the Dag:
+
+ .. code-block:: bash
+
+ airflow dags test example_llm_survey_scheduled
+
+The ``check_schema`` XCom holds the comparison, ``generate_sql`` holds the
query the model
+wrote, and ``send_result`` logs or mails the rows.
+
+The Dag
+-------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py
+ :language: python
+ :start-after: [START example_llm_survey_scheduled]
+ :end-before: [END example_llm_survey_scheduled]
+
+Variants
+--------
+
+``example_llm_survey_interactive`` in the same file takes ad hoc questions: a
+``HITLEntryOperator`` edits the question, an ``ApprovalOperator`` reviews the
rows. No
+download, no schedule; the CSV is assumed in place.
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py
+ :language: python
+ :start-after: [START example_llm_survey_interactive]
+ :end-before: [END example_llm_survey_interactive]
+
+When one question is not enough, ``example_llm_survey_agentic`` splits a
research question
+into sub-questions, maps SQL generation and execution over them, and
synthesizes the
+results behind an approval gate.
+
+Adapting it
+-----------
+
+* Make the schema check block: add a ``@task.branch`` on ``compatible`` between
+ ``check_schema`` and ``generate_sql``, as :doc:`gate_loads_on_schema_drift`
does.
+* Point ``SURVEY_CSV_ENDPOINT`` and the HTTP connection at your own published
file, and
+ replace ``SURVEY_SCHEMA`` with a description of its columns; that
description is what
+ the model writes SQL against.
+* Change ``SCHEDULED_PROMPT`` to the question the report answers.
+* Swap the local-file ``DataSourceConfig`` for a warehouse connection when the
data
+ lives in a database rather than a file; ``LLMSQLQueryOperator`` accepts
either.
diff --git a/providers/common/ai/docs/use_cases/research_agent_with_review.rst
b/providers/common/ai/docs/use_cases/research_agent_with_review.rst
new file mode 100644
index 00000000000..a92d839540e
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/research_agent_with_review.rst
@@ -0,0 +1,88 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Research agent with human review
+================================
+
+Someone has a question that needs a knowledge base, a dataset and a web search
to answer,
+and the right sequence of lookups depends on the question. This Dag runs a
LangChain ReAct
+agent that decides for itself which tools to call and in what order, then
hands the raw
+findings to a separate formatting step and to a reviewer. Airflow puts a
person in front
+of the agent to edit the question, exposes the findings and tool calls as
XCom, makes
+formatting its own retryable task, and holds the report for approval.
+
+This is the shape for teams with existing LangChain tools. The agent loop is
LangChain's;
+the connection, the formatting call and the review gates are Airflow's.
+
+What this demonstrates
+----------------------
+
+* :doc:`../hooks/langchain` -- ``LangChainHook`` supplies the chat model and
embeddings
+ from one Airflow connection.
+* :doc:`../operators/llm` -- ``LLMOperator`` formats the findings with a Jinja
template
+ that reads the agent's XCom.
+* :doc:`apache-airflow-providers-standard:operators/hitl` --
``HITLEntryOperator`` to edit
+ the question, ``ApprovalOperator`` to sign off the report.
+* :doc:`../toolsets/langchain` -- the reverse direction, giving a LangChain
agent tools
+ built from Airflow connections.
+
+Run it
+------
+
+1. Install the LangChain extra and the packages the tools use:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[langchain]" \
+ langchain-openai langchain-text-splitters langchain-community
faiss-cpu
+
+2. Create a ``langchain`` connection named ``langchain_default`` with your API
key in the
+ password field (see :doc:`../connections/langchain`).
+
+3. Optionally put documents under ``DOCS_PATH`` and the survey CSV at
``SURVEY_CSV_PATH``;
+ without them the Dag writes sample pages so the tools have something to
search.
+
+4. Trigger the Dag:
+
+ .. code-block:: bash
+
+ airflow dags test example_langchain_tool_agent
+
+The run pauses at ``prompt_review`` and ``report_approval``; answer from
Required Actions
+in the UI (see the note on :doc:`index`). The ``run_research_agent`` log shows
every tool
+call, and its XCom keeps the findings and call list.
+
+The Dag
+-------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py
+ :language: python
+ :start-after: [START example_langchain_tool_agent]
+ :end-before: [END example_langchain_tool_agent]
+
+The tools (knowledge-base search, survey query, a stubbed web search, a clock)
are ordinary
+LangChain ``@tool`` functions built in ``_build_tools``.
+
+Adapting it
+-----------
+
+* Replace ``_build_tools`` with your own LangChain tools. Anything from the
LangChain
+ ecosystem works unchanged.
+* To run the same shape on pydantic-ai instead, use :doc:`../operators/agent`
with
+ :doc:`../toolsets/index`; the surrounding review and formatting tasks stay
as they are.
+* For routine questions, drop ``prompt_review`` and take the question from a
Dag param;
+ keep ``report_approval``.
diff --git a/providers/common/ai/docs/use_cases/route_pipeline_failures.rst
b/providers/common/ai/docs/use_cases/route_pipeline_failures.rst
new file mode 100644
index 00000000000..046961a8cbf
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/route_pipeline_failures.rst
@@ -0,0 +1,88 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Route pipeline failures to a fix or a person
+============================================
+
+A task failed overnight and the on-call engineer has to read the error, decide
whether it
+was a blip worth rerunning, something a person has to fix now, or noise to
ignore, and
+then act. This Dag hands the reading to a model that returns a pick and how
sure it is.
+Airflow runs only the chosen branch, sends uncertain picks to a human with a
deadline, and
+demands more confidence to page someone than to rerun a task.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/llm_branch` -- ``LLMBranchOperator`` picks one downstream
task id.
+* :doc:`../classifier_models` -- a classifier model, so the gate has a
confidence to read.
+* :doc:`../approval_gates` -- ``on_uncertain="review"`` routes low-confidence
picks to a
+ human, with ``approval_timeout`` bounding the wait.
+
+Run it
+------
+
+1. Install the provider with the classifier extra:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[typesafe]"
+
+2. Set ``{"model": "typesafe:jev-1.13.0"}`` in the ``pydanticai_default``
extra. The second
+ Dag below reads the same kind of connection as ``jev_default``.
+
+3. Trigger the Dag:
+
+ .. code-block:: bash
+
+ airflow dags test example_llm_branch_decision_policy
+
+The ``triage_failure`` log shows the pick and its confidence. Exactly one of
``rerun``,
+``page_oncall`` and ``ignore`` runs. Below the confidence bar the run pauses
for review
+instead (Airflow 3.1 or later); answer from Required Actions in the UI (see
the note on
+:doc:`index`).
+
+The Dag
+-------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_branch.py
+ :language: python
+ :start-after: [START howto_operator_llm_branch_decision_policy]
+ :end-before: [END howto_operator_llm_branch_decision_policy]
+
+Classify-then-act variant
+-------------------------
+
+When the action depends on the score itself, classify in one task and act in
the next
+(:doc:`../classifier_models` explains reading the score). It uses the
``jev_default``
+connection; run it as ``airflow dags test
example_classifier_model_confidence``:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_classifier_model.py
+ :language: python
+ :start-after: [START howto_classifier_model_confidence]
+ :end-before: [END howto_classifier_model_confidence]
+
+Adapting it
+-----------
+
+* Replace the hard-coded ``prompt`` with the failed task's own error. An
+ ``on_failure_callback`` on the production Dag can trigger this one with the
exception
+ text in ``conf``, and the prompt reads ``{{ dag_run.conf["error"] }}``.
+* Make the branch tasks do the work: ``rerun`` clears the failed task instance
through
+ the REST API, ``page_oncall`` posts to your alerting provider.
+* Tune ``min_confidence`` per branch. A wrong page costs more than an extra
rerun, so
+ ``page_oncall`` demands more certainty.
+* To decide retry versus fail inside Airflow's retry loop, see
:doc:`../retry_policies`.
diff --git a/providers/common/ai/docs/use_cases/triage_support_tickets.rst
b/providers/common/ai/docs/use_cases/triage_support_tickets.rst
new file mode 100644
index 00000000000..1d14310736c
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/triage_support_tickets.rst
@@ -0,0 +1,72 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Triage support tickets
+======================
+
+A queue of free-text support tickets arrives every hour and someone has to
read each one,
+decide how urgent it is, and hand it to the right team. This Dag has the model
do the
+reading and produce a typed record per ticket: priority, category, a one-line
summary and
+a suggested next action. Airflow runs one task per ticket, so a bad ticket
retries alone,
+every result is in XCom, and one ``schedule`` argument makes it hourly.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/llm` -- ``@task.llm`` returns a prompt string; the
operator makes
+ the call and pushes the result.
+* :doc:`../structured_output` -- ``output_type=TicketAnalysis`` gives the
downstream task
+ a Pydantic instance, not a string to parse.
+* :doc:`apache-airflow:authoring-and-scheduling/dynamic-task-mapping` --
``.expand()``
+ fans one model call out per ticket.
+
+Run it
+------
+
+1. Install the provider with the extra for your vendor:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai]"
+
+2. Trigger the Dag:
+
+ .. code-block:: bash
+
+ airflow dags test example_llm_analysis_pipeline
+
+The ``store_results`` log has one ``[PRIORITY] category: summary`` line per
ticket; each
+``analyze_ticket`` map index holds a ``TicketAnalysis`` XCom.
+
+The Dag
+-------
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_llm_analysis_pipeline.py
+ :language: python
+ :start-after: [START howto_decorator_llm_pipeline]
+ :end-before: [END howto_decorator_llm_pipeline]
+
+Adapting it
+-----------
+
+* Replace the list in ``get_support_tickets`` with a query against your
ticketing system,
+ for example an ``SQLExecuteQueryOperator`` or an ``HttpOperator`` upstream.
+* Set ``schedule="@hourly"`` on the ``@dag`` and use ``{{ data_interval_start
}}`` in
+ the query so each run picks up only new tickets.
+* Make ``priority`` and ``category`` ``Literal`` types on ``TicketAnalysis``
so the
+ model cannot invent a label. :doc:`route_pipeline_failures` shows how to
branch on the
+ result.
diff --git a/providers/common/ai/docs/use_cases/weekly_status_report.rst
b/providers/common/ai/docs/use_cases/weekly_status_report.rst
new file mode 100644
index 00000000000..db17cd60046
--- /dev/null
+++ b/providers/common/ai/docs/use_cases/weekly_status_report.rst
@@ -0,0 +1,120 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ .. http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+
+Weekly status report with a hallucination check
+===============================================
+
+A program manager wants a weekly status on a set of proposals: promised,
landed, still
+open. The evidence is spread across a wiki and a repository, and people act on
the report,
+so it has to be right. This Dag gathers the evidence, has the model assess
each proposal,
+synthesizes a report, then has a second model call judge it against the
evidence and a
+plain-Python step apply only the flagged corrections. A person reviews the
result.
+
+Airflow maps gathering and assessment per proposal, keeps every intermediate
in XCom, and
+makes the correction step deterministic so the model cannot rewrite the report
while
+fixing it.
+
+The example tracks Airflow Improvement Proposals against Confluence and
GitHub, both
+public, and solves the job a second way with one autonomous agent for
comparison.
+
+What this demonstrates
+----------------------
+
+* :doc:`../operators/llm` -- mapped ``LLMOperator`` calls with structured
output for the
+ per-proposal analysis, one bounded by ``UsageLimits`` for the synthesis, and
one whose
+ only job is validation.
+* :doc:`../structured_output` -- ``ValidationResult`` lists each claim with a
verdict, so
+ the correction step has something exact to act on.
+* :doc:`apache-airflow-providers-standard:operators/hitl` --
``ApprovalOperator`` before the
+ report goes out.
+* :doc:`../toolsets/skills` -- the agent variant loads a ``SKILL.md`` bundle
that tells it
+ how to assess progress.
+
+Run it
+------
+
+1. Install the provider. The agent variant also needs the skills extra:
+
+ .. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[openai,skills]"
+
+2. Optionally set ``GITHUB_TOKEN`` in the environment. Without it the Dag
paces itself
+ to the unauthenticated rate limit and takes longer.
+
+3. Trigger either Dag. Both take an ``aip_numbers`` param:
+
+ .. code-block:: bash
+
+ airflow dags test example_aip_progress_tracker
+ airflow dags test example_aip_progress_tracker_skills
+
+The ``validate_report`` XCom shows the disputed claims and
``apply_validation`` shows what
+changed. The run pauses at ``review_report``; answer from Required Actions in
the UI (see the
+note on :doc:`index`).
+
+The Dag
+-------
+
+The validation output type and the step that applies it:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
+ :language: python
+ :start-after: [START aip_tracker_validation_output]
+ :end-before: [END aip_tracker_validation_output]
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
+ :language: python
+ :start-after: [START aip_tracker_validation]
+ :end-before: [END aip_tracker_validation]
+
+The mapped analysis that feeds it:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
+ :language: python
+ :start-after: [START aip_tracker_dtm_analysis]
+ :end-before: [END aip_tracker_dtm_analysis]
+
+The full Dag, evidence gathering included, is ``example_aip_progress_tracker``
in
+`example_aip_progress_tracker.py
<https://github.com/apache/airflow/blob/providers-common-ai/|version|/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py>`__.
+It is long because the evidence gathering is real.
+
+The agent variant
+-----------------
+
+The same job as one operator call. The agent reads a skill that explains how
to assess a
+proposal, gets the wiki and repository as tool functions, and decides its own
order of
+work:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
+ :language: python
+ :start-after: [START aip_tracker_skills_operator]
+ :end-before: [END aip_tracker_skills_operator]
+
+Use the pipeline when every step must be auditable. Use the agent when a fixed
task graph
+would only re-encode what the model can work out from the skill.
+
+Adapting it
+-----------
+
+* Replace the Confluence and GitHub fetchers with your own sources: a project
tracker, a
+ design-doc folder, a deployment log. Everything from ``analyze_aip`` down is
+ source-agnostic.
+* Keep the validation step even if you drop the rest. It is cheap and catches
the failures
+ that make people stop trusting generated reports.
+* Put it on ``schedule="@weekly"`` and send the approved report to a channel
from a task
+ after ``review_report``.
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_10k.py
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_10k.py
index 3a7b99556a1..fd2d8ccb199 100644
---
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_10k.py
+++
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_10k.py
@@ -45,7 +45,7 @@ detail while Airflow provides the orchestration.
.. code-block:: text
- analyst_question (HITLEntryOperator)
+ analyst_input (HITLEntryOperator)
-> get_question (@task)
-> get_tickers (@task)
-> decompose_question (@task.llm, structured output)
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
index 4142ac0bbd3..11e04b86197 100644
---
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
+++
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llamaindex_10k.py
@@ -38,7 +38,7 @@ is supported -- configure via the ``tickers`` Dag parameter.
.. code-block:: text
- analyst_question (HITLEntryOperator)
+ analyst_input (HITLEntryOperator)
-> get_question (@task)
-> get_tickers (@task)
-> decompose_question (@task.llm, structured output)
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_analysis_pipeline.py
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_analysis_pipeline.py
index f396b53e4d0..b13dc811ae9 100644
---
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_analysis_pipeline.py
+++
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_analysis_pipeline.py
@@ -23,6 +23,7 @@ from pydantic import BaseModel
from airflow.providers.common.compat.sdk import dag, task
+# [START howto_decorator_llm_pipeline]
# Pydantic output classes must be defined at module scope so they can be
# imported by name when downstream tasks deserialize the XCom payload.
class TicketAnalysis(BaseModel):
@@ -34,9 +35,10 @@ class TicketAnalysis(BaseModel):
suggested_action: str
-# [START howto_decorator_llm_pipeline]
@dag(tags=["example"])
def example_llm_analysis_pipeline():
+ """Triage a queue of support tickets: one model call per ticket, typed
results, ready for a schedule."""
+
@task
def get_support_tickets():
"""Fetch unprocessed support tickets."""
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py
index 1c239c9ab82..598ab7a59ec 100644
---
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py
+++
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_llm_survey_analysis.py
@@ -299,8 +299,8 @@ if LLMSQLQueryOperator is not None:
csv_ready = prepare_csv(download_survey.output)
# ------------------------------------------------------------------
- # Step 3: Validate the downloaded CSV schema against the reference.
- # Raises if critical columns are missing or renamed.
+ # Step 3: Compare the downloaded CSV schema against the reference.
+ # Reports what changed; add a @task.branch on ``compatible`` to block
on drift.
# ------------------------------------------------------------------
check_schema = LLMSchemaCompareOperator(
task_id="check_schema",