tracing-upstream-lineage

Rastreia a linhagem upstream de dados para identificar fontes, DAGs e dependências que alimentam uma tabela ou coluna. Suporta rastreamento de três tipos de destino: tabelas, colunas e DAGs; utiliza o código-fonte do DAG do Airflow e inspeção de tarefas para encontrar pipelines produtores. Lida com fontes SQL (cláusulas FROM), sistemas externos (S3, Postgres, Salesforce, APIs HTTP) e fontes baseadas em arquivos; rastreia recursivamente cadeias upstream. Inclui rastreamento em nível de coluna por meio de mapeamentos diretos, transformações e agregações no código do DAG...

npx skills add https://github.com/astronomer/agents --skill tracing-upstream-lineage

Upstream Lineage: Sources

Trace the origins of data - answer "Where does this data come from?"

Lineage Investigation

Step 1: Identify the Target Type

Determine what we're tracing:

  • Table: Trace what populates this table
  • Column: Trace where this specific column comes from
  • DAG: Trace what data sources this DAG reads from

Step 2: Find the Producing DAG

Tables are typically populated by Airflow DAGs. Find the connection:

  1. Search DAGs by name: Use af dags list and look for DAG names matching the table name

    • load_customers -> customers table
    • etl_daily_orders -> orders table
  2. Explore DAG source code: Use af dags source <dag_id> to read the DAG definition

    • Look for INSERT, MERGE, CREATE TABLE statements
    • Find the target table in the code
  3. Check DAG tasks: Use af tasks list <dag_id> to see what operations the DAG performs

On Astro

If you're running on Astro, the Lineage tab in the Astro UI provides visual lineage exploration across DAGs and datasets. Use it to quickly trace upstream dependencies without manually searching DAG source code.

On OSS Airflow

Use DAG source code and task logs to trace lineage (no built-in cross-DAG UI).

Step 3: Trace Data Sources

From the DAG code, identify source tables and systems:

SQL Sources (look for FROM clauses):

# In DAG code:
SELECT * FROM source_schema.source_table  # <- This is an upstream source

External Sources (look for connection references):

  • S3Operator -> S3 bucket source
  • PostgresOperator -> Postgres database source
  • SalesforceOperator -> Salesforce API source
  • HttpOperator -> REST API source

File Sources:

  • CSV/Parquet files in object storage
  • SFTP drops
  • Local file paths

Step 4: Build the Lineage Chain

Recursively trace each source:

TARGET: analytics.orders_daily
    ^
    +-- DAG: etl_daily_orders
            ^
            +-- SOURCE: raw.orders (table)
            |       ^
            |       +-- DAG: ingest_orders
            |               ^
            |               +-- SOURCE: Salesforce API (external)
            |
            +-- SOURCE: dim.customers (table)
                    ^
                    +-- DAG: load_customers
                            ^
                            +-- SOURCE: PostgreSQL (external DB)

Step 5: Check Source Health

For each upstream source:

  • Tables: Check freshness with the checking-freshness skill
  • DAGs: Check recent run status with af dags stats
  • External systems: Note connection info from DAG code

Lineage for Columns

When tracing a specific column:

  1. Find the column in the target table schema
  2. Search DAG source code for references to that column name
  3. Trace through transformations:
    • Direct mappings: source.col AS target_col
    • Transformations: COALESCE(a.col, b.col) AS target_col
    • Aggregations: SUM(detail.amount) AS total_amount

Output: Lineage Report

Summary

One-line answer: "This table is populated by DAG X from sources Y and Z"

Lineage Diagram

[Salesforce] --> [raw.opportunities] --> [stg.opportunities] --> [fct.sales]
                        |                        |
                   DAG: ingest_sfdc         DAG: transform_sales

Source Details

SourceTypeConnectionFreshnessOwner
raw.ordersTableInternal2h agodata-team
SalesforceAPIsalesforce_connReal-timesales-ops

Transformation Chain

Describe how data flows and transforms:

  1. Raw data lands in raw.orders via Salesforce API sync
  2. DAG transform_orders cleans and dedupes into stg.orders
  3. DAG build_order_facts joins with dimensions into fct.orders

Data Quality Implications

  • Single points of failure?
  • Stale upstream sources?
  • Complex transformation chains that could break?

Related Skills

  • Check source freshness: checking-freshness skill
  • Debug source DAG: debugging-dags skill
  • Trace downstream impacts: tracing-downstream-lineage skill
  • Add manual lineage annotations: annotating-task-lineage skill
  • Build custom lineage extractors: creating-openlineage-extractors skill

Mais skills de astronomer

airflow-state-store
astronomer
Persists task and asset state across retries and DAG runs using Airflow 3.3's AIP-103 key/value stores (`task_state_store`, `asset_state_store`) and the…
creating-openlineage-extractors
astronomer
Extratores OpenLineage personalizados para operadores Airflow não suportados e cenários complexos de linhagem. Duas abordagens: adicionar métodos OpenLineage diretamente aos operadores que você possui (recomendado), ou criar extratores personalizados para operadores de terceiros que você não pode modificar. Os extratores interceptam a execução do operador em três pontos: antes da execução para linhagem estática, após o sucesso para saídas determinadas em tempo de execução e, opcionalmente, após falha para linhagem parcial. Registre os extratores via airflow.cfg ou ambiente...
debugging-dags
astronomer
Análise sistemática de causa raiz e remediação para DAGs do Airflow com falhas, utilizando fluxos de investigação estruturados. Orienta por um processo de diagnóstico em quatro etapas: identificar a falha, extrair detalhes do erro, reunir informações contextuais e fornecer etapas de remediação acionáveis. Classifica as falhas em quatro tipos (dados, código, infraestrutura, dependência) para focar a investigação e sugerir correções apropriadas. Fornece comandos CLI prontos para recuperação de logs, comparação de execuções, limpeza de tarefas e DAG...
delegating-to-otto
astronomer
Direciona o agente Otto da
deploying-airflow
astronomer
Implantar DAGs e projetos do Airflow. Use quando o usuário quiser implantar código, enviar DAGs, configurar CI/CD, implantar em produção ou perguntar sobre estratégias de implantação…
deploying-go-sdk-bundles
astronomer
Compila, empacota e implanta pacotes compilados do Airflow Go SDK para que o ExecutableCoordinator possa executá-los. Use quando o usuário quiser compilar um pacote de tarefas Go, pedir…
testing-dags
astronomer
Ciclos iterativos de teste-depuração-correção para DAGs do Airflow com diagnóstico abrangente de falhas. Comece com af runs trigger-wait <dag_id> para executar um DAG e aguardar a conclusão; não são necessárias verificações prévias. Em caso de falha, use af runs diagnose para um resumo abrangente de falhas e af tasks logs para inspecionar detalhes de erros de tarefas específicas. Suporta configuração personalizada, timeouts e tentativas de repetição; lida com cenários de sucesso, falha e timeout com interpretação clara da resposta. Validação rápida disponível...
tracing-downstream-lineage
astronomer
Rastreie a linhagem de dados downstream para avaliar o impacto de alterações antes de modificar tabelas ou DAGs. Identifica consumidores diretos de uma tabela ou DAG alvo por meio de busca em código-fonte, dependências de views e conexões com ferramentas de BI. Constrói uma árvore de dependências completa mapeando todos os impactos downstream, desde tabelas até dashboards e modelos de ML. Categoriza dependências por criticidade (crítica, alta, média, baixa) para priorizar comunicação com stakeholders e testes. Gera um relatório de impacto com avaliação de risco, afetados...