Dagster & Prefect
This feature is considered in a preview stage, and is under active development, and not considered ready for production use. You may encounter feature gaps, and the APIs may change. For more information, see the API lifecycle stages documentation.
The dagster-prefect library uses Dagster Pipes to launch work on Prefect from a Dagster asset. Dagster stays the control plane, handling scheduling, partitioning, lineage, and retries, while the work itself runs on Prefect's infrastructure. Existing deployments work as is, with no change to your Prefect code.
This is useful when a workflow already runs on Prefect and you would rather orchestrate it than rewrite it, and when a step needs the durability of a workflow engine but should still be a node in the asset graph.
Two ways to launch
| Client | Launches | Executed by |
|---|---|---|
PipesPrefectDeploymentClient | a deployment run | a worker on the deployment's work pool, a push work pool, or a Prefect Managed work pool |
PipesPrefectTaskClient | a background task run | a task worker (prefect task serve) |
Prefer deployments. A deployment runs with no change to the flow, and it is the only option that can run on Prefect-managed infrastructure. Background tasks have no environment channel, so the payload has to travel as a task argument.
Installation
- uv
- pip
uv add dagster-prefect
pip install dagster-prefect
dagster-prefect cannot be installed alongside dagster-airflow or dagster-airlift when those are pinned to Airflow 2.x. Prefect 3 requires sqlalchemy>=2, and Airflow 2.x requires sqlalchemy<2. Use separate environments for Prefect and Airflow 2.x integrations.
Launching a deployment
Launch an existing deployment by its flow-name/deployment-name. The flow doesn't need any change:
import dagster as dg
from dagster_prefect import PipesPrefectDeploymentClient, PrefectResource
@dg.asset(kinds={"prefect"})
def orders_summary(
context: dg.AssetExecutionContext, prefect_deployments: PipesPrefectDeploymentClient
):
return prefect_deployments.run(
context=context,
deployment="refresh-orders/production",
parameters={"as_of": "latest"},
).get_materialize_result()
defs = dg.Definitions(
assets=[orders_summary],
resources={
"prefect_deployments": PipesPrefectDeploymentClient(
prefect=PrefectResource(
api_url=dg.EnvVar("PREFECT_API_URL"),
api_key=dg.EnvVar("PREFECT_API_KEY"),
)
)
},
)
The Dagster step blocks until the flow run reaches a terminal state, and materializes the asset if it succeeded, with a Prefect Run URL linking to the run in Prefect. The flow run's logs show up in the Dagster step's compute logs.
Without a Pipes session in the flow, that's all Dagster gets: the flow can't report metadata or asset checks, and assets are only materialized once the whole flow run finishes. The step also logs a warning that no Pipes messages were received, which is expected in this case.
api_url is the Prefect API, for example http://127.0.0.1:4200/api for an open source server. Set ui_url as well on Prefect Cloud, whose UI is served from a different host than its API.
Reporting from the flow
Opening a Pipes session in the flow is optional. It lets the flow report metadata and asset checks back as it runs, for example one materialization per asset of a multi-asset as each one finishes:
from dagster_pipes import PipesPrefectLogsMessageWriter, open_dagster_pipes
from prefect import flow, task
@task(retries=2)
def extract(as_of: str) -> list[dict]: ...
@flow
def refresh_orders(as_of: str = "latest") -> None:
rows = extract(as_of)
with open_dagster_pipes(message_writer=PipesPrefectLogsMessageWriter()) as pipes:
pipes.report_asset_materialization(metadata={"rows": len(rows)})
Anything the flow reports lands on the materialization. The block is safe outside Dagster: run the flow on its own and open_dagster_pipes warns and returns a no-op context.
Partitioning
Dagster can tell the flow which slice of data to compute, instead of the flow working it out at runtime. Set partition_parameter to the name of the flow parameter that should receive the partition key:
@dg.asset(
partitions_def=dg.DailyPartitionsDefinition(start_date="2026-01-01"),
kinds={"prefect"},
)
def daily_report(
context: dg.AssetExecutionContext, prefect_deployments: PipesPrefectDeploymentClient
):
return prefect_deployments.run(
context=context,
deployment="daily-report/production",
partition_parameter="day",
).get_materialize_result()
Each partition becomes one Prefect flow run with day set to that partition's key, so filling in a range of dates is an ordinary Dagster backfill rather than a script that iterates flow runs and works out which ones are missing. Use partition_window_parameters=("start", "end") to pass a time-partitioned window instead; both are sent as ISO 8601 strings, since Prefect parameters must be JSON-serializable.
A Pipes-aware flow can also read pipes.partition_key, pipes.partition_key_range, and pipes.partition_time_window off the Pipes context without any of this.
Multi-dimensional partitions are not supported by partition_parameter, because their keys are composite. Pass the dimensions you want as ordinary parameters.
Launching a background task
A background task has no environment channel, so it must accept the Pipes payload as a dagster_pipes_params argument. Opening a Pipes session with it is optional, to report metadata and asset checks back:
from dagster_pipes import (
PipesMappingParamsLoader,
PipesPrefectLogsMessageWriter,
open_dagster_pipes,
)
from prefect import task
@task
def summarize(as_of: str, dagster_pipes_params: dict[str, str] | None = None) -> None:
with open_dagster_pipes(
params_loader=PipesMappingParamsLoader(dagster_pipes_params or {}),
message_writer=PipesPrefectLogsMessageWriter(),
) as pipes:
pipes.report_asset_materialization(metadata={"rows": 100})
@dg.asset(kinds={"prefect"})
def orders_summary(context: dg.AssetExecutionContext, prefect_tasks: PipesPrefectTaskClient):
return prefect_tasks.run(
context=context, task=summarize, parameters={"as_of": "latest"}
).get_materialize_result()
A task worker must be serving the task, otherwise the task run is created and never picked up.
Cancellation
Terminating the Dagster run cancels the Prefect flow run it was waiting on. Set forward_termination=False on the client to leave the Prefect run alone.
Background tasks are the exception: Prefect's task worker runs a task to completion regardless of a cancellation request, so terminating the Dagster run logs a warning and the task keeps running.
A Prefect run cancelled from Prefect's side fails the Dagster step with a message naming the cancellation.
Reporting messages back
Both clients read Pipes messages from the Prefect run's logs, through the same Prefect API they use to launch and poll the run. No shared filesystem, bucket, or extra credentials are needed, so this works for a worker in a container, on another machine, or on Prefect-managed infrastructure.
For messages to arrive:
- The flow or task opens a Pipes session with
message_writer=PipesPrefectLogsMessageWriter(). It needsdagster-pipes1.13.25 or later in the environment the flow or task runs in. - Prefect sends logs to its API, which it does by default. Setting
PREFECT_LOGGING_TO_API_ENABLED=falseturns that off. PREFECT_LOGGING_LEVELisINFOor lower, since messages are logged atINFO.
The messages show up in the Prefect run's logs as JSON lines, under the prefect.flow_runs.dagster_pipes logger. If none arrive, as with a flow that doesn't open a Pipes session, the Dagster step still materializes the asset when the Prefect run succeeds and logs a warning listing what to check.
Large payloads
A single message can be at most about 900 KB, below Prefect's 1 MB limit on a log. A larger one is replaced by an error in the Dagster step's logs rather than delivered. Each message also counts against Prefect Cloud's log rate limit. For large metadata or heavy reporting, use a blob store instead, with the matching reader and writer on each side:
import boto3
from dagster_aws.pipes import PipesS3MessageReader
PipesPrefectDeploymentClient(
prefect=PrefectResource(api_url=dg.EnvVar("PREFECT_API_URL")),
message_reader=PipesS3MessageReader(bucket="my-bucket", client=boto3.client("s3")),
)
import boto3
from dagster_pipes import PipesS3MessageWriter, open_dagster_pipes
with open_dagster_pipes(message_writer=PipesS3MessageWriter(client=boto3.client("s3"))) as pipes:
...