AI Engineering Academy · Lección

Streaming de salida en LangChain

Implementará streaming de tokens mediante chains de LCEL para que su aplicación muestre cada palabra a medida que llega, en lugar de esperar a la respuesta completa, mejorando la latencia percibida.

Lección 4 de 413 pasos

Streaming de salida en LangChain es una lección gratuita de AI Engineering Academy en CoddyKit. Esta es la lección 4 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de AI Engineering Academy, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de AI Engineering Academy incluye 4 lecciones en total.

Por qué es importante la transmisión de resultados

Sin transmisión de resultados, los usuarios observan una pantalla en blanco mientras esperan a que el LLM termine de generar la respuesta, lo que puede tardar entre 5 y 30 segundos en respuestas largas. Con la transmisión, los tokens aparecen a medida que se generan y proporcionan información inmediata. Esto mejora notablemente la capacidad de respuesta percibida. El LCEL de LangChain propaga automáticamente la transmisión por toda la cadena cuando se llama a .stream().

Transmisión básica con .stream()

Toda cadena de LCEL expone un método .stream() que devuelve un iterador de fragmentos. En una cadena que termina en StrOutputParser, cada fragmento es una parte de una cadena de texto. Puede iterar sobre los fragmentos e imprimirlos o generarlos a medida que llegan. La transmisión se produce en el nivel HTTP: cada token de la API de OpenAI se reenvía a través del analizador en cuanto llega.

from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser

chain = (
    ChatPromptTemplate.from_template('Explain {topic} in detail.')
    | ChatOpenAI(model='gpt-4o-mini')
    | StrOutputParser()
)

# Stream tokens to stdout
for chunk in chain.stream({'topic': 'quantum entanglement'}):
    print(chunk, end='', flush=True)
print()  # final newline

Transmisión asíncrona con .astream()

.astream() es la versión asíncrona de .stream(). Devuelve un iterador asíncrono que se consume con async for. Este es el enfoque correcto en FastAPI, Starlette y otros frameworks web asíncronos donde el controlador de solicitudes es una corrutina. Usar transmisión síncrona en un controlador asíncrono bloquearía el bucle de eventos.

import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate

chain = (
    ChatPromptTemplate.from_template('Write a poem about {subject}')
    | ChatOpenAI(model='gpt-4o-mini')
    | StrOutputParser()
)

async def stream_response():
    async for chunk in chain.astream({'subject': 'the ocean'}):
        print(chunk, end='', flush=True)

asyncio.run(stream_response())

Transmisión en FastAPI con StreamingResponse

En FastAPI, envuelva un generador asíncrono en StreamingResponse con media_type='text/plain' para transmitir tokens de texto al navegador. Para los eventos enviados por el servidor (SSE), use media_type='text/event-stream' y dé formato a cada fragmento como data: ...\n\n. De este modo, el navegador recibe los tokens a medida que se generan sin tener que esperar a que se complete la respuesta.

from fastapi import FastAPI
from fastapi.responses import StreamingResponse

app = FastAPI()

async def generate_stream(topic: str):
    async for chunk in chain.astream({'topic': topic}):
        yield chunk

@app.get('/stream')
async def stream_endpoint(topic: str):
    return StreamingResponse(
        generate_stream(topic),
        media_type='text/plain'
    )

# SSE format for frontend EventSource
async def sse_stream(topic: str):
    async for chunk in chain.astream({'topic': topic}):
        yield f'data: {chunk}\n\n'

Transmisión a través de pasos intermedios

Las cadenas de LCEL propagan la transmisión a través de cada paso que la admite. StrOutputParser reconoce la transmisión y pasa los fragmentos de inmediato. Sin embargo, algunos analizadores, como JsonOutputParser, deben almacenar toda la salida antes de analizarla, lo que interrumpe la transmisión. LangChain lo hace explícito: si un paso no es compatible con la transmisión, acumula la salida antes de pasarla al siguiente.

from langchain_core.output_parsers import JsonOutputParser

# This chain does NOT stream token by token
# JsonOutputParser must buffer the full response before parsing JSON
json_chain = (
    ChatPromptTemplate.from_template('Return JSON: {task}')
    | ChatOpenAI(model='gpt-4o-mini')
    | JsonOutputParser()  # buffers until complete
)

# But partial JSON streaming IS possible with streaming_json_parser
for partial in json_chain.stream({'task': 'list 3 colors'}):
    print(partial)  # prints partial dict as it fills in

astream_events para un control detallado

.astream_events() proporciona una API de transmisión más granular que emite eventos para cada paso de la cadena, no solo para la salida final. Cada evento tiene un campo kind (on_chain_start, on_llm_stream, on_chain_end) y una carga data. Esto permite transmitir por separado los resultados de las llamadas a herramientas, el razonamiento intermedio y la salida final a distintas partes de una interfaz de usuario.

async def stream_with_events(question: str):
    async for event in chain.astream_events(
        {'question': question},
        version='v2'
    ):
        kind = event['event']
        if kind == 'on_llm_stream':
            chunk = event['data']['chunk'].content
            print(chunk, end='', flush=True)
        elif kind == 'on_chain_end':
            print('\n[Done]')
        elif kind == 'on_tool_start':
            print(f'\n[Tool: {event["name"]}]')

Almacenamiento en búfer de la salida transmitida

A veces necesita transmitir tokens al usuario y capturar la respuesta completa para registrarla o procesarla posteriormente. Use .astream() con un acumulador basado en una lista. Una los fragmentos después del bucle para obtener el texto completo. Este patrón permite mostrar la salida en tiempo real y, al mismo tiempo, guardar la respuesta completa para análisis, almacenamiento en caché o evaluación.

async def stream_and_capture(question: str) -> str:
    full_response = []
    async for chunk in chain.astream({'question': question}):
        print(chunk, end='', flush=True)  # stream to user
        full_response.append(chunk)        # also collect
    print()  # newline
    complete = ''.join(full_response)
    await log_response(question, complete)  # log full text
    return complete

Transmisión con llamadas a herramientas

Cuando un modelo genera una llamada a una herramienta en una respuesta transmitida, los argumentos de la función llegan como fragmentos de tokens. Debe almacenar en búfer la cadena de argumentos JSON hasta que se complete la llamada a la herramienta antes de ejecutarla. LangChain gestiona esto automáticamente en sus ejecutores de agentes, pero si está creando un bucle de transmisión personalizado, debe comprobar finish_reason y acumular los fragmentos de tool_call.function.arguments.

from openai import AsyncOpenAI

client = AsyncOpenAI()

async def stream_with_tools(prompt: str):
    tool_call_buffer = {}
    async with client.chat.completions.stream(
        model='gpt-4o-mini',
        messages=[{'role': 'user', 'content': prompt}],
        tools=[weather_tool_schema]
    ) as stream:
        async for chunk in stream:
            delta = chunk.choices[0].delta
            if delta.tool_calls:
                for tc in delta.tool_calls:
                    idx = tc.index
                    if idx not in tool_call_buffer:
                        tool_call_buffer[idx] = ''
                    if tc.function.arguments:
                        tool_call_buffer[idx] += tc.function.arguments

Cancelación y tiempo de espera con transmisión

Las respuestas transmitidas de larga duración necesitan compatibilidad con la cancelación. En Python asíncrono, puede cancelar un asyncio.Task que envuelva la transmisión. En FastAPI, el framework gestiona automáticamente la cancelación cuando el cliente se desconecta si se utiliza StreamingResponse. Establezca un tiempo de espera mediante el parámetro timeout del cliente de OpenAI o envuelva la transmisión con asyncio.wait_for() para cancelarla después de una duración máxima.

import asyncio

async def stream_with_timeout(question: str, timeout: float = 30.0):
    async def _stream():
        async for chunk in chain.astream({'question': question}):
            yield chunk

    try:
        async for chunk in asyncio.timeout(_stream(), timeout):
            print(chunk, end='', flush=True)
    except asyncio.TimeoutError:
        print('\n[Stream timed out after 30 seconds]')
    except asyncio.CancelledError:
        print('\n[Stream cancelled by client disconnect]')

