Aller au contenu principal

Décorateurs de fonctions

Le décorateur @mlflow.trace vous permet de créer un span pour n'importe quelle fonction. Les décorateurs de fonctions offrent le moyen le plus simple d'ajouter le traçage avec un minimum de modifications de code.

  • MLflow détecte les relations parent-enfant entre les fonctions, le rendant compatible avec les intégrations de traçage automatique.
  • Capture les exceptions lors de l'exécution de fonctions et les enregistre comme événements d'étendue.
  • Logs automatiquement le nom de la fonction, les entrées, les sorties et le temps d'exécution.
  • Peut être utilisé avec les fonctionnalités de suivi automatique.

Prérequis

Cette page nécessite les packages suivants :

  • mlflow[databricks] 3.1 et supérieur : Fonctionnalités MLflow de base avec des fonctionnalités GenAI et la connectivité Databricks.
  • openai 1.0.0 et supérieur : (Facultatif) Uniquement si votre code personnalisé interagit avec OpenAI ; remplacez-le par d'autres SDK si nécessaire.

Installer les exigences de base :

Python
%pip install --upgrade "mlflow[databricks]>=3.1"
# %pip install --upgrade "openai>=1.0.0" # Install if needed

Exemple de base

Le code suivant est un exemple minimal d'utilisation du décorateur pour le traçage des fonctions Python.

astuce

Ordre du décorateur

Pour assurer une observabilité complète, le décorateur @mlflow.trace doit généralement être le plus externe si vous utilisez plusieurs décorateurs. Consultez Utilisation de @mlflow.trace avec d’autres décorateurs pour une explication détaillée et des exemples.

Python
import mlflow


@mlflow.trace(span_type="func", attributes={"key": "value"})
def add_1(x):
return x + 1


@mlflow.trace(span_type="func", attributes={"key1": "value1"})
def minus_1(x):
return x - 1


@mlflow.trace(name="Trace Test")
def trace_test(x):
step1 = add_1(x)
return minus_1(step1)


trace_test(4)

Décorateur de traçage

remarque

Lorsqu'une trace contient plusieurs spans avec le même nom, MLflow leur ajoute un suffixe à incrémentation automatique, tel que _1, _2.

Personnaliser les portées

Le décorateur @mlflow.trace accepte les arguments suivants pour personnaliser la portée à créer :

  • name parameter pour remplacer le nom de l'étendue par le default (le nom de la fonction décorée)
  • span_type parameter pour définir le type d'étendue. Définissez l'un des Types d'étendue intégrés ou une chaîne.
  • attributes paramètre pour ajouter des attributs personnalisés à l'étendue.
astuce

Ordre du décorateur

Lors de la combinaison de @mlflow.trace avec d'autres décorateurs (par exemple, issus de frameworks web), il est crucial qu'il soit le plus extérieur. Pour un exemple clair d'ordre correct par rapport à un ordre incorrect, veuillez vous référer à Utilisation de @mlflow.trace avec d'autres décorateurs.

Python
@mlflow.trace(
name="call-local-llm", span_type=SpanType.LLM, attributes={"model": "gpt-4o-mini"}
)
def invoke(prompt: str):
return client.invoke(
messages=[{"role": "user", "content": prompt}], model="gpt-4o-mini"
)

Vous pouvez également mettre à jour un intervalle actif ou en direct dynamiquement à l'intérieur de la fonction en utilisant l'API mlflow.get_current_active_span.

Python
@mlflow.trace(span_type=SpanType.LLM)
def invoke(prompt: str):
model_id = "gpt-4o-mini"
# Get the current span (created by the @mlflow.trace decorator)
span = mlflow.get_current_active_span()
# Set the attribute to the span
span.set_attributes({"model": model_id})
return client.invoke(messages=[{"role": "user", "content": prompt}], model=model_id)

Pour plus d’exemples de modification LiveSpan d’objets, consultez le traçage d'étendue avec les gestionnaires de contexte.

Utilisez @mlflow.trace avec d'autres décorateurs

Lorsque vous appliquez plusieurs décorateurs à une seule fonction, il est crucial de placer @mlflow.trace comme décorateur **le plus externe** (celui qui est tout en haut). Cela garantit que MLflow peut capturer l'exécution complète de la fonction, y compris le comportement de tout décorateur interne.

