Skip to content

Examples

Real-world examples of using Latent for evaluation workflows.

Basic Pipeline

A simple data processing pipeline:

import asyncio
from latent.prefect import flow, task, params, logger
from latent.mlflow import log_metric, log_metrics
import pandas as pd

@task("extract", input="raw_data", output="extracted_data")
async def extract_task(raw_data: pd.DataFrame) -> pd.DataFrame:
    logger.info(f"Extracting {len(raw_data)} rows")
    extracted = raw_data[["id", "value", "timestamp"]]
    log_metric("extracted_rows", len(extracted))
    return extracted

@task("transform", input="extracted_data", output="transformed_data")
async def transform_task(data: pd.DataFrame) -> pd.DataFrame:
    logger.info(f"Transforming data")
    data["value_normalized"] = (data["value"] - data["value"].mean()) / data["value"].std()
    log_metric("mean_value", data["value"].mean())
    return data

@task("load", input="transformed_data", output="final_results")
async def load_task(data: pd.DataFrame) -> dict:
    results = {
        "total_rows": len(data),
        "mean_normalized": float(data["value_normalized"].mean())
    }
    logger.info(f"Pipeline complete: {results}")
    log_metrics(results)
    return results

@flow("etl_pipeline")
async def etl_pipeline():
    extracted = await extract_task()
    transformed = await transform_task(extracted)
    results = await load_task(transformed)
    return results

if __name__ == "__main__":
    asyncio.run(etl_pipeline())

LLM Evaluation Pipeline

Evaluate LLM responses:

import asyncio
from typing import Annotated
from latent.prefect import flow, task, params, logger
from latent.mlflow import log_metric, log_metrics, log_param
from latent.agents import ReActAgent, Judge, ScoredModel, OrdinalScore

class ResponseQuality(ScoredModel):
    quality: Annotated[int, OrdinalScore(scale=(1, 2, 3, 4, 5), pass_threshold=3)]

@task("generate_responses", input="prompts", output="responses", cache=False)
async def generate_responses_task(prompts: list) -> list:
    agent = ReActAgent("responder", model=params.model_name)
    responses = []

    for i, prompt in enumerate(prompts):
        response = await agent.ask(prompt)
        responses.append(response)
        logger.info(f"Generated response {i+1}/{len(prompts)}")

    log_param("model", params.model_name)
    log_metric("num_responses", len(responses))
    return responses

@task("evaluate_responses", input="responses", output="evaluations")
async def evaluate_responses_task(responses: list) -> dict:
    judge = Judge("quality_judge", model=params.model_name, output_type=ResponseQuality)
    scores = []

    for response in responses:
        result = await judge.evaluate({"conversation": response})
        scores.append(result.quality)

    results = {
        "mean_score": sum(scores) / len(scores),
        "min_score": min(scores),
        "max_score": max(scores)
    }

    log_metrics(results)
    logger.info(f"Evaluation results: {results}")
    return results

@flow("llm_evaluation")
async def llm_evaluation_flow():
    responses = await generate_responses_task()
    evaluation = await evaluate_responses_task(responses)
    return evaluation

if __name__ == "__main__":
    asyncio.run(llm_evaluation_flow())

Parallel Processing

Process items in parallel with concurrency control:

import asyncio

@task("process_item")
async def process_item_task(item: dict) -> dict:
    logger.info(f"Processing item {item['id']}")
    result = expensive_computation(item)
    log_metric(f"item_{item['id']}_score", result["score"])
    return result

@task("aggregate_results", output="results")
async def aggregate_results_task(results: list) -> dict:
    total_score = sum(r["score"] for r in results)
    avg_score = total_score / len(results)

    log_metric("average_score", avg_score)
    return {"average_score": avg_score, "results": results}

@flow("parallel_processing", input="input_items")
async def parallel_processing_flow(input_items: list):
    # Process all items in parallel with concurrency limit; .map() returns a results list
    results = await process_item_task.map(input_items, concurrency=5)

    # Aggregate
    final_results = await aggregate_results_task(results)
    return final_results

With DataFrame input:

@task("process_row")
async def process_row_task(row: dict) -> dict:
    return {"id": row["id"], "processed": row["value"] * 2}

@flow("dataframe_processing", input="data")
async def dataframe_flow(data: pd.DataFrame):
    # DataFrame rows are automatically converted to dicts; .map() returns a results list
    results = await process_row_task.map(data, concurrency=3)
    return results

Multi-Stage Pipeline

Chain multiple flows:

import asyncio

# Flow 1: Data Collection
@flow("data_collection", output=["raw_conversations"])
async def data_collection_flow():
    conversations = collect_conversations()
    return conversations

# Flow 2: Processing (depends on Flow 1)
@flow(
    "data_processing",
    input=["data_collection.raw_conversations"],
    output=["processed_conversations"]
)
async def data_processing_flow(raw_conversations):
    processed = process(raw_conversations)
    return processed

# Flow 3: Evaluation (depends on Flow 2)
@flow(
    "evaluation",
    input=["data_processing.processed_conversations"],
    output=["evaluation_results"]
)
async def evaluation_flow(processed_conversations):
    results = evaluate(processed_conversations)
    return results

# Run all flows in sequence
async def main():
    await data_collection_flow()
    await data_processing_flow()
    await evaluation_flow()

if __name__ == "__main__":
    asyncio.run(main())

See Also