latent.prefect.decorators¶
Prefect flow and task decorators with automagic configuration.
Classes¶
LatentTask¶
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 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 the task for async execution.