Run dbt CLI commands as Flyte tasks and capture node results and reports.
DBT
The dbt plugin lets you run dbt CLI invocations as Flyte tasks. A DbtTask maps one dbtRunner.invoke(...) call to one Flyte task, so each dbt command appears as its own node in the Flyte run graph.
Use this plugin when you want to orchestrate dbt commands alongside Python tasks, chain dbt invocations with other workflow steps, capture dbt node results as typed task outputs, and optionally render a dbt report in the Flyte UI.
Installation
pip install flyteplugins-dbtInstall the dbt adapter required by your project in the task image as well, such as dbt-duckdb, dbt-bigquery, or dbt-snowflake.
import flyte
image = (
flyte.Image.from_debian_base(python_version=(3, 12))
.with_pip_packages(
"flyteplugins-dbt",
"dbt-duckdb",
)
)Quick start
Put your dbt project and profiles directory in the task image, then create a DbtTask with project-level configuration. Pass the dbt command when you call the task.
from pathlib import Path
import flyte
from flyteplugins.dbt import DbtNodeResult, DbtTask
DBT_PROJECT_DIR = "jaffle_shop"
DBT_PROFILES_DIR = "dbt-profiles"
image = (
flyte.Image.from_debian_base(python_version=(3, 12))
.with_pip_packages("flyteplugins-dbt", "dbt-duckdb")
.with_source_folder(Path(DBT_PROJECT_DIR))
.with_source_folder(Path(DBT_PROFILES_DIR))
)
env = flyte.TaskEnvironment(name="dbt-jaffle-shop", image=image)
dbt_build = DbtTask(
name="dbt-build",
task_environment=env,
project_dir=DBT_PROJECT_DIR,
profiles_dir=DBT_PROFILES_DIR,
profile=DBT_PROJECT_DIR,
report=True,
)
@env.task
async def main() -> list[DbtNodeResult]:
return await dbt_build.aio(command="build")project_dir, profiles_dir, profile, and target_path are configured on the task because they describe the dbt project layout. The actual dbt invocation is a task input, so the same DbtTask can run different commands.
Commands and arguments
Pass command as either a shell-style string or an explicit list of tokens:
await dbt_build.aio(command="build")
await dbt_build.aio(command="docs generate")
await dbt_build.aio(command=["source", "freshness"])Use the typed task inputs for common dbt selectors and targets:
await dbt_build.aio(
command="test",
select=["stg_orders+"],
exclude=["tag:slow"],
target="prod",
)Use extra_args for dbt flags that are not modeled directly:
await dbt_build.aio(
command="run",
extra_args=["--full-refresh"],
)DbtTask manages these flags itself, so do not pass them in extra_args: --project-dir, --profiles-dir, --profile, --target, --target-path, --select, and --exclude.
Chaining dbt tasks
dbt does not support parallel programmatic invocations in the same Python process. When orchestrating multiple dbt commands from one Flyte task, run them sequentially:
@env.task
async def dbt_pipeline() -> list[DbtNodeResult]:
await dbt_build.aio(command="run")
return await dbt_build.aio(command="test")You can override task execution settings per call:
@env.task
async def dbt_pipeline():
await dbt_build.override(short_name="dbt-run").aio(command="run")
await dbt_build.override(short_name="dbt-test").aio(command="test")Outputs
DbtTask returns a list of DbtNodeResult values summarized from dbt node results.
from flyteplugins.dbt import DbtNodeResult
@env.task
async def count_failed_tests() -> int:
results: list[DbtNodeResult] = await dbt_build.aio(command="test")
return sum(result.status.lower() == "fail" for result in results)Each DbtNodeResult includes:
| Field | Description |
|---|---|
unique_id |
dbt unique node ID |
name |
dbt node name |
resource_type |
dbt resource type, such as model, test, or source |
status |
dbt node status |
message |
dbt message, when available |
failures |
Failure count, when dbt provides one |
execution_time |
Node execution time, when available |
relation_name |
Relation name, when dbt provides one |
Some dbt commands, such as commands that return plain strings instead of node result objects, may return an empty list. The dbt command output is still available in task logs.
Reports and failures
Set report=True to write a dbt report tab for the task. The report includes node names, resource types, statuses, failures, execution times, messages, and relation names.
dbt_test = DbtTask(
name="dbt-test",
task_environment=env,
project_dir=DBT_PROJECT_DIR,
profiles_dir=DBT_PROFILES_DIR,
profile=DBT_PROJECT_DIR,
report=True,
)If dbt reports failure, the task fails. If dbt provides an exception, the plugin raises that exception. Otherwise, it raises DbtInvocationError with a summary of the failing node results.
The report is written before the failure is raised, so a failed dbt task can still show which nodes succeeded and which nodes failed.
For example, a failing dbt test task can still show the full node result table in the report:
Custom dbt callbacks
Pass dbt event callbacks to DbtTask when you need custom logging or event handling.
from dbt_common.events.base_types import EventMsg
def print_version_callback(event: EventMsg):
if event.info.name == "MainReportVersion":
print(f"We are thrilled to be running dbt{event.data.version}")
dbt_build = DbtTask(
name="dbt-build",
task_environment=env,
project_dir=DBT_PROJECT_DIR,
profiles_dir=DBT_PROFILES_DIR,
callbacks=[print_version_callback],
)The callback receives dbt’s structured EventMsg, including event.info metadata and event-specific event.data fields. For remote execution, callback functions must be importable. If you define a callback in the same file as your DbtTask, the plugin serializes it by import path when the workflow is registered. Import-path strings are also supported:
dbt_build = DbtTask(
name="dbt-build",
task_environment=env,
project_dir=DBT_PROJECT_DIR,
callbacks=["my_project.callbacks.print_version_callback"],
)Reference
DbtTask
| Parameter | Default | Description |
|---|---|---|
name |
- | Flyte task name |
task_environment |
None |
Task environment that owns the task |
project_dir |
None |
Path to the dbt project directory |
profiles_dir |
None |
Path to the dbt profiles directory |
profile |
None |
dbt profile name |
target_path |
None |
dbt target path |
callbacks |
[] |
dbt event callbacks or import-path strings |
report |
False |
Write a dbt report tab in the Flyte UI |
DbtTask also accepts normal Flyte task settings such as resources, secrets, queue, interruptible, and short_name. Caching is not supported yet; leave cache unset or set it to "disable".
Invocation inputs
| Input | Default | Description |
|---|---|---|
command |
- | dbt command as a string or token list |
select |
None |
dbt selectors passed as --select |
exclude |
None |
dbt selectors passed as --exclude |
target |
None |
dbt target passed as --target |
extra_args |
None |
Additional dbt CLI arguments |