cosmos-dbt-fusion

tarafından astronomer

Astronomer Cosmos'u Snowflake, Databricks, BigQuery veya Redshift üzerindeki dbt Fusion projeleri için yerel yürütme ile yapılandırın. Cosmos 1.11.0+ gerektirir, dbt Fusion ikili dosyası Airflow çalışma zamanında ayrıca kurulmalıdır ve alt süreç çağrısı ile ExecutionMode.LOCAL kullanır. Üç ayrıştırma stratejisini destekler: dbt_manifest (büyük projeler için en hızlı), dbt_ls (karmaşık seçiciler için) veya otomatik (basit kurulumlar). Depo bağlantıları için ProfileConfig kurulumunu, dbt proje yolları için ProjectConfig'i ve... kapsar.

npx skills add https://github.com/astronomer/agents --skill cosmos-dbt-fusion

Cosmos + dbt Fusion: Implementation Checklist

Execute steps in order. This skill covers Fusion-specific constraints only.

Version note: dbt Fusion support was introduced in Cosmos 1.11.0. Requires Cosmos ≥1.11.

Reference: See reference/cosmos-config.md for ProfileConfig, operator_args, and Airflow 3 compatibility details.

Before starting, confirm: (1) dbt engine = Fusion (not Core → use cosmos-dbt-core), (2) warehouse = Snowflake, Databricks, Bigquery and Redshift only.

Fusion-Specific Constraints

ConstraintDetails
No asyncAIRFLOW_ASYNC not supported
No virtualenvFusion is a binary, not a Python package
Warehouse supportSnowflake, Databricks, Bigquery and Redshift support while in preview

1. Confirm Cosmos Version

CRITICAL: Cosmos 1.11.0 introduced dbt Fusion compatibility.

# Check installed version
pip show astronomer-cosmos

# Install/upgrade if needed
pip install "astronomer-cosmos>=1.11.0"

Validate: pip show astronomer-cosmos reports version ≥ 1.11.0


2. Install the dbt Fusion Binary (REQUIRED)

dbt Fusion is NOT bundled with Cosmos or dbt Core. Install it into the Airflow runtime/image.

Determine where to install the Fusion binary (Dockerfile / base image / runtime).

Example Dockerfile Install

USER root
RUN apt-get update && apt-get install -y curl
ENV SHELL=/bin/bash
RUN curl -fsSL https://public.cdn.getdbt.com/fs/install/install.sh | sh -s -- --update
USER astro

Common Install Paths

EnvironmentTypical path
Astro Runtime/home/astro/.local/bin/dbt
System-wide/usr/local/bin/dbt

Validate: The dbt binary exists at the chosen path and dbt --version succeeds.


3. Choose Parsing Strategy (RenderConfig)

Parsing strategy is the same as dbt Core. Pick ONE:

Load modeWhen to useRequired inputs
dbt_manifestLarge projects; fastest parsingProjectConfig.manifest_path
dbt_lsComplex selectors; need dbt-native selectionFusion binary accessible to scheduler
automaticSimple setups; let Cosmos pick(none)
from cosmos import RenderConfig, LoadMode

_render_config = RenderConfig(
    load_method=LoadMode.AUTOMATIC,  # or DBT_MANIFEST, DBT_LS
)

4. Configure Warehouse Connection (ProfileConfig)

Reference: See reference/cosmos-config.md for full ProfileConfig options and examples.

from cosmos import ProfileConfig
from cosmos.profiles import SnowflakeUserPasswordProfileMapping

_profile_config = ProfileConfig(
    profile_name="default",
    target_name="dev",
    profile_mapping=SnowflakeUserPasswordProfileMapping(
        conn_id="snowflake_default",
    ),
)

5. Configure ExecutionConfig (LOCAL Only)

CRITICAL: dbt Fusion with Cosmos requires ExecutionMode.LOCAL with dbt_executable_path pointing to the Fusion binary.

from cosmos import ExecutionConfig
from cosmos.constants import InvocationMode

_execution_config = ExecutionConfig(
    invocation_mode=InvocationMode.SUBPROCESS,
    dbt_executable_path="/home/astro/.local/bin/dbt",  # REQUIRED: path to Fusion binary
    # execution_mode is LOCAL by default - do not change
)

6. Configure Project (ProjectConfig)

from cosmos import ProjectConfig

_project_config = ProjectConfig(
    dbt_project_path="/path/to/dbt/project",
    # manifest_path="/path/to/manifest.json",  # for dbt_manifest load mode
    # install_dbt_deps=False,  # if deps precomputed in CI
)

7. Assemble DAG / TaskGroup

Option A: DbtDag (Standalone)

from cosmos import DbtDag, ProjectConfig, ProfileConfig, ExecutionConfig, RenderConfig
from cosmos.profiles import SnowflakeUserPasswordProfileMapping
from pendulum import datetime

_project_config = ProjectConfig(
    dbt_project_path="/usr/local/airflow/dbt/my_project",
)

_profile_config = ProfileConfig(
    profile_name="default",
    target_name="dev",
    profile_mapping=SnowflakeUserPasswordProfileMapping(
        conn_id="snowflake_default",
    ),
)

_execution_config = ExecutionConfig(
    dbt_executable_path="/home/astro/.local/bin/dbt",  # Fusion binary
)

_render_config = RenderConfig()

my_fusion_dag = DbtDag(
    dag_id="my_fusion_cosmos_dag",
    project_config=_project_config,
    profile_config=_profile_config,
    execution_config=_execution_config,
    render_config=_render_config,
    start_date=datetime(2025, 1, 1),
    schedule="@daily",
)

Option B: DbtTaskGroup (Inside Existing DAG)

from airflow.sdk import dag, task  # Airflow 3.x
# from airflow.decorators import dag, task  # Airflow 2.x
from airflow.models.baseoperator import chain
from cosmos import DbtTaskGroup, ProjectConfig, ProfileConfig, ExecutionConfig
from pendulum import datetime

_project_config = ProjectConfig(dbt_project_path="/usr/local/airflow/dbt/my_project")
_profile_config = ProfileConfig(profile_name="default", target_name="dev")
_execution_config = ExecutionConfig(dbt_executable_path="/home/astro/.local/bin/dbt")

@dag(start_date=datetime(2025, 1, 1), schedule="@daily")
def my_dag():
    @task
    def pre_dbt():
        return "some_value"

    dbt = DbtTaskGroup(
        group_id="dbt_fusion_project",
        project_config=_project_config,
        profile_config=_profile_config,
        execution_config=_execution_config,
    )

    @task
    def post_dbt():
        pass

    chain(pre_dbt(), dbt, post_dbt())

my_dag()

8. Final Validation

Before finalizing, verify:

  • Cosmos version: ≥1.11.0
  • Fusion binary installed: Path exists and is executable
  • Warehouse supported: Snowflake, Databricks, Bigquery or Redshift only
  • Secrets handling: Airflow connections or env vars, NOT plaintext

Troubleshooting

If user reports dbt Core regressions after enabling Fusion:

AIRFLOW__COSMOS__PRE_DBT_FUSION=1

User Must Test

  • The DAG parses in the Airflow UI (no import/parse-time errors)
  • A manual run succeeds against the target warehouse (at least one model)

Reference


Related Skills

  • cosmos-dbt-core: For dbt Core projects (not Fusion)
  • authoring-dags: General DAG authoring patterns
  • testing-dags: Testing DAGs after creation

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-downstream-lineage
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