tracing-downstream-lineage

tarafından astronomer

Tablo veya DAG'leri değiştirmeden önce aşağı yönlü veri soyunu izleyerek değişiklik etkisini değerlendirir. Kaynak kod araması, görünüm bağımlılıkları ve BI araç bağlantıları aracılığıyla hedef tablo veya DAG'in doğrudan tüketicilerini belirler. Tablolardan panolara ve ML modellerine kadar tüm aşağı yönlü etkileri haritalayan tam bir bağımlılık ağacı oluşturur. Bağımlılıkları kritiklik düzeyine (kritik, yüksek, orta, düşük) göre kategorize ederek paydaş iletişimi ve test önceliklendirmesini sağlar. Risk değerlendirmesi ve etkilenen... içeren

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

Downstream Lineage: Impacts

Answer the critical question: "What breaks if I change this?"

Use this BEFORE making changes to understand the blast radius.

Impact Analysis

Step 1: Identify Direct Consumers

Find everything that reads from this target:

For Tables:

  1. Search DAG source code: Look for DAGs that SELECT from this table

    • Use af dags list to get all DAGs
    • Use af dags source <dag_id> to search for table references
    • Look for: FROM target_table, JOIN target_table
  2. Check for dependent views:

    -- Snowflake
    SELECT * FROM information_schema.view_table_usage
    WHERE table_name = '<target_table>'
    
    -- Or check SHOW VIEWS and search definitions
    
  3. Look for BI tool connections:

    • Dashboards often query tables directly
    • Check for common BI patterns in table naming (rpt_, dashboard_)

On Astro

If you're running on Astro, the Lineage tab in the Astro UI provides visual dependency graphs across DAGs and datasets, making downstream impact analysis faster. It shows which DAGs consume a given dataset and their current status, reducing the need for manual source code searches.

For DAGs:

  1. Check what the DAG produces: Use af dags source <dag_id> to find output tables
  2. Then trace those tables' consumers (recursive)

Step 2: Build Dependency Tree

Map the full downstream impact:

SOURCE: fct.orders
    |
    +-- TABLE: agg.daily_sales --> Dashboard: Executive KPIs
    |       |
    |       +-- TABLE: rpt.monthly_summary --> Email: Monthly Report
    |
    +-- TABLE: ml.order_features --> Model: Demand Forecasting
    |
    +-- DIRECT: Looker Dashboard "Sales Overview"

Step 3: Categorize by Criticality

Critical (breaks production):

  • Production dashboards
  • Customer-facing applications
  • Automated reports to executives
  • ML models in production
  • Regulatory/compliance reports

High (causes significant issues):

  • Internal operational dashboards
  • Analyst workflows
  • Data science experiments
  • Downstream ETL jobs

Medium (inconvenient):

  • Ad-hoc analysis tables
  • Development/staging copies
  • Historical archives

Low (minimal impact):

  • Deprecated tables
  • Unused datasets
  • Test data

Step 4: Assess Change Risk

For the proposed change, evaluate:

Schema Changes (adding/removing/renaming columns):

  • Which downstream queries will break?
  • Are there SELECT * patterns that will pick up new columns?
  • Which transformations reference the changing columns?

Data Changes (values, volumes, timing):

  • Will downstream aggregations still be valid?
  • Are there NULL handling assumptions that will break?
  • Will timing changes affect SLAs?

Deletion/Deprecation:

  • Full dependency tree must be migrated first
  • Communication needed for all stakeholders

Step 5: Find Stakeholders

Identify who owns downstream assets:

  1. DAG owners: Check owners field in DAG definitions
  2. Dashboard owners: Usually in BI tool metadata
  3. Team ownership: Look for team naming patterns or documentation

Output: Impact Report

Summary

"Changing fct.orders will impact X tables, Y DAGs, and Z dashboards"

Impact Diagram

                    +--> [agg.daily_sales] --> [Executive Dashboard]
                    |
[fct.orders] -------+--> [rpt.order_details] --> [Ops Team Email]
                    |
                    +--> [ml.features] --> [Demand Model]

Detailed Impacts

DownstreamTypeCriticalityOwnerNotes
agg.daily_salesTableCriticaldata-engUpdated hourly
Executive DashboardDashboardCriticalanalyticsCEO views daily
ml.order_featuresTableHighml-teamRetraining weekly

Risk Assessment

