skip to content

Why does an LCEL chain stop streaming token-by-token, and how do you fix it?

level: seniorimportance: should knowfreq 50%

answer

  1. Streaming is cooperative between steps
  2. Default transform buffers, then yields once
  3. Generator-style lambdas keep chunks flowing
  4. astream_events shows where it stalls
  5. Whole-output steps and streaming are exclusive

basics

~20 s

A step that only implements invoke must consume the whole upstream output before producing anything, so the chain emits one big chunk. Fix it by making that step a generator-based RunnableLambda, or move it out of the streaming path.

solid answer

~50 s

Streaming in LCEL is cooperative. A `RunnableSequence` streams by asking each step to transform an *iterator* of chunks into an iterator of chunks. Any step that does not implement that — the default falls back to buffering the whole input, calling `invoke`, and yielding one chunk — becomes a dam: everything upstream of it still streams internally, but the caller sees a single chunk at the end. The usual culprit is a `RunnableLambda` wrapping an ordinary function, or any step that genuinely needs the complete value. Two fixes: write the function as a **generator that takes an iterator and yields**, which `RunnableLambda` will drive in transform mode, or restructure so the blocking step sits before the final streaming step rather than after it. When you need visibility into what is actually streaming, `astream_events` emits per-step events such as `on_chat_model_stream`, so you can see which step first buffers.

code

python · 15 lines
python
from typing import Iterator
from langchain_core.runnables import RunnableLambda

# Blocking: needs the whole string before it can return anything
shout_blocking = RunnableLambda(lambda s: s.upper())

# Streaming: consumes and yields chunk by chunk
def shout_stream(chunks: Iterator[str]) -> Iterator[str]:
    for chunk in chunks:
        yield chunk.upper()

shout_streaming = RunnableLambda(shout_stream)

print(list(shout_blocking.transform(iter(["ab", "cd"]))))    # ['ABCD']   one chunk
print(list(shout_streaming.transform(iter(["ab", "cd"]))))   # ['AB', 'CD'] chunks kept

go deeper

for a junior

Know that calling stream on a chain gives you chunks, and that if you only ever get one chunk something in the chain is waiting for the whole output.

for a middle

Explain the buffering default: a step without an incremental implementation accumulates the upstream chunks, calls invoke, and yields once — so one non-streaming step flattens the whole chain.

for a senior

Show the diagnosis and the fix: name your steps, read astream_events to find the dam, rewrite incremental steps as generators, and move audit or logging work onto listeners instead of into the data path.

for a principal

Own the product tradeoff. Whole-answer validation and token streaming cannot both sit on the user-facing path, so decide deliberately whether perceived latency or post-hoc guarantees wins, and design the fallback for when the streamed answer must be retracted.

## How streaming actually propagates `Runnable` has a second, lower-level pair of methods behind `stream`/`astream`: `transform` and `atransform`, which take an *iterator of inputs* and yield an *iterator of outputs*. `RunnableSequence.stream` wires the steps together with these, so chunks flow through the pipeline as they are produced. The base class provides a default `transform` for any component that has not implemented one: accumulate everything the upstream yields, add the chunks together into the final value, call `invoke` on it, and yield exactly one output chunk. This default is why `stream` never *errors* on a non-streaming component — and why it silently degrades. The chain still returns an iterator; it just yields once, at the end. ## Diagnosing it Symptoms: a UI that shows nothing for eight seconds and then the whole answer at once, while the model provider's dashboard shows a streaming request. The chain is streaming; something after the model is buffering. To find the dam, use `astream_events`, which emits structured events for every step in the run — `on_chain_start`, `on_chat_model_stream`, `on_parser_stream` and so on, each tagged with the step's name and any tags you attached via `.with_config(tags=[...])`. If you see a healthy series of `on_chat_model_stream` events but the sequence's own output events arrive only once at the end, the step immediately downstream of the model is the one buffering. Naming your steps (`@chain` on a named function, or `.with_config(run_name="...")`) is what makes this trace readable at all — a wall of anonymous `RunnableLambda` entries tells you nothing. ## The fixes **1. Make the lambda a transform-style generator.** `RunnableLambda` supports functions written as generators that consume an iterator and yield: the function receives the upstream chunk iterator and yields output chunks as it goes. `RunnableGenerator` is the explicit class for the same job. This works whenever your transformation is *incremental* — uppercasing, prefixing, filtering, redacting token by token. **2. Reorder so the blocking work is upstream.** Some steps genuinely cannot be incremental: anything that needs the complete text to make a decision. If that step is post-processing the final answer, you have a real choice between streaming and that transformation — you cannot have both on the same output. Common resolution: stream the raw model output to the user, and run the blocking transformation on the accumulated result afterwards, out of band. **3. Move it off the streaming path.** Auditing, logging, and metrics do not need to be pipeline steps at all; `.with_listeners()` and callbacks observe a run without sitting in the data path and without blocking chunk flow. **4. Accept partial results where the format allows.** Some structured formats can be parsed incrementally and emit growing partial objects; a parser that requires the complete document cannot. If your consumer can handle partial payloads, an incremental format is what makes end-to-end streaming possible; if it cannot, the buffering is not a bug, it is the requirement. ## Async and back-pressure Use `astream` end to end when serving over HTTP. A sync `stream` from inside an event loop blocks the loop, and a sync-only step inside an async chain runs in a thread — correct, but it costs a worker thread per concurrent request, which is exactly the resource you were trying to protect by streaming. There is no explicit back-pressure protocol beyond ordinary iterator pull semantics: if your consumer is slow to iterate, upstream chunk production slows with it, which is usually what you want. ## Parallel blocks and streaming A fan-out step producing a dict of branch results generally cannot emit that dict incrementally — downstream steps see it once branches complete. Keep fan-out early and let the last step be the token producer. ## What to say in an interview The strong answer names the mechanism (`transform` over chunk iterators, with a buffering default), names the diagnostic (`astream_events` plus named steps), and states the honest tradeoff: a transformation that needs the whole answer and token streaming to the user are mutually exclusive on the same path, so you decide which one the product actually needs.

  • How would you confirm which step is buffering without adding print statements everywhere?
    Run the chain with `astream_events` and watch the event stream: you will see `on_chat_model_stream` events firing from the model while the surrounding step emits nothing until the end. Give steps real names first — `@chain` on named functions or `.with_config(run_name=...)` — because the events are keyed by step name and a chain of anonymous lambdas is unreadable in the trace.
  • Your final step must validate the complete answer before it is shown. Can you still stream?
    Not to the same consumer on the same path — validation over the whole answer forces buffering by definition. The usual resolutions are to stream the raw output optimistically and validate afterwards, accepting that you may have to retract; to stream to a hidden buffer and reveal on success, which is just buffering with extra steps; or to validate incrementally on partial output if the check permits it.
  • Does calling stream instead of invoke change what the model provider is billed for?
    No. Streaming changes how the response is delivered, not what is generated: the same prompt and the same completion tokens are billed either way. What it changes is perceived latency and your ability to abort early — if you stop consuming the iterator, generation can be cancelled, which is the only case where streaming genuinely saves tokens.

saying these in an interview costs you the question

  • Thinking stream returning an iterator means real streaming
  • Assuming any lambda in a chain streams automatically
  • Believing streaming reduces token cost by itself
  • Expecting a parallel fan-out to emit its dict incrementally
  • Adding logging as a pipeline step and blocking chunk flow

context