Si @mlflow.trace n'est pas le décorateur le plus externe, sa visibilité sur l'exécution de la fonction peut être limitée ou incorrecte, ce qui peut entraîner des traces incomplètes ou une mauvaise représentation des entrées, des sorties et du temps d'exécution de la fonction.

Considérez l'exemple conceptuel suivant :

Python
import mlflow
import functools
import time

# A hypothetical additional decorator
def simple_timing_decorator(func):
@functools.wraps(func)
def wrapper(*args, **kwargs):
start_time = time.time()
result = func(*args, **kwargs)
end_time = time.time()
print(f"{func.__name__} executed in {end_time - start_time:.4f} seconds by simple_timing_decorator.")
return result
return wrapper

# Correct order: @mlflow.trace is outermost
@mlflow.trace(name="my_decorated_function_correct_order")
@simple_timing_decorator
# @another_framework_decorator # e.g., @app.route("/mypath") from Flask
def my_complex_function(x, y):
# Function logic here
time.sleep(0.1) # Simulate work
return x + y

# Incorrect order: @mlflow.trace is NOT outermost
@simple_timing_decorator
@mlflow.trace(name="my_decorated_function_incorrect_order")
# @another_framework_decorator
def my_other_complex_function(x, y):
time.sleep(0.1)
return x * y

# Example calls
if __name__ == "__main__":
print("Calling function with correct decorator order:")
my_complex_function(5, 3)

print("\nCalling function with incorrect decorator order:")
my_other_complex_function(5, 3)

Dans l'exemple my_complex_function (ordre correct), @mlflow.trace capturera l'exécution complète, y compris le temps ajouté par simple_timing_decorator. Dans my_other_complex_function (ordre incorrect), la trace capturée par MLflow pourrait ne pas refléter avec précision le temps d'exécution total ou pourrait manquer des modifications aux entrées/sorties effectuées par simple_timing_decorator avant que @mlflow.trace ne les voie.

Ajoutez des tags de trace

Des balises peuvent être ajoutées aux traces pour fournir des métadonnées supplémentaires au niveau de la trace. Il existe plusieurs façons de définir des balises sur une trace. Veuillez vous référer au guide d'ajout de tags personnalisés pour les autres méthodes.

Python
@mlflow.trace
def my_func(x):
mlflow.update_current_trace(tags={"fruit": "apple"})
return x + 1

Personnaliser les aperçus de requêtes et de réponses dans l'interface utilisateur

La tab Traces dans l'interface utilisateur MLflow affiche une liste de traces, et les colonnes Request et Response affichent un aperçu de l'entrée et de la sortie de bout en bout de chaque trace. Cela vous permet de comprendre rapidement ce que chaque trace représente.

default, ces aperçus sont tronqués à un nombre fixe de caractères. Cependant, vous pouvez personnaliser ce qui est affiché dans ces colonnes en utilisant les paramètres request_preview et response_preview au sein de la fonction mlflow.update_current_trace(). Ceci est particulièrement utile pour les entrées ou sorties complexes où la troncature par default pourrait ne pas afficher les informations les plus pertinentes.

Vous trouverez ci-dessous un exemple de définition d'un aperçu de requête personnalisé pour une trace qui traite un document long et les instructions de l'utilisateur, visant à afficher les informations les plus pertinentes dans la colonne Request de l'interface utilisateur :

Python
import mlflow

@mlflow.trace(name="Summarization Pipeline")
def summarize_document(document_content: str, user_instructions: str):
# Construct a custom preview for the request column
# For example, show beginning of document and user instructions
request_p = f"Doc: {document_content[:30]}... Instr: {user_instructions[:30]}..."
mlflow.update_current_trace(request_preview=request_p)

# Simulate LLM call
# messages = [
# {"role": "system", "content": "Summarize the following document based on user instructions."},
# {"role": "user", "content": f"Document: {document_content}\nInstructions: {user_instructions}"}
# ]
# completion = client.chat.completions.create(model="gpt-4o-mini", messages=messages)
# summary = completion.choices[0].message.content
summary = f"Summary of document starting with '{document_content[:20]}...' based on '{user_instructions}'"

# Customize the response preview
response_p = f"Summary: {summary[:50]}..."
mlflow.update_current_trace(response_preview=response_p)

return summary

# Example Call
long_document = "This is a very long document that contains many details about various topics..." * 10
instructions = "Focus on the key takeaways regarding topic X."
summary_result = summarize_document(long_document, instructions)
# print(summary_result)