Change TypeRisk LevelMitigation
Add columnLowNo action needed
Rename columnHighUpdate 3 DAGs, 2 dashboards
Delete columnCriticalFull migration plan required
Change data typeMediumTest downstream aggregations

Recommended Actions

Before making changes:

  1. Notify owners: @data-eng, @analytics, @ml-team
  2. Update downstream DAG: transform_daily_sales
  3. Test dashboard: Executive KPIs
  4. Schedule change during low-impact window

Related Skills

  • Trace where data comes from: tracing-upstream-lineage skill
  • Check downstream freshness: checking-freshness skill
  • Debug any broken DAGs: debugging-dags skill
  • Add manual lineage annotations: annotating-task-lineage skill
  • Build custom lineage extractors: creating-openlineage-extractors skill

astronomer tarafından daha fazla skill

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
Özel OpenLineage çıkarıcıları, desteklenmeyen Airflow operatörleri ve karmaşık lineage senaryoları için. İki yaklaşım: sahip olduğunuz operatörlere doğrudan OpenLineage yöntemleri eklemek (önerilen) veya değiştiremeyeceğiniz üçüncü taraf operatörler için özel çıkarıcılar oluşturmak. Çıkarıcılar, operatör yürütmesini üç noktada keser: statik lineage için yürütmeden önce, çalışma zamanında belirlenen çıktılar için başarıdan sonra ve isteğe bağlı olarak kısmi lineage için başarısızlıktan sonra. Çıkarıcıları airflow.cfg veya ortam değişkenleri aracı
debugging-dags
astronomer
Sistematik kök neden analizi ve başarısız Airflow DAG'leri için yapılandırılmış soruşturma iş akışlarıyla düzeltme. Dört adımlı teşhis sürecinde rehberlik eder: hatayı belirleme, hata ayrıntılarını çıkarma, bağlamsal bilgi toplama ve uygulanabilir düzeltme adımları sunma. Hataları dört türe (veri, kod, altyapı, bağımlılık) kategorize ederek soruşturmayı odaklar ve uygun düzeltmeler önerir. Günlük alma, çalıştırma karşılaştırması, görev temizleme ve DAG... için kullanıma hazır CLI komutları sağlar.
delegating-to-otto
astronomer
Drives Astronomer's Otto agent (`astro otto`) as a delegated sub-agent for Airflow, dbt, and data-engineering work. Use when the user explicitly asks to "use…
deploying-airflow
astronomer
Airflow DAG'larını ve projelerini dağıtın. Kullanıcı kod dağıtmak, DAG'ları göndermek, CI/CD kurmak, üretime dağıtmak veya dağıtım stratejileri hakkında soru sorduğunda kullanın…
deploying-go-sdk-bundles
astronomer
Derlenmiş Airflow Go SDK paketlerini derler, paketler ve dağıtır, böylece ExecutableCoordinator bunları çalıştırabilir. Kullanıcı bir Go görev paketini derlemek istediğinde, sorduğunda…
testing-dags
astronomer
Airflow DAG'leri için kapsamlı hata teşhisi ile yinelemeli test-hata ayıklama-düzeltme döngüleri. Bir DAG'ı çalıştırmak ve tamamlanmasını beklemek için af runs trigger-wait <dag_id> ile başlayın; ön kontrol gerekmez. Hata durumunda, kapsamlı hata özeti için af runs diagnose ve belirli görevlerden hata detaylarını incelemek için af tasks logs kullanın. Özel yapılandırma, zaman aşımları ve yeniden deneme girişimlerini destekler; başarı, hata ve zaman aşımı senaryolarını net yanıt yorumlama ile ele alır. Hızlı doğrulama mevcuttur...
tracing-upstream-lineage
astronomer
Bir tabloyu veya sütunu besleyen kaynakları, DAG'leri ve bağımlılıkları belirlemek için upstream veri soyunu izler. Üç hedef türünün izlenmesini destekler: tablolar, sütunlar ve DAG'ler; üreten pipeline'ları bulmak için Airflow DAG kaynak kodu ve görev incelemesi kullanır. SQL kaynaklarını (FROM cümleleri), harici sistemleri (S3, Postgres, Salesforce, HTTP API'leri) ve dosya tabanlı kaynakları işler; upstream zincirlerini yinelemeli olarak izler. DAG kodundaki doğrudan eşlemeler, dönüşümler ve toplamalar aracılığıyla sütun düzeyinde izleme içerir...