tracing-upstream-lineage

โดย astronomer

ติดตามสายข้อมูลต้นทางเพื่อระบุแหล่งที่มา DAG และ dependencies ที่ป้อนเข้าสู่ตารางหรือคอลัมน์ รองรับการติดตามเป้าหมายสามประเภท: ตาราง คอลัมน์ และ DAG; ใช้ซอร์สโค้ด Airflow DAG และการตรวจสอบงานเพื่อค้นหาไปป์ไลน์ที่ผลิตข้อมูล จัดการกับแหล่ง SQL (คำสั่ง FROM), ระบบภายนอก (S3, Postgres, Salesforce, HTTP APIs) และแหล่งที่ใช้ไฟล์; ติดตามสายโซ่ต้นทางแบบเรียกซ้ำ รวมถึงการติดตามระดับคอลัมน์ผ่านการแมปโดยตรง การแปลง และการรวมในโค้ด 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

Skills เพิ่มเติมจาก 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
ตัวแยกข้อมูล OpenLineage แบบกำหนดเองสำหรับโอเปอเรเตอร์ Airflow ที่ไม่รองรับและสถานการณ์สายเลือดที่ซับซ้อน สองแนวทาง: เพิ่มเมธอด OpenLineage ลงในโอเปอเรเตอร์ที่คุณเป็นเจ้าของโดยตรง (แนะนำ) หรือสร้างตัวแยกข้อมูลแบบกำหนดเองสำหรับโอเปอเรเตอร์ของบุคคลที่สามที่คุณไม่สามารถแก้ไขได้ ตัวแยกข้อมูลจะสกัดกั้นการทำงานของโอเปอเรเตอร์ที่สามจุด: ก่อนการดำเนินการสำหรับสายเลือดแบบคงที่ หลังจากสำเร็จสำหรับเอาต์พุตที่กำหนดในรันไทม์ และหลังจากล้มเหลวสำหรับสายเลือดบางส่วน ลงทะเบียนตัวแยกข้อมูลผ่าน airflow.cfg หรือสภาพแวดล้อม...
debugging-dags
astronomer
การวิเคราะห์สาเหตุที่แท้จริงอย่างเป็นระบบและการแก้ไขสำหรับ Airflow DAGs ที่ล้มเหลว พร้อมขั้นตอนการตรวจสอบที่มีโครงสร้าง ชี้แนะผ่านกระบวนการวินิจฉัยสี่ขั้นตอน: ระบุความล้มเหลว ดึงรายละเอียดข้อผิดพลาด รวบรวมข้อมูลบริบท และส่งมอบขั้นตอนการแก้ไขที่สามารถดำเนินการได้ จัดหมวดหมู่ความล้มเหลวออกเป็นสี่ประเภท (ข้อมูล โค้ด โครงสร้างพื้นฐาน การพึ่งพา) เพื่อมุ่งเน้นการตรวจสอบและแนะนำการแก้ไขที่เหมาะสม ให้คำสั่ง CLI ที่พร้อมใช้งานสำหรับการดึงข้อมูลบันทึก การเปรียบเทียบการรัน การล้างงาน และ DAG...
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 DAGs และโปรเจกต์ ใช้เมื่อผู้ใช้ต้องการปรับใช้โค้ด, ส่ง DAGs, ตั้งค่า CI/CD, ปรับใช้สู่ระบบผลิต หรือสอบถามเกี่ยวกับกลยุทธ์การปรับใช้…
deploying-go-sdk-bundles
astronomer
สร้าง, แพ็ก, และปรับใช้ชุดรวม Airflow Go SDK ที่คอมไพล์แล้ว เพื่อให้ ExecutableCoordinator สามารถรันได้ ใช้เมื่อผู้ใช้ต้องการคอมไพล์ชุดรวมงาน Go, ถาม…
testing-dags
astronomer
วงจรการทดสอบ-ดีบัก-แก้ไขแบบวนซ้ำสำหรับ Airflow DAGs พร้อมการวินิจฉัยข้อผิดพลาดอย่างครอบคลุม เริ่มต้นด้วย af runs trigger-wait <dag_id> เพื่อรัน DAG และรอให้เสร็จสมบูรณ์ ไม่จำเป็นต้องตรวจสอบก่อนเริ่มต้น เมื่อเกิดข้อผิดพลาด ให้ใช้ af runs diagnose เพื่อสรุปข้อผิดพลาดอย่างครอบคลุม และ af tasks logs เพื่อตรวจสอบรายละเอียดข้อผิดพลาดจากงานเฉพาะ รองรับการกำหนดค่าเอง การหมดเวลา และการลองใหม่ จัดการสถานการณ์สำเร็จ ล้มเหลว และหมดเวลาพร้อมการตีความผลลัพธ์ที่ชัดเจน มีการตรวจสอบความถูกต้องอย่างรวดเร็ว...
tracing-downstream-lineage
astronomer
ติดตามสายข้อมูลปลายน้ำเพื่อประเมินผลกระทบจากการเปลี่ยนแปลงก่อนปรับแก้ตารางหรือ DAG ระบุผู้บริโภคโดยตรงของตารางเป้าหมายหรือ DAG ผ่านการค้นหาในซอร์สโค้ด การขึ้นต่อกันของวิว และการเชื่อมต่อเครื่องมือ BI สร้างแผนผังการขึ้นต่อกันแบบสมบูรณ์ที่แสดงผลกระทบปลายน้ำทั้งหมด ตั้งแต่ตารางไปจนถึงแดชบอร์ดและโมเดล ML จัดหมวดหมู่การขึ้นต่อกันตามความสำคัญ (วิกฤต สูง ปานกลาง ต่ำ) เพื่อจัดลำดับความสำคัญในการสื่อสารกับผู้มีส่วนได้ส่วนเสียและการทดสอบ สร้างรายงานผลกระทบพร้อมการประเมินความเสี่ยง ผลกระทบที่ได้รับ...