En définissant request_preview et response_preview sur la trace (généralement le span racine), vous contrôlez la manière dont l'interaction globale est résumée dans la vue de la liste principale des traces, ce qui facilite l'identification et la compréhension des traces en un coup d'œil.

Gestion automatique des exceptions

Si une Exception est levée pendant le traitement d'une opération instrumentée par trace, une indication sera affichée dans l'interface utilisateur que l'invocation n'a pas été fructueuse et une capture partielle des données sera disponible pour faciliter le debugging. De plus, les détails concernant l'Exception qui a été levée seront inclus dans Events de l'étendue partiellement terminée, ce qui aidera davantage à identifier les problèmes survenant dans votre code.

Erreur de trace

Combiner avec le traçage automatique

Le traçage manuel s'intègre en toute transparence aux capacités de traçage automatique de MLflow. Voir Combiner le traçage manuel et automatique.

Traçage des workflows complexes

Pour les workflows complexes comportant plusieurs étapes, utilisez des étendues imbriquées pour capturer le flux d'exécution détaillé :

Python
@mlflow.trace(name="data_pipeline")
def process_data_pipeline(data_source: str):
# Extract phase
with mlflow.start_span(name="extract") as extract_span:
raw_data = extract_from_source(data_source)
extract_span.set_outputs({"record_count": len(raw_data)})

# Transform phase
with mlflow.start_span(name="transform") as transform_span:
transformed = apply_transformations(raw_data)
transform_span.set_outputs({"transformed_count": len(transformed)})

# Load phase
with mlflow.start_span(name="load") as load_span:
result = load_to_destination(transformed)
load_span.set_outputs({"status": "success"})

return result

Multi-threading

MLflow Tracing est compatible avec les threads, les traces sont isolées par default par thread. Mais vous pouvez également créer une trace qui s'étend sur plusieurs threads avec quelques étapes supplémentaires.

MLflow utilise le mécanisme ContextVar intégré de Python pour assurer la sécurité des threads, qui n'est pas propagée entre les threads par default. Par conséquent, vous devez copier manuellement le contexte du thread principal vers le thread worker, comme indiqué dans l'exemple ci-dessous.

Python
import contextvars
from concurrent.futures import ThreadPoolExecutor, as_completed
import mlflow
from mlflow.entities import SpanType
import openai

client = openai.OpenAI()

# Enable MLflow Tracing for OpenAI
mlflow.openai.autolog()


@mlflow.trace
def worker(question: str) -> str:
messages = [
{"role": "system", "content": "You are a helpful assistant."},
{"role": "user", "content": question},
]
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=messages,
temperature=0.1,
max_tokens=100,
)
return response.choices[0].message.content


@mlflow.trace
def main(questions: list[str]) -> list[str]:
results = []
# Almost same as how you would use ThreadPoolExecutor, but two additional steps
# 1. Copy the context in the main thread using copy_context()
# 2. Use ctx.run() to run the worker in the copied context
with ThreadPoolExecutor(max_workers=2) as executor:
futures = []
for question in questions:
ctx = contextvars.copy_context()
futures.append(executor.submit(ctx.run, worker, question))
for future in as_completed(futures):
results.append(future.result())
return results


questions = [
"What is the capital of France?",
"What is the capital of Germany?",
]

main(questions)

Traçage multi-threadé

astuce

En revanche, ContextVar est copié vers des tâches **asynchrones** by default. Par conséquent, vous n'avez pas besoin de copier manuellement le contexte lors de l'utilisation de asyncio, ce qui pourrait être un moyen plus simple de gérer les tâches concurrentes liées aux E/S en Python avec MLflow Tracing.

Sorties de streaming

Le décorateur @mlflow.trace peut être utilisé pour tracer les fonctions qui renvoient un générateur ou un itérateur, depuis MLflow 2.20.2.

Python
@mlflow.trace
def stream_data():
for i in range(5):
yield i

L'exemple ci-dessus générera une trace avec un seul span pour la fonction stream_data. Par default, MLflow capturera tous les éléments générés par le générateur sous forme de liste dans la sortie du span. Dans l'exemple ci-dessus, la sortie de la portée sera [0, 1, 2, 3, 4].

remarque

Une étendue pour une fonction de Stream débutera lorsque l'itérateur renvoyé commencera à être **consommé**, et se terminera lorsque l'itérateur sera épuisé, ou qu'une exception sera levée pendant l'itération.

Utilisation des réducteurs de sortie

