How AgentExecutor in LCEL makes LLM streaming output in the final step #61

Closed
opened 2026-02-16 00:18:24 -05:00 by yindo · 1 comment
Owner

Originally created by @wangcailin on GitHub (Dec 8, 2023).

How can AgentExecutor in LCEL in the last step how to make LLM streaming output, my business logic is

  1. call OpenAI function calling
  2. output the final result according to the content
    Currently the business processes are all achievable, but the final answer cannot be streamed through astream, astream streams the steps of AgentExecutor, and what I want to achieve is to stream the LLM answer of the last step.
    @eyurtsev Please help guide me on how to achieve this, thanks
from langchain.agents.output_parsers import JSONAgentOutputParser
from typing import AsyncIterator
from typing import List, Tuple, Dict
from datetime import datetime
from langchain.agents import (
    AgentExecutor,
)
from langchain.callbacks.streaming_stdout_final_only import (
    FinalStreamingStdOutCallbackHandler,
)
from langchain.schema.runnable import Runnable, RunnableLambda, RunnableParallel
from langchain.agents.format_scratchpad import format_to_openai_function_messages
from langchain.agents.output_parsers import OpenAIFunctionsAgentOutputParser
from langchain.chat_models import ChatOpenAI
from langchain.prompts import ChatPromptTemplate, MessagesPlaceholder
from langchain.pydantic_v1 import BaseModel, Field
from langchain.schema.messages import AIMessage, HumanMessage
from langchain.tools.render import format_tool_to_openai_function
from .tools import agent_tools, _get_tools
from .callbacks import StreamingStdOutCallbackHandler

# gpt-4-1106-preview
llm = ChatOpenAI(streaming=True, callbacks=[
    FinalStreamingStdOutCallbackHandler(
        answer_prefix_tokens=["The", "answer", ":"])
], temperature=0.7, model_name="gpt-3.5-turbo-1106")
llm = llm.with_fallbacks(
    [ChatOpenAI(streaming=True, temperature=0.7, model_name="gpt-3.5-turbo-1106")])

assistant_system_message = """{custom_prompt}

当前用户ID:{openid}
当前时间:{current_time}
当前时区:东八区UTC/GMT+08:00

"""
prompt = ChatPromptTemplate.from_messages(
    [
        ("system", assistant_system_message),
        MessagesPlaceholder(variable_name="chat_history"),
        ("user", "{input}"),
        MessagesPlaceholder(variable_name="agent_scratchpad"),
    ]
)


def _llm_with_tools(input: Dict) -> Runnable:
    return RunnableLambda(lambda x: x["input"]) | llm.bind(
        functions=input["functions"]
    )


def _format_chat_history(chat_history: List[Tuple[str, str]]):
    buffer = []
    for human, ai in chat_history:
        buffer.append(HumanMessage(content=human))
        buffer.append(AIMessage(content=ai))
    return buffer


def _get_current_time() -> str:
    current_time = datetime.now()
    formatted_time = current_time.strftime('%Y年%m月%d日 %H:%M')
    return formatted_time


agent = (
    RunnableParallel({
        "input": lambda x: x["input"],
        "custom_prompt": lambda x: x["custom_prompt"],
        "chat_history": lambda x: _format_chat_history(x["chat_history"]),
        "agent_scratchpad": lambda x: format_to_openai_function_messages(
            x["intermediate_steps"]
        ),
        "functions": lambda x: [
            format_tool_to_openai_function(tool) for tool in _get_tools(x["rule_id"])
        ],
        "openid": lambda x: x["openid"],
        "current_time": lambda x: _get_current_time(),
    })
    | {
        "input": prompt,
        "functions": lambda x: x["functions"],
    }
    | _llm_with_tools
    | OpenAIFunctionsAgentOutputParser()
)


class AgentInput(BaseModel):
    input: str
    chat_history: List[Tuple[str, str]] = Field(
        ..., extra={"widget": {"type": "chat", "input": "input", "output": "output"}}
    )
    openid: str
    rule_id: str
    custom_prompt: str


agent_executor = AgentExecutor(agent=agent, tools=agent_tools, handle_parsing_errors=True, verbose=False).with_types(
    input_type=AgentInput
) | (lambda x: x["output"])


async def async_stream(message, openid, rule_id, chat_history, custom_prompt):
    # agent_executor.agent.llm_chain.llm.callbacks = [
    #     StreamingStdOutCallbackHandler()]
    async for token in agent_executor.astream(
            {"chat_history": chat_history, "rule_id": rule_id, "openid": openid, "input": message, "custom_prompt": custom_prompt}):
        print('token#', token)
        yield token

