Skip to content

latent.prefect.decorators

Prefect flow and task decorators with automagic configuration.

Classes

LatentTask

LatentTask(prefect_task, name: str, original_func: Callable)

Wrapper around Prefect task providing enhanced .map() functionality.

Supports: - Concurrency control via an asyncio.Semaphore bound - Automatic DataFrame row iteration (converts to list of dicts) - Automatic list/array handling

Example: @task("process_doc") async def process_doc(doc): return await transform(doc)

@flow("my_flow")
async def my_flow():
    docs = ["doc1.md", "doc2.md", "doc3.md"]
    # Map with a concurrency bound of 5; returns results in order
    results = await process_doc.map(docs, concurrency=5)

    # Works with DataFrames too (iterates over rows as dicts)
    df = pd.DataFrame({"text": ["a", "b"], "id": [1, 2]})
    results = await process_doc.map(df, concurrency=3)

Functions

latent_flow

latent_flow(flow_name: str, input: str | list[str] | None = None, output: str | list[str] | None = None, config_schema: type[T] | None = None, description: str | None = None, tags: list[str] | None = None, version: str | None = None, flow_kwargs = {}) -> Callable[[Callable[P, R]], Callable[P, R]]

Decorator that wraps @flow with automagic config loading, logging, and MLFlow setup.

Note: This decorator is now available as @latent.flow for a cleaner API. Both names work identically, but @latent.flow is the preferred syntax.

Args: flow_name: Name of the flow (must match directory name in flows/) input: Optional dataset name(s) to load and inject as function arguments output: Optional dataset name(s) to save return values to config_schema: Optional Pydantic model to validate parameters against description: Optional description for CLI discovery tags: Optional tags for CLI listing/filtering version: Optional version string **flow_kwargs: Additional arguments to pass to @flow decorator

Returns: Decorated flow function with configuration loaded

Example (new style): from latent.prefect import flow

class MyConfig(BaseModel):
    batch_size: int

@flow("my_flow", config_schema=MyConfig)
def my_flow():
    # Use params...
    pass

latent_task

latent_task(name: str, input: str | list[str] | None = None, output: str | list[str] | None = None, cache: bool = True, retry: bool = True, span_type: str | None = None, task_kwargs = {}) -> Callable[[Callable[P, R]], LatentTask]

Task decorator with automatic catalog loading, saving, and observability.

Note: This decorator is available as @latent.task for a cleaner API.

Automatic Features: - Caching keyed on inputs + task identity (disable with cache=False) - Retries on failure (disable with retry=False) - MLflow Span generation (auto-detected or explicit span_type) - Automatic input loading from catalog (from filesystem) - Automatic output saving to MLflow artifacts - Enhanced .map() with concurrency control

Caching: When cache=True (default), uses Prefect's native caching: - cache key = arguments + this task's identity (name, module.qualname of the decorated function, and its source). A caller-supplied cache_policy= is composed with the identity terms rather than replacing them, so it cannot put two distinct tasks back on one cache record; cache_policy=None and cache_policy=NO_CACHE still mean "do not cache". - catalog input= data is fingerprinted into the key, so an edited input file re-runs the task (_CatalogInputCachePolicy). A dataset that cannot be fingerprinted — a floating remote version, a missing or unreadable file — makes the run uncacheable rather than assumed unchanged. An input the caller passes in explicitly is left to the ordinary argument hashing, exactly as the wrapper leaves its catalog entry unread. - persist_result = True (Prefect's default result storage, ~/.prefect/storage/; point PREFECT_LOCAL_STORAGE_PATH elsewhere to move it) - cache_expiration = 7 days (default)

Args: name: Name of the task input: Dataset name(s) to load and inject as function arguments output: Dataset name(s) to save return values to (saved to MLflow artifacts) cache: Enable automatic caching (default: True) retry: Enable automatic retries on failure (default: True) span_type: Type of MLflow span to create (e.g. "AGENT", "JUDGE") **task_kwargs: Additional arguments to pass to @task decorator

Returns: LatentTask: A wrapped task with enhanced .map(concurrency=N) support

Example: @latent.task("process_data", span_type="TOOL") def process(data): ...

# Map with concurrency control
@latent.flow("my_flow")
def my_flow():
    items = [1, 2, 3, 4, 5]
    futures = process.map(items, concurrency=3)
    results = [f.result() for f in futures]

set_cli_kwargs

set_cli_kwargs(kwargs: dict[str, Any]) -> contextvars.Token

Set CLI kwargs for the next flow invocation. Called by the CLI runner.

Methods

LatentTask.map

map(items: Iterable[Any], args = (), concurrency: int | None = None, per_task_timeout_s: float | None = None, kwargs = {})

Map this task over an iterable, awaiting all results concurrently.

Each invocation runs as a coroutine on the calling event loop (await task(item)) and they are gathered together. concurrency bounds the fan-out with an asyncio.Semaphore; per_task_timeout_s wraps each call in asyncio.wait_for. Returns results (not futures) in input order.

Args: items: Iterable to map over (first argument). Can be: - List/array of items - pandas DataFrame (will iterate over rows as dicts) args: Additional positional arguments passed to each task invocation (unmapped) concurrency: Maximum number of concurrent task executions. If None, all are gathered without a bound. per_task_timeout_s: If set, each call is wrapped in asyncio.wait_for with this deadline. *kwargs: Additional keyword arguments passed to each task invocation (unmapped)

Returns: List of results in the same order as items.

Example: @task("add_value") def add_value(x: int, offset: int) -> int: return x + offset

@flow("my_flow")
async def my_flow():
    results = await add_value.map([1, 2, 3], offset=10, concurrency=5)
    # results == [11, 12, 13]

LatentTask.submit

submit(args = (), kwargs = {})

Submit the task for async execution.

Attributes

P

R