Si vous souhaitez agréger les éléments pour obtenir une sortie à étendue unique, vous pouvez utiliser le paramètre output_reducer pour spécifier une fonction personnalisée afin d'agréger les éléments. La fonction personnalisée doit prendre une liste d'éléments générés comme entrées.

Python
from typing import List, Any

@mlflow.trace(output_reducer=lambda x: ",".join(x))
def stream_data():
for c in "hello":
yield c

Dans l'exemple ci-dessus, la sortie de l'étendue sera "h,e,l,l,o". Les segments bruts peuvent toujours être trouvés dans la Events tab du span dans l'interface utilisateur de MLflow Trace, vous permettant d'inspecter les valeurs individuelles générées lors du debugging.

Modèles courants de réducteurs de sortie

Voici quelques modèles courants pour l'implémentation de réducteurs de sortie.

Agrégation de jetons

Python
from typing import List, Dict, Any

def aggregate_tokens(chunks: List[str]) -> str:
"""Concatenate streaming tokens into complete text"""
return "".join(chunks)

@mlflow.trace(output_reducer=aggregate_tokens)
def stream_text():
for word in ["Hello", " ", "World", "!"]:
yield word

Agrégation des métriques

Python
def aggregate_metrics(chunks: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Aggregate streaming metrics into summary statistics"""
values = [c["value"] for c in chunks if "value" in c]
return {
"count": len(values),
"sum": sum(values),
"average": sum(values) / len(values) if values else 0,
"max": max(values) if values else None,
"min": min(values) if values else None
}

@mlflow.trace(output_reducer=aggregate_metrics)
def stream_metrics():
for i in range(10):
yield {"value": i * 2, "timestamp": time.time()}

Collecte d'erreurs

Python
def collect_results_and_errors(chunks: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Separate successful results from errors"""
results = []
errors = []

for chunk in chunks:
if chunk.get("error"):
errors.append(chunk["error"])
else:
results.append(chunk.get("data"))

return {
"results": results,
"errors": errors,
"success_rate": len(results) / len(chunks) if chunks else 0,
"has_errors": len(errors) > 0
}

Exemple avancé : streaming OpenAI

Voici un exemple avancé qui utilise le output_reducer pour consolider la sortie ChatCompletionChunk d'un LLM OpenAI en un seul objet message.

astuce

Nous vous recommandons d'utiliser l'auto-traçage pour OpenAI pour les cas d'utilisation en production, qui gère cela automatiquement. L'exemple ci-dessous est présenté à des fins de démonstration.

Python
import mlflow
import openai
from openai.types.chat import *
from typing import Optional


def aggregate_chunks(outputs: list[ChatCompletionChunk]) -> Optional[ChatCompletion]:
"""Consolidate ChatCompletionChunks to a single ChatCompletion"""
if not outputs:
return None

first_chunk = outputs[0]
delta = first_chunk.choices[0].delta
message = ChatCompletionMessage(
role=delta.role, content=delta.content, tool_calls=delta.tool_calls or []
)
finish_reason = first_chunk.choices[0].finish_reason
for chunk in outputs[1:]:
delta = chunk.choices[0].delta
message.content += delta.content or ""
message.tool_calls += delta.tool_calls or []
finish_reason = finish_reason or chunk.choices[0].finish_reason

base = ChatCompletion(
id=first_chunk.id,
choices=[Choice(index=0, message=message, finish_reason=finish_reason)],
created=first_chunk.created,
model=first_chunk.model,
object="chat.completion",
)
return base


@mlflow.trace(output_reducer=aggregate_chunks)
def predict(messages: list[dict]):
client = openai.OpenAI()
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=messages,
stream=True,
)
for chunk in stream:
yield chunk


for chunk in predict([{"role": "user", "content": "Hello"}]):
print(chunk)

Dans l'exemple ci-dessus, l'étendue predict générée aura un seul message de complétion de chat comme sortie, qui est agrégé par la fonction de réduction personnalisée.

Cas d’usage réels

Voici des exemples supplémentaires de réducteurs de sortie pour les scénarios GenAI courants :

Réponse LLM avec analyse JSON

Python
from typing import List, Dict, Any
import json

def parse_json_from_llm(content: str) -> str:
"""Extract and clean JSON from LLM responses that may include markdown"""
# Remove common markdown code block wrappers
if content.startswith("```json") and content.endswith("```"):
content = content[7:-3] # Remove ```json prefix and ``` suffix
elif content.startswith("```") and content.endswith("```"):
content = content[3:-3] # Remove generic ``` wrappers
return content.strip()

