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 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:DocumentLoaderOperatorparses local files,LlamaIndexEmbeddingOperatorchunks and persists the index,LlamaIndexRetrievalOperatorretrieves for a fixed question, and anLLMOperatorsynthesizes 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 aquestionparam) 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– twoDocumentLoaderOperatorcalls tag documents from different sources viametadata_fieldsbefore merging and embedding them into one index, for filtered retrieval downstream.
See DocumentLoaderOperator, LlamaIndex LlamaIndexEmbeddingOperator, and LlamaIndex LlamaIndexRetrievalOperator for how the individual operators work.
SEC 10-K financial analysis¶
example_llamaindex_10k.py and 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 LangChainHook and FAISS
directly.
Analysis Dag flow:
A human confirms or edits the comparison question and the tickers to compare (
HITLEntryOperator).@task.llmdecomposes the question into one sub-question per company – the LLM decides how many sub-questions are needed, not the Dag author:@task.llm( llm_conn_id=LLM_CONN_ID, system_prompt=DECOMPOSE_SYSTEM_PROMPT, output_type=DecomposedQuestion, # Push the structured output to XCom as a dict so the example runs on # every supported Airflow version (the model-instance form needs 3.3+). serialize_output=True, ) def decompose_question(question: str, tickers: str) -> str: return ( f"Decompose this question into company-specific sub-questions.\n" f"Available companies (by ticker): {tickers}\n\n" f"Question: {question}" ) decomposed = decompose_question(question, tickers)
Dynamic Task Mapping fans retrieval out, one mapped task per sub-question, each querying that company’s index:
retrieval_results = LlamaIndexRetrievalOperator.partial( task_id="retrieve", embed_model="text-embedding-3-small", llm_conn_id=LLAMAINDEX_CONN_ID, top_k=5, ).expand_kwargs(retrieval_kwargs)
The mapped results are zipped back to their sub-questions (mapped task outputs preserve input order) and synthesized into a structured
AnalysisReportby anLLMOperatorbounded byUsageLimits.ApprovalOperatorgates the report before hand-off.
See LlamaIndex LlamaIndexEmbeddingOperator, LlamaIndex LlamaIndexRetrievalOperator, and LangChainHook for how the individual operators/hooks behave.
Natural-language survey analysis¶
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 LLMSQLQueryOperator 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 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:
decompose_questionreturns one sub-question per dimension.LLMSQLQueryOperator.partial(...).expand(prompt=sub_questions)generates SQL for each dimension as an independent mapped task instance.AnalyticsOperator.partial(...).expand(...)runs each query. If one dimension’s query fails, only that mapped instance retries – the other three keep their results.collect_resultszips the dimension labels back onto the (order-preserved) results.An
LLMOperatorsynthesizes the four labeled result sets into one narrative, thenApprovalOperatorgates 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 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-mappedLLMOperatorcalls analyze each AIP with structured output; the per-AIP analyses are synthesized into one report; a secondLLMOperatorvalidates 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 singleAgentOperatorloaded 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:
validate = LLMOperator(
task_id="validate_report",
llm_conn_id=LLM_CONN_ID,
system_prompt=VALIDATION_SYSTEM_PROMPT,
prompt="""\
Verify the following synthesized report against the raw per-AIP evidence.
Flag any claims not grounded in the evidence.
=== SYNTHESIZED REPORT ===
{{ ti.xcom_pull(task_ids='synthesize_report') }}
=== RAW PER-AIP EVIDENCE ===
{{ ti.xcom_pull(task_ids='format_report') }}""",
output_type=ValidationResult,
serialize_output=True,
usage_limits=UsageLimits(
request_limit=5,
input_tokens_limit=30_000,
output_tokens_limit=8_000,
),
agent_params={"model_settings": {"temperature": 0}},
)
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:
report = AgentOperator(
task_id="track_aip_progress",
llm_conn_id=LLM_CONN_ID,
system_prompt=AGENT_SYSTEM_PROMPT,
prompt=prompt,
toolsets=[
AgentSkillsToolset(sources=[SKILLS_DIR]),
aip_toolset,
],
agent_params={"model_settings": {"temperature": 0}},
usage_limits=UsageLimits(
request_limit=30,
input_tokens_limit=200_000,
output_tokens_limit=16_000,
),
)
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.