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())