From 6c171e2838a2bef7aa63aa774bb1badd2ad3754a Mon Sep 17 00:00:00 2001 From: Vadym Barda Date: Wed, 25 Sep 2024 11:23:31 -0400 Subject: [PATCH] docs: update streaming content how-to (#1841) --- docs/docs/how-tos/index.md | 2 +- docs/docs/how-tos/streaming-content.ipynb | 198 ++++++++++++++++------ docs/mkdocs.yml | 2 +- 3 files changed, 147 insertions(+), 55 deletions(-) diff --git a/docs/docs/how-tos/index.md b/docs/docs/how-tos/index.md index 518064348..fc1f14de5 100644 --- a/docs/docs/how-tos/index.md +++ b/docs/docs/how-tos/index.md @@ -50,7 +50,7 @@ These guides show how to use different streaming modes. - [How to stream state updates of your graph](stream-updates.ipynb) - [How to stream LLM tokens](streaming-tokens.ipynb) - [How to stream LLM tokens without LangChain models](streaming-tokens-without-langchain.ipynb) -- [How to stream arbitrarily nested content](streaming-content.ipynb) +- [How to stream custom data](streaming-content.ipynb) - [How to configure multiple streaming modes at the same time](stream-multiple.ipynb) - [How to stream events from within a tool](streaming-events-from-within-tools.ipynb) - [How to stream events from within a tool without LangChain models](streaming-events-from-within-tools-without-langchain.ipynb) diff --git a/docs/docs/how-tos/streaming-content.ipynb b/docs/docs/how-tos/streaming-content.ipynb index 84f60a19f..2f7b70a4a 100644 --- a/docs/docs/how-tos/streaming-content.ipynb +++ b/docs/docs/how-tos/streaming-content.ipynb @@ -5,13 +5,15 @@ "id": "15c4bd28", "metadata": {}, "source": [ - "# How to stream arbitrary nested content\n", + "# How to stream custom data\n", "\n", - "The most common use case for streaming from inside a node is to stream LLM tokens, but you may have other long-running streaming functions you wish to render for the user. While individual nodes in LangGraph cannot return generators (since they are executed to completion for each [superstep](https://langchain-ai.github.io/langgraph/concepts/low_level)), we can still stream arbitrary custom functions from within a node using a similar tact and calling `astream_events` on the graph.\n", + "The most common use case for streaming from inside a node is to stream LLM tokens, but you may also want to stream custom data. For example, you might have some long-running streaming functions you may wish to render for the user. \n", "\n", - "We do so using a [RunnableGenerator](https://api.python.langchain.com/en/latest/runnables/langchain_core.runnables.base.RunnableGenerator.html#langchain-core-runnables-base-runnablegenerator) (which your function will automatically behave as if wrapped as a [RunnableLambda](https://api.python.langchain.com/en/latest/runnables/langchain_core.runnables.base.RunnableLambda.html#langchain_core.runnables.base.RunnableLambda)).\n", + "You can do so in two ways:\n", + "* using graph's `.stream` / `.astream` methods with `stream_mode=\"custom\"`\n", + "* emitting custom events using [adispatch_custom_events](https://python.langchain.com/docs/how_to/callbacks_custom_events/).\n", "\n", - "Below is a simple toy example.\n", + "Below is a simple toy example that shows both.\n", "\n", "## Setup\n", "\n", @@ -20,7 +22,7 @@ }, { "cell_type": "code", - "execution_count": null, + "execution_count": 1, "id": "e1a20f31", "metadata": {}, "outputs": [], @@ -44,46 +46,36 @@ }, { "cell_type": "markdown", - "id": "95301021-1db9-426f-807c-ec5b37bd5a9d", + "id": "29814253-ca9b-4844-a8a5-d6b19fbdbdba", "metadata": {}, "source": [ - "
\n", - "

ASYNC IN PYTHON<=3.10

\n", - "

\n", - "Any Langchain RunnableLambda, a RunnableGenerator, or Tool that invokes other runnables and is running async in python<=3.10, will have to propagate callbacks to child objects manually. This is because LangChain cannot automatically propagate callbacks to child objects in this case.\n", - " \n", - "This is a common reason why you may fail to see events being emitted from custom runnables or tools.\n", - "

\n", - "
" + "## Stream custom data using `.stream / .astream`" ] }, { "cell_type": "markdown", - "id": "07d5779e", + "id": "b729644a-b65f-4e69-ad45-f2e88ffb4e9d", "metadata": {}, "source": [ - "## Define the graph" + "### Define the graph" ] }, { "cell_type": "code", - "execution_count": 1, - "id": "486a01a0", + "execution_count": 2, + "id": "9731c40f-5ce7-460d-b2ad-33185529c99d", "metadata": {}, "outputs": [], "source": [ "from langchain_core.messages import AIMessage\n", - "from langchain_core.runnables import RunnableGenerator\n", - "from langchain_core.runnables import RunnableConfig\n", - "\n", "from langgraph.graph import START, StateGraph, MessagesState, END\n", + "from langgraph.types import StreamWriter\n", "\n", - "# Define a new graph\n", - "workflow = StateGraph(MessagesState)\n", - "\n", - "\n", - "async def my_generator(state: MessagesState):\n", - " messages = [\n", + "async def my_node(\n", + " state: MessagesState, \n", + " writer: StreamWriter # <-- provide StreamWriter to write chunks to be streamed\n", + "):\n", + " chunks = [\n", " \"Four\",\n", " \"score\",\n", " \"and\",\n", @@ -94,42 +86,146 @@ " \"fathers\",\n", " \"...\",\n", " ]\n", - " for message in messages:\n", - " yield message\n", + " for chunk in chunks:\n", + " # write the chunk to be streamed using stream_mode=custom \n", + " writer(chunk)\n", "\n", + " return {\"messages\": [AIMessage(content=\" \".join(chunks))]}\n", "\n", - "async def my_node(state: MessagesState, config: RunnableConfig):\n", - " messages = []\n", - " # Tagging a node makes it easy to filter out which events to include in your stream\n", - " # It's completely optional, but useful if you have many functions with similar names\n", - " gen = RunnableGenerator(my_generator).with_config(\n", - " tags=[\"should_stream\"],\n", - " callbacks=config.get(\n", - " \"callbacks\", []\n", - " ), # <-- Propagate callbacks (Python <= 3.10)\n", - " )\n", - " async for message in gen.astream(state):\n", - " messages.append(message)\n", - " return {\"messages\": [AIMessage(content=\" \".join(messages))]}\n", - "\n", + "# Define a new graph\n", + "workflow = StateGraph(MessagesState)\n", "\n", "workflow.add_node(\"model\", my_node)\n", "workflow.add_edge(START, \"model\")\n", "workflow.add_edge(\"model\", END)\n", + "\n", "app = workflow.compile()" ] }, { "cell_type": "markdown", - "id": "2af9e94e", + "id": "ecd69eed-9624-4640-b0af-c9f82b190900", "metadata": {}, "source": [ - "## Stream arbitrarily nested content" + "### Stream content" ] }, { "cell_type": "code", - "execution_count": 2, + "execution_count": 3, + "id": "00a91b15-82c7-443c-acb6-a7406df15cee", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "Four|score|and|seven|years|ago|our|fathers|...|" + ] + } + ], + "source": [ + "from langchain_core.messages import HumanMessage\n", + "\n", + "inputs = [HumanMessage(content=\"What are you thinking about?\")]\n", + "async for chunk in app.astream({\"messages\": inputs}, stream_mode=\"custom\"):\n", + " print(chunk, end=\"|\", flush=True)" + ] + }, + { + "cell_type": "markdown", + "id": "29035302-3111-45bf-ac69-50ab940f8cb4", + "metadata": {}, + "source": [ + "## Stream custom data using `.astream_events`" + ] + }, + { + "cell_type": "markdown", + "id": "822e91c3-03be-4778-9fa5-a6ec57be3e52", + "metadata": {}, + "source": [ + "If you are already using graph's `.astream_events` method in your workflow, you can also stream custom data by emitting custom events using `adispatch_custom_event`" + ] + }, + { + "cell_type": "markdown", + "id": "0fb6c3e5-7377-4f93-a8c6-44582ee3bc1a", + "metadata": {}, + "source": [ + "
\n", + "

ASYNC IN PYTHON<=3.10

\n", + "

\n", + "\n", + "LangChain cannot automatically propagate configuration, including callbacks necessary for `astream_events()`, to child runnables if you are running async code in python<=3.10. This is a common reason why you may fail to see events being emitted from custom runnables or tools.\n", + "\n", + "If you are running python<=3.10, you will need to manually propagate the `RunnableConfig` object to the child runnable in async environments. For an example of how to manually propagate the config, see the implementation of the node below with `adispatch_custom_event`.\n", + "\n", + "If you are running python>=3.11, the `RunnableConfig` will automatically propagate to child runnables in async environment. However, it is still a good idea to propagate the `RunnableConfig` manually if your code may run in other Python versions.\n", + "

\n", + "
" + ] + }, + { + "cell_type": "markdown", + "id": "b390a9fe-2d5f-4e82-a1ea-c7c0186b8559", + "metadata": {}, + "source": [ + "### Define the graph" + ] + }, + { + "cell_type": "code", + "execution_count": 6, + "id": "486a01a0", + "metadata": {}, + "outputs": [], + "source": [ + "from langchain_core.runnables import RunnableConfig, RunnableLambda\n", + "from langchain_core.callbacks.manager import adispatch_custom_event\n", + "\n", + "async def my_node(state: MessagesState, config: RunnableConfig):\n", + " chunks = [\n", + " \"Four\",\n", + " \"score\",\n", + " \"and\",\n", + " \"seven\",\n", + " \"years\",\n", + " \"ago\",\n", + " \"our\",\n", + " \"fathers\",\n", + " \"...\",\n", + " ]\n", + " for chunk in chunks:\n", + " await adispatch_custom_event(\n", + " \"my_custom_event\",\n", + " {\"chunk\": chunk},\n", + " config=config # <-- propagate config\n", + " )\n", + "\n", + " return {\"messages\": [AIMessage(content=\" \".join(chunks))]}\n", + "\n", + "# Define a new graph\n", + "workflow = StateGraph(MessagesState)\n", + "\n", + "workflow.add_node(\"model\", my_node)\n", + "workflow.add_edge(START, \"model\")\n", + "workflow.add_edge(\"model\", END)\n", + "\n", + "app = workflow.compile()" + ] + }, + { + "cell_type": "markdown", + "id": "7dcded03-6776-405e-afae-005a3212d3e4", + "metadata": {}, + "source": [ + "### Stream content" + ] + }, + { + "cell_type": "code", + "execution_count": 7, "id": "ce773a40", "metadata": {}, "outputs": [ @@ -137,7 +233,7 @@ "name": "stdout", "output_type": "stream", "text": [ - "{'chunk': 'Four'}|{'chunk': 'score'}|{'chunk': 'and'}|{'chunk': 'seven'}|{'chunk': 'years'}|{'chunk': 'ago'}|{'chunk': 'our'}|{'chunk': 'fathers'}|{'chunk': '...'}|" + "Four|score|and|seven|years|ago|our|fathers|...|" ] } ], @@ -146,15 +242,11 @@ "\n", "inputs = [HumanMessage(content=\"What are you thinking about?\")]\n", "async for event in app.astream_events({\"messages\": inputs}, version=\"v2\"):\n", - " kind = event[\"event\"]\n", " tags = event.get(\"tags\", [])\n", - " if kind == \"on_chain_stream\" and \"should_stream\" in tags:\n", + " if event[\"event\"] == \"on_custom_event\" and event[\"name\"] == \"my_custom_event\":\n", " data = event[\"data\"]\n", " if data:\n", - " # Empty content in the context of OpenAI or Anthropic usually means\n", - " # that the model is asking for a tool to be invoked.\n", - " # So we only print non-empty content\n", - " print(data, end=\"|\")" + " print(data[\"chunk\"], end=\"|\", flush=True)" ] } ], diff --git a/docs/mkdocs.yml b/docs/mkdocs.yml index 4745bd491..b91d8034a 100644 --- a/docs/mkdocs.yml +++ b/docs/mkdocs.yml @@ -153,7 +153,7 @@ nav: - Stream state updates: how-tos/stream-updates.ipynb - Stream LLM tokens: how-tos/streaming-tokens.ipynb - Stream LLM tokens without LangChain models: how-tos/streaming-tokens-without-langchain.ipynb - - Stream arbitrarily nested content: how-tos/streaming-content.ipynb + - Stream custom data: how-tos/streaming-content.ipynb - Configure multiple streaming modes: how-tos/stream-multiple.ipynb - Stream events from within tools: how-tos/streaming-events-from-within-tools.ipynb - Stream events from within tools without LangChain models: how-tos/streaming-events-from-within-tools-without-langchain.ipynb