def json_stream_reducer(chunks: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Aggregate LLM streaming output and parse JSON response"""
full_content = ""
metadata = {}
errors = []

# Process different chunk types
for chunk in chunks:
chunk_type = chunk.get("type", "content")

if chunk_type == "content" or chunk_type == "token":
full_content += chunk.get("content", "")
elif chunk_type == "metadata":
metadata.update(chunk.get("data", {}))
elif chunk_type == "error":
errors.append(chunk.get("error"))

# Return early if errors occurred
if errors:
return {
"status": "error",
"errors": errors,
"raw_content": full_content,
**metadata
}

# Try to parse accumulated content as JSON
try:
cleaned_content = parse_json_from_llm(full_content)
parsed_data = json.loads(cleaned_content)

return {
"status": "success",
"data": parsed_data,
"raw_content": full_content,
**metadata
}
except json.JSONDecodeError as e:
return {
"status": "parse_error",
"error": f"Failed to parse JSON: {str(e)}",
"raw_content": full_content,
**metadata
}

@mlflow.trace(output_reducer=json_stream_reducer)
def generate_structured_output(prompt: str, schema: dict):
"""Generate structured JSON output from an LLM"""
# Simulate streaming JSON generation
yield {"type": "content", "content": '{"name": "John", '}
yield {"type": "content", "content": '"email": "john@example.com", '}
yield {"type": "content", "content": '"age": 30}'}

# Add metadata
trace_id = mlflow.get_current_active_span().request_id if mlflow.get_current_active_span() else None
yield {"type": "metadata", "data": {"trace_id": trace_id, "model": "gpt-4"}}

Génération de sortie structurée avec OpenAI

Voici un exemple complet d'utilisation de réducteurs de sortie avec OpenAI pour générer et analyser des réponses JSON structurées :

Python
import json
import mlflow
import openai
from typing import List, Dict, Any, Optional

def structured_output_reducer(chunks: List[Dict[str, Any]]) -> Dict[str, Any]:
"""
Aggregate streaming chunks into structured output with comprehensive error handling.
Handles token streaming, metadata collection, and JSON parsing.
"""
content_parts = []
trace_id = None
model_info = None
errors = []

for chunk in chunks:
chunk_type = chunk.get("type", "token")

if chunk_type == "token":
content_parts.append(chunk.get("content", ""))
elif chunk_type == "trace_info":
trace_id = chunk.get("trace_id")
model_info = chunk.get("model")
elif chunk_type == "error":
errors.append(chunk.get("message"))

# Join all content parts
full_content = "".join(content_parts)

# Base response
response = {
"trace_id": trace_id,
"model": model_info,
"raw_content": full_content
}

# Handle errors
if errors:
response["status"] = "error"
response["errors"] = errors
return response

# Try to extract and parse JSON
try:
# Clean markdown wrappers if present
json_content = full_content.strip()
if json_content.startswith("```json") and json_content.endswith("```"):
json_content = json_content[7:-3].strip()
elif json_content.startswith("```") and json_content.endswith("```"):
json_content = json_content[3:-3].strip()

parsed_data = json.loads(json_content)
response["status"] = "success"
response["data"] = parsed_data

except json.JSONDecodeError as e:
response["status"] = "parse_error"
response["error"] = f"JSON parsing failed: {str(e)}"
response["error_position"] = e.pos if hasattr(e, 'pos') else None

return response

@mlflow.trace(output_reducer=structured_output_reducer)
async def generate_customer_email(
customer_name: str,
issue: str,
sentiment: str = "professional"
) -> None:
"""
Generate a structured customer service email response.
Demonstrates real-world streaming with OpenAI and structured output parsing.
"""
client = openai.AsyncOpenAI()

system_prompt = """You are a customer service assistant. Generate a professional email response in JSON format:
{
"subject": "email subject line",
"greeting": "personalized greeting",
"body": "main email content addressing the issue",
"closing": "professional closing",
"priority": "high|medium|low"
}"""

user_prompt = f"Customer: {customer_name}\nIssue: {issue}\nTone: {sentiment}"

try:
# Stream the response
stream = await client.chat.completions.create(
model="gpt-4o-mini",
messages=[
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt}
],
stream=True,
temperature=0.7
)

# Yield streaming tokens
async for chunk in stream:
if chunk.choices[0].delta.content:
yield {
"type": "token",
"content": chunk.choices[0].delta.content
}

# Add trace metadata
if current_span := mlflow.get_current_active_span():
yield {
"type": "trace_info",
"trace_id": current_span.request_id,
"model": "gpt-4o-mini"
}

except Exception as e:
yield {
"type": "error",
"message": f"OpenAI API error: {str(e)}"
}

# Example usage
async def main():
# This will automatically aggregate the streamed output into structured JSON
async for chunk in generate_customer_email(
customer_name="John Doe",
issue="Product arrived damaged",
sentiment="empathetic"
):
# In practice, you might send these chunks to a frontend
print(chunk.get("content", ""), end="", flush=True)
remarque

Avantages de l'intégration

Cet exemple illustre plusieurs modèles concrets :

  • Mises à jour de l'interface utilisateur en streaming : Les jetons peuvent être affichés au fur et à mesure de leur arrivée
  • Validation de la sortie structurée : l'analyse JSON garantit le format de la réponse
  • **Résilience aux erreurs** : Gère avec élégance les erreurs d'API et les échecs d'analyse
  • Corrélation des traces : Lie la sortie streaming aux traces MLflow pour le debugging.

Agrégation de réponses multi-modèles

Python
def multi_model_reducer(chunks: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Aggregate responses from multiple models"""
responses = {}
latencies = {}

for chunk in chunks:
model = chunk.get("model")
if model:
responses[model] = chunk.get("response", "")
latencies[model] = chunk.get("latency", 0)

return {
"responses": responses,
"latencies": latencies,
"fastest_model": min(latencies, key=latencies.get) if latencies else None,
"consensus": len(set(responses.values())) == 1
}

Test des réducteurs de sortie

Les réducteurs de sortie peuvent être testés indépendamment du framework de traçage, ce qui facilite la vérification de leur gestion correcte des cas limites :

Python
import unittest
from typing import List, Dict, Any

def my_reducer(chunks: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Example reducer to be tested"""
if not chunks:
return {"status": "empty", "total": 0}

total = sum(c.get("value", 0) for c in chunks)
errors = [c for c in chunks if c.get("error")]

return {
"status": "error" if errors else "success",
"total": total,
"count": len(chunks),
"average": total / len(chunks) if chunks else 0,
"error_count": len(errors)
}

class TestOutputReducer(unittest.TestCase):
def test_normal_case(self):
chunks = [
{"value": 10},
{"value": 20},
{"value": 30}
]
result = my_reducer(chunks)
self.assertEqual(result["status"], "success")
self.assertEqual(result["total"], 60)
self.assertEqual(result["average"], 20.0)

def test_empty_input(self):
result = my_reducer([])
self.assertEqual(result["status"], "empty")
self.assertEqual(result["total"], 0)

def test_error_handling(self):
chunks = [
{"value": 10},
{"error": "Network timeout"},
{"value": 20}
]
result = my_reducer(chunks)
self.assertEqual(result["status"], "error")
self.assertEqual(result["total"], 30)
self.assertEqual(result["error_count"], 1)

def test_missing_values(self):
chunks = [
{"value": 10},
{"metadata": "some info"}, # No value field
{"value": 20}
]
result = my_reducer(chunks)
self.assertEqual(result["total"], 30)
self.assertEqual(result["count"], 3)
astuce

Considérations sur les performances

  • Les réducteurs de sortie reçoivent tous les segments en mémoire en une seule fois. Pour les très grands streams, envisagez d’implémenter des alternatives de streaming ou des stratégies de découpage.
  • L'étendue reste ouverte jusqu'à ce que le générateur soit entièrement consommé, ce qui a un impact sur les métriques de latence.
  • Les réducteurs doivent être stateless et éviter les effets secondaires pour un comportement prévisible.

Types de fonctions pris en charge

Le décorateur @mlflow.trace prend actuellement en charge les types de fonctions suivants :

Type de fonction

Pris en charge

Synchroniser

Oui

Asynchrone

Oui (MLflow >= 2.16.0)

Générateur

Oui (MLflow >= 2.20.2)

Générateur asynchrone

Oui (MLflow >= 2.20.2)

Type de fonction

Pris en charge

Synchroniser

Oui

Asynchrone

Oui (MLflow >= 2.16.0)

Générateur

Oui (MLflow >= 2.20.2)

Générateur asynchrone

Oui (MLflow >= 2.20.2)

Ressources supplémentaires