SSE del lado del cliente con JavaScript

En el frontend, la API EventSource nativa del navegador consume los eventos enviados por el servidor. Cuando el endpoint de FastAPI emite fragmentos data: token\n\n, EventSource activa un evento message para cada uno. Añada cada token al DOM a medida que llega para crear un efecto de máquina de escribir. Para tener más control, fetch() con response.body.getReader() proporciona acceso completo a la transmisión.

// Frontend JavaScript (not Python)
const source = new EventSource('/stream?topic=quantum+computing');
const outputDiv = document.getElementById('output');

source.onmessage = (event) => {
    outputDiv.textContent += event.data;
};

source.onerror = () => {
    source.close();
    outputDiv.textContent += ' [done]';
};

// Alternative: fetch with ReadableStream
const response = await fetch('/stream?topic=ai');
const reader = response.body.getReader();
while (true) {
    const {done, value} = await reader.read();
    if (done) break;
    outputDiv.textContent += new TextDecoder().decode(value);
}

Prácticas recomendadas para la transmisión

Siga estas prácticas recomendadas al implementar la transmisión: use siempre flush=True al imprimir en stdout para evitar el almacenamiento en búfer. Establezca stream_usage=True si necesita recuentos precisos de tokens durante la transmisión. Emita un marcador data: [DONE]\n\n al final de las transmisiones SSE para que el cliente sepa cuándo cerrar la conexión. Pruebe los endpoints de transmisión con curl --no-buffer para verificar que los tokens llegan progresivamente.

# Complete SSE endpoint with DONE sentinel
async def sse_generator(question: str):
    try:
        async for chunk in chain.astream({'question': question}):
            # Escape any newlines in the chunk
            safe_chunk = chunk.replace('\n', ' ')
            yield f'data: {safe_chunk}\n\n'
    finally:
        yield 'data: [DONE]\n\n'

@app.get('/chat/stream')
async def chat_stream(question: str):
    return StreamingResponse(
        sse_generator(question),
        media_type='text/event-stream',
        headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}
    )

Comprobación rápida

Compruebe sus conocimientos sobre la transmisión de resultados en LangChain.

Resumen de la lección

En esta lección ha aprendido que stream() y astream() permiten iterar sobre fragmentos de tokens a medida que se generan, eliminando la larga espera de la respuesta completa; StreamingResponse en FastAPI, con formato SSE, entrega tokens a los clientes del navegador en tiempo real; y astream_events() proporciona hooks de eventos detallados para cada paso de la cadena, incluidas las llamadas a herramientas y las salidas intermedias. A continuación, exploraremos la gestión de memoria para conversaciones de varios turnos.

Gratis para empezar

Aprende Python con un tutor de IA — gratis

Escribe y ejecuta código real en tu navegador, obtén ayuda instantánea de un tutor de IA disponible 24/7 y continúa donde lo dejaste en la web o en la aplicación.

Cursos
30
Lecciones
120

Preguntas frecuentes

¿La lección «Streaming de salida en LangChain» es gratis?

Sí — el texto completo de «Streaming de salida en LangChain» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de AI Engineering Academy, actualiza a CoddyKit PRO. El curso de AI Engineering Academy incluye 4 lecciones en total.

¿Qué aprenderé en «Streaming de salida en LangChain»?

Implementará streaming de tokens mediante chains de LCEL para que su aplicación muestre cada palabra a medida que llega, en lugar de esperar a la respuesta completa, mejorando la latencia percibida. Practicas AI Engineering Academy con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.

¿Necesito experiencia previa para empezar AI Engineering Academy?

No se requiere experiencia previa. AI Engineering Academy en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 4 de 4.

¿Cuánto tiempo toma la lección «Streaming de salida en LangChain»?

La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.

¿Puedo escribir y ejecutar código en esta lección de AI Engineering Academy?

Sí. Cada lección de AI Engineering Academy incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.

Todas las lecciones de este curso

  1. Arquitectura de LangChain y abstracciones principales
  2. Creación de chains con LCEL
  3. Chains ramificadas y paralelas
  4. Streaming de salida en LangChain
← Volver a AI Engineering Academy