```from langchain.agents.output_parsers import JSONAgentOutputParser
from typing import AsyncIterator
from typing import List, Tuple, Dict
from datetime import datetime
from langchain.agents import (
    AgentExecutor,
)
from langchain.callbacks.streaming_stdout_final_only import (
    FinalStreamingStdOutCallbackHandler,
)
from langchain.schema.runnable import Runnable, RunnableLambda, RunnableParallel
from langchain.agents.format_scratchpad import format_to_openai_function_messages
from langchain.agents.output_parsers import OpenAIFunctionsAgentOutputParser
from langchain.chat_models import ChatOpenAI
from langchain.prompts import ChatPromptTemplate, MessagesPlaceholder
from langchain.pydantic_v1 import BaseModel, Field
from langchain.schema.messages import AIMessage, HumanMessage
from langchain.tools.render import format_tool_to_openai_function
from .tools import agent_tools, _get_tools
from .callbacks import StreamingStdOutCallbackHandler

# gpt-4-1106-preview
llm = ChatOpenAI(streaming=True, callbacks=[
    FinalStreamingStdOutCallbackHandler(
        answer_prefix_tokens=["The", "answer", ":"])
], temperature=0.7, model_name="gpt-3.5-turbo-1106")
llm = llm.with_fallbacks(
    [ChatOpenAI(streaming=True, temperature=0.7, model_name="gpt-3.5-turbo-1106")])

assistant_system_message = """{custom_prompt}

当前用户ID:{openid}
当前时间:{current_time}
当前时区:东八区UTC/GMT+08:00

"""
prompt = ChatPromptTemplate.from_messages(
    [
        ("system", assistant_system_message),
        MessagesPlaceholder(variable_name="chat_history"),
        ("user", "{input}"),
        MessagesPlaceholder(variable_name="agent_scratchpad"),
    ]
)


def _llm_with_tools(input: Dict) -> Runnable:
    return RunnableLambda(lambda x: x["input"]) | llm.bind(
        functions=input["functions"]
    )


def _format_chat_history(chat_history: List[Tuple[str, str]]):
    buffer = []
    for human, ai in chat_history:
        buffer.append(HumanMessage(content=human))
        buffer.append(AIMessage(content=ai))
    return buffer


def _get_current_time() -> str:
    current_time = datetime.now()
    formatted_time = current_time.strftime('%Y年%m月%d日 %H:%M')
    return formatted_time


agent = (
    RunnableParallel({
        "input": lambda x: x["input"],
        "custom_prompt": lambda x: x["custom_prompt"],
        "chat_history": lambda x: _format_chat_history(x["chat_history"]),
        "agent_scratchpad": lambda x: format_to_openai_function_messages(
            x["intermediate_steps"]
        ),
        "functions": lambda x: [
            format_tool_to_openai_function(tool) for tool in _get_tools(x["rule_id"])
        ],
        "openid": lambda x: x["openid"],
        "current_time": lambda x: _get_current_time(),
    })
    | {
        "input": prompt,
        "functions": lambda x: x["functions"],
    }
    | _llm_with_tools
    | OpenAIFunctionsAgentOutputParser()
)


class AgentInput(BaseModel):
    input: str
    chat_history: List[Tuple[str, str]] = Field(
        ..., extra={"widget": {"type": "chat", "input": "input", "output": "output"}}
    )
    openid: str
    rule_id: str
    custom_prompt: str


agent_executor = AgentExecutor(agent=agent, tools=agent_tools, handle_parsing_errors=True, verbose=False).with_types(
    input_type=AgentInput
) | (lambda x: x["output"])


async def async_stream(message, openid, rule_id, chat_history, custom_prompt):
    # agent_executor.agent.llm_chain.llm.callbacks = [
    #     StreamingStdOutCallbackHandler()]
    async for token in agent_executor.astream(
            {"chat_history": chat_history, "rule_id": rule_id, "openid": openid, "input": message, "custom_prompt": custom_prompt}):
        print('token#', token)
        yield token
Originally created by @wangcailin on GitHub (Dec 8, 2023). How can AgentExecutor in LCEL in the last step how to make LLM streaming output, my business logic is 1. call OpenAI function calling 2. output the final result according to the content Currently the business processes are all achievable, but the final answer cannot be streamed through astream, astream streams the steps of AgentExecutor, and what I want to achieve is to stream the LLM answer of the last step. @eyurtsev Please help guide me on how to achieve this, thanks ``` from langchain.agents.output_parsers import JSONAgentOutputParser from typing import AsyncIterator from typing import List, Tuple, Dict from datetime import datetime from langchain.agents import ( AgentExecutor, ) from langchain.callbacks.streaming_stdout_final_only import ( FinalStreamingStdOutCallbackHandler, ) from langchain.schema.runnable import Runnable, RunnableLambda, RunnableParallel from langchain.agents.format_scratchpad import format_to_openai_function_messages from langchain.agents.output_parsers import OpenAIFunctionsAgentOutputParser from langchain.chat_models import ChatOpenAI from langchain.prompts import ChatPromptTemplate, MessagesPlaceholder from langchain.pydantic_v1 import BaseModel, Field from langchain.schema.messages import AIMessage, HumanMessage from langchain.tools.render import format_tool_to_openai_function from .tools import agent_tools, _get_tools from .callbacks import StreamingStdOutCallbackHandler # gpt-4-1106-preview llm = ChatOpenAI(streaming=True, callbacks=[ FinalStreamingStdOutCallbackHandler( answer_prefix_tokens=["The", "answer", ":"]) ], temperature=0.7, model_name="gpt-3.5-turbo-1106") llm = llm.with_fallbacks( [ChatOpenAI(streaming=True, temperature=0.7, model_name="gpt-3.5-turbo-1106")]) assistant_system_message = """{custom_prompt} 当前用户ID:{openid} 当前时间:{current_time} 当前时区:东八区UTC/GMT+08:00 """ prompt = ChatPromptTemplate.from_messages( [ ("system", assistant_system_message), MessagesPlaceholder(variable_name="chat_history"), ("user", "{input}"), MessagesPlaceholder(variable_name="agent_scratchpad"), ] ) def _llm_with_tools(input: Dict) -> Runnable: return RunnableLambda(lambda x: x["input"]) | llm.bind( functions=input["functions"] ) def _format_chat_history(chat_history: List[Tuple[str, str]]): buffer = [] for human, ai in chat_history: buffer.append(HumanMessage(content=human)) buffer.append(AIMessage(content=ai)) return buffer def _get_current_time() -> str: current_time = datetime.now() formatted_time = current_time.strftime('%Y年%m月%d日 %H:%M') return formatted_time agent = ( RunnableParallel({ "input": lambda x: x["input"], "custom_prompt": lambda x: x["custom_prompt"], "chat_history": lambda x: _format_chat_history(x["chat_history"]), "agent_scratchpad": lambda x: format_to_openai_function_messages( x["intermediate_steps"] ), "functions": lambda x: [ format_tool_to_openai_function(tool) for tool in _get_tools(x["rule_id"]) ], "openid": lambda x: x["openid"], "current_time": lambda x: _get_current_time(), }) | { "input": prompt, "functions": lambda x: x["functions"], } | _llm_with_tools | OpenAIFunctionsAgentOutputParser() ) class AgentInput(BaseModel): input: str chat_history: List[Tuple[str, str]] = Field( ..., extra={"widget": {"type": "chat", "input": "input", "output": "output"}} ) openid: str rule_id: str custom_prompt: str agent_executor = AgentExecutor(agent=agent, tools=agent_tools, handle_parsing_errors=True, verbose=False).with_types( input_type=AgentInput ) | (lambda x: x["output"]) async def async_stream(message, openid, rule_id, chat_history, custom_prompt): # agent_executor.agent.llm_chain.llm.callbacks = [ # StreamingStdOutCallbackHandler()] async for token in agent_executor.astream( {"chat_history": chat_history, "rule_id": rule_id, "openid": openid, "input": message, "custom_prompt": custom_prompt}): print('token#', token) yield token ```from langchain.agents.output_parsers import JSONAgentOutputParser from typing import AsyncIterator from typing import List, Tuple, Dict from datetime import datetime from langchain.agents import ( AgentExecutor, ) from langchain.callbacks.streaming_stdout_final_only import ( FinalStreamingStdOutCallbackHandler, ) from langchain.schema.runnable import Runnable, RunnableLambda, RunnableParallel from langchain.agents.format_scratchpad import format_to_openai_function_messages from langchain.agents.output_parsers import OpenAIFunctionsAgentOutputParser from langchain.chat_models import ChatOpenAI from langchain.prompts import ChatPromptTemplate, MessagesPlaceholder from langchain.pydantic_v1 import BaseModel, Field from langchain.schema.messages import AIMessage, HumanMessage from langchain.tools.render import format_tool_to_openai_function from .tools import agent_tools, _get_tools from .callbacks import StreamingStdOutCallbackHandler # gpt-4-1106-preview llm = ChatOpenAI(streaming=True, callbacks=[ FinalStreamingStdOutCallbackHandler( answer_prefix_tokens=["The", "answer", ":"]) ], temperature=0.7, model_name="gpt-3.5-turbo-1106") llm = llm.with_fallbacks( [ChatOpenAI(streaming=True, temperature=0.7, model_name="gpt-3.5-turbo-1106")]) assistant_system_message = """{custom_prompt} 当前用户ID:{openid} 当前时间:{current_time} 当前时区:东八区UTC/GMT+08:00 """ prompt = ChatPromptTemplate.from_messages( [ ("system", assistant_system_message), MessagesPlaceholder(variable_name="chat_history"), ("user", "{input}"), MessagesPlaceholder(variable_name="agent_scratchpad"), ] ) def _llm_with_tools(input: Dict) -> Runnable: return RunnableLambda(lambda x: x["input"]) | llm.bind( functions=input["functions"] ) def _format_chat_history(chat_history: List[Tuple[str, str]]): buffer = [] for human, ai in chat_history: buffer.append(HumanMessage(content=human)) buffer.append(AIMessage(content=ai)) return buffer def _get_current_time() -> str: current_time = datetime.now() formatted_time = current_time.strftime('%Y年%m月%d日 %H:%M') return formatted_time agent = ( RunnableParallel({ "input": lambda x: x["input"], "custom_prompt": lambda x: x["custom_prompt"], "chat_history": lambda x: _format_chat_history(x["chat_history"]), "agent_scratchpad": lambda x: format_to_openai_function_messages( x["intermediate_steps"] ), "functions": lambda x: [ format_tool_to_openai_function(tool) for tool in _get_tools(x["rule_id"]) ], "openid": lambda x: x["openid"], "current_time": lambda x: _get_current_time(), }) | { "input": prompt, "functions": lambda x: x["functions"], } | _llm_with_tools | OpenAIFunctionsAgentOutputParser() ) class AgentInput(BaseModel): input: str chat_history: List[Tuple[str, str]] = Field( ..., extra={"widget": {"type": "chat", "input": "input", "output": "output"}} ) openid: str rule_id: str custom_prompt: str agent_executor = AgentExecutor(agent=agent, tools=agent_tools, handle_parsing_errors=True, verbose=False).with_types( input_type=AgentInput ) | (lambda x: x["output"]) async def async_stream(message, openid, rule_id, chat_history, custom_prompt): # agent_executor.agent.llm_chain.llm.callbacks = [ # StreamingStdOutCallbackHandler()] async for token in agent_executor.astream( {"chat_history": chat_history, "rule_id": rule_id, "openid": openid, "input": message, "custom_prompt": custom_prompt}): print('token#', token) yield token
yindo closed this issue 2026-02-16 00:18:24 -05:00
Author
Owner

@eyurtsev commented on GitHub (Dec 8, 2023):

This does not appear to be a langserve issue, but an LCEL question. Please let me know if i got that wrong.

Likely the issue is that the code is using RunnableLambda which does not preserve streaming capabilities. Instead use RunnableGenerator

from langchain.chat_models import ChatAnthropic
from langchain.schema.runnable import RunnablePassthrough, RunnableGenerator

chain = ChatAnthropic() # Streams

for chunk in chain.stream('hello'):
    print(chunk)   

chain  = ChatAnthropic() | (lambda x: x) # Does not stream

for chunk in chain.stream('hello'):
    print(chunk)

chain = ChatAnthropic() | RunnablePassthrough() # Will stream because it defines `transform`

for chunk in chain.stream('hello'):
    print(chunk)

# Let's define our own

def _transform(input_stream):
    for chunk in input_stream:
        yield chunk['output']

chain = {
    'output': ChatAnthropic()
} | RunnableGenerator(_transform)

for chunk in chain.stream('hello'):
    print(chunk)
@eyurtsev commented on GitHub (Dec 8, 2023): This does not appear to be a langserve issue, but an LCEL question. Please let me know if i got that wrong. Likely the issue is that the code is using RunnableLambda which does not preserve streaming capabilities. Instead use `RunnableGenerator` ```python from langchain.chat_models import ChatAnthropic from langchain.schema.runnable import RunnablePassthrough, RunnableGenerator chain = ChatAnthropic() # Streams for chunk in chain.stream('hello'): print(chunk) chain = ChatAnthropic() | (lambda x: x) # Does not stream for chunk in chain.stream('hello'): print(chunk) chain = ChatAnthropic() | RunnablePassthrough() # Will stream because it defines `transform` for chunk in chain.stream('hello'): print(chunk) # Let's define our own def _transform(input_stream): for chunk in input_stream: yield chunk['output'] chain = { 'output': ChatAnthropic() } | RunnableGenerator(_transform) for chunk in chain.stream('hello'): print(chunk) ```
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: langchain-ai/langserve#61