Building with LLMs

Building Modular Pipelines


Here is a function that exists, in some form, in almost every LLM codebase older than three months. It started as twenty lines.

Python
def process_document(path):    text = extract_text(path)    if len(text) < 50:        return {"error": "too short"}    if detect_language(text) != "en":        text = translate(text)    summary = llm.invoke(f"Summarise:\n{text}").content    if "```" in summary:        summary = summary.split("```")[1]    entities = json.loads(llm.invoke(f"Extract entities as JSON:\n{text}").content)    sentiment = llm.invoke(f"Sentiment of:\n{text}").content.strip().lower()    if sentiment not in ("positive", "negative", "neutral"):        sentiment = llm.invoke(f"Reply with one word only. Sentiment of:\n{text}").content    # ... 180 more lines

Now try to answer three questions about it.

Does the entity extractor work? You cannot tell without running the whole thing, which means an API call, which means money and four seconds. Why did last Tuesday's run fail? The traceback says json.loads raised, but not what the model returned or which document caused it. Can you run the three model calls concurrently? Not without rewriting the function, because the sequencing is baked into the control flow rather than declared anywhere.

The function is not badly written. It is badly shaped. Everything is fused: the steps, the ordering, the error handling and the I/O all live in one scope, so nothing can be exercised or replaced on its own. A modular pipeline is the alternative — small pieces with declared inputs and outputs, composed into a whole, where every piece is testable alone and the composition is data rather than control flow.

One verb per component, wiring as dataLoad andnormaliseBuild the promptCall the modelParse andvalidateRoute onthe resultEach arrow is a typed value, so any stage can be tested without the one before it.
The seam is where errors get caught and logged, which is why the seams are the design.

Single responsibility: one component, one verb

The first principle is the oldest one in software, and it is worth restating in LLM terms because the temptation to violate it is stronger here. A component should do one thing you can name with one verb: summarise, extract entities, classify sentiment, translate. Not "process".

Compare a fused component with the same work split:

FusedSplit
Testing the extractorRun everything, pay for 3 callsTest one component, mock one call
Improving the summary promptRisk breaking extractionTouch one file
Using a cheap model for sentimentNot possible — one model per functionBind a different model to that component
Running steps concurrentlyRewrite the functionChange the composition expression
Locating a failureRead the traceback and guessThe failing component names itself

Split, the same work looks like this:

Python
import osfrom langchain.chat_models import init_chat_modelfrom langchain_core.prompts import ChatPromptTemplatefrom langchain_core.output_parsers import StrOutputParser, JsonOutputParserfast  = init_chat_model(os.environ.get("FAST_MODEL", "anthropic:claude-haiku-4-5"))smart = init_chat_model(os.environ.get("SMART_MODEL", "anthropic:claude-sonnet-5"))summarise = (    ChatPromptTemplate.from_template("Summarise in 3 sentences:\n\n{text}")    | smart | StrOutputParser())extract_entities = (    ChatPromptTemplate.from_template(        'Extract entities. Return JSON: {{"people": [], "orgs": [], "dates": []}}\n\n{text}')    | smart | JsonOutputParser())classify_sentiment = (    ChatPromptTemplate.from_template(        "Sentiment: positive, negative, or neutral. One word only.\n\n{text}")    | fast | StrOutputParser())

Note what just became possible for free: classify_sentiment runs on the cheap model because that decision now lives in one place. If sentiment is 40% of your calls and the cheap model is five times cheaper, you have cut a meaningful slice of the bill with a one-word edit.

A component you cannot test without a network call is not a component. It is a section of a larger function that happens to have a name.

Composability: make the wiring data, not control flow

Once components share an interface, composition becomes an expression you can read, change and reason about. Steps that depend on each other go in a pipe; steps that do not go in parallel.

Python
from langchain_core.runnables import RunnableParallel, RunnablePassthroughanalyse = RunnableParallel(    summary=summarise,    entities=extract_entities,    sentiment=classify_sentiment,    source=RunnablePassthrough(),)result = analyse.invoke({"text": document_text})

The latency arithmetic is worth doing because it is the clearest payoff. Say the summariser takes 2.1 s, entity extraction 1.8 s and sentiment 0.6 s.

  • Sequential: 2.1 + 1.8 + 0.6 = 4.5 s
  • Parallel: bounded by the slowest branch = 2.1 s

A 53% reduction in wall-clock time for zero extra cost — the same three API calls are made either way. The fused function could not do this, not because its author did not know about concurrency, but because the shape forbade it.

The adapter problem

Real pipelines are rarely a clean chain of matching shapes. One step emits a string, the next wants a dictionary; a step needs both the original input and the previous output. Two small tools cover almost every case:

Python
from langchain_core.runnables import RunnableLambda# Reshape between stepsto_input = RunnableLambda(lambda s: {"text": s})# Keep the original alongside a computed valuewith_summary = RunnablePassthrough.assign(summary=summarise)# {"text": "..."}  ->  {"text": "...", "summary": "..."}

RunnablePassthrough.assign is the one people miss most often, and it removes a whole class of awkward lambdas: it adds a key to the dictionary flowing through, keeping everything already there. When a chain fails with a KeyError, the bug is nearly always a missing adapter rather than a broken prompt. Check the shapes first.

Error handling at the component boundary

In the fused function, one bad JSON response killed the whole document. That is the wrong blast radius. Each component should decide for itself what happens when it fails, and the default should almost always be "degrade, do not explode".

Python
import loggingdef safe(component, fallback, name):    """Wrap a component so failures are logged and a default is returned."""    def run(x):        try:            return component.invoke(x)        except Exception as exc:            logging.warning("component %s failed: %s", name, exc, exc_info=True)            return fallback    return RunnableLambda(run)analyse = RunnableParallel(    summary=safe(summarise, "", "summarise"),    entities=safe(extract_entities, {"people": [], "orgs": [], "dates": []}, "entities"),    sentiment=safe(classify_sentiment, "neutral", "sentiment"),)

Now a malformed JSON response costs you the entity list for one document. The summary and sentiment still arrive, the user still gets something useful, and the log tells you exactly which component failed and on what.

The judgement call is which failures deserve which treatment. A rough rule:

FailureBehaviourWhy
Optional enrichment (sentiment, tags)Log and return a defaultThe result is still useful without it
The primary output (the answer itself)Retry, then fail loudlyA silent empty answer is worse than an error
Malformed input (empty text, wrong type)Fail immediately, do not call the modelYou will pay for a call that cannot succeed
Authentication or quota errorsFail loudly, do not retryRetrying a 401 wastes time and hides the real problem

Decide the blast radius of every failure deliberately. The default in fused code is "everything dies", and that default is almost never what you would have chosen.

Conditional routing

Not every document should take the same path. A 200-word support email does not need the treatment a 40-page report needs, and paying for it is waste.

Python
from langchain_core.runnables import RunnableBranchshort_path = ChatPromptTemplate.from_template(    "Reply in one sentence:\n\n{text}") | fast | StrOutputParser()long_path = ChatPromptTemplate.from_template(    "Produce a structured analysis with headings:\n\n{text}") | smart | StrOutputParser()router = RunnableBranch(    (lambda x: len(x["text"].split()) < 300, short_path),    (lambda x: len(x["text"].split()) < 5000, long_path),    fallback_chain,          # the final argument is the default)

RunnableBranch evaluates conditions in order and runs the first match; the last argument is the default and is not optional. Two traps live here. The conditions are checked top to bottom, so a broad condition placed first silently shadows every narrower one below it. And the condition functions receive the raw input dictionary, not the output of anything — writing x["summary"] in a condition when summary is produced downstream gives you a KeyError, not a helpful message.

Routing pays for itself quickly. If 70% of your traffic is short documents and the cheap model costs a fifth of the expensive one, routing cuts the model spend to 0.7 × 0.2 + 0.3 × 1.0 = 0.44 of the unrouted cost — a 56% saving from one branch.

Retries and fallbacks: two different failures

These get conflated constantly, and they are not the same thing. A retry assumes the same call will work if you try again — it is for transient failures. A fallback assumes the same call will keep failing, so it tries a different one.

Python
import anthropic, openai# The provider SDKs raise their own classes, not Python's TimeoutErrorTRANSIENT = (openai.APIConnectionError, openai.RateLimitError, openai.InternalServerError,             anthropic.APIConnectionError, anthropic.RateLimitError,             anthropic.InternalServerError)# Retry: same model, transient failure (429, 5xx, timeout)robust = summarise.with_retry(    retry_if_exception_type=TRANSIENT,    wait_exponential_jitter=True,    stop_after_attempt=3,)# Fallback: a different model or a different strategyresilient = summarise.with_fallbacks([    summarise_on_claude,    RunnableLambda(lambda x: x["text"][:400] + "..."),   # extractive last resort])

Two details in that retry configuration matter more than they look. wait_exponential_jitter=True spreads retries randomly in time; without jitter, every client that failed during an outage retries at exactly the same instants, and the synchronised wave keeps the service down — a stampede your own retry logic created. And retry_if_exception_type is a whitelist for a reason: retrying a ValidationError or a 401 just burns three times the latency before failing anyway. Name the SDK's own exception classes in it: an OpenAI or Anthropic timeout is not a subclass of Python's built-in TimeoutError, so a whitelist of built-ins never matches and the retry silently does nothing.

One more layer to know about: the OpenAI and Anthropic clients already retry transient errors twice on their own (max_retries=2 by default). Stack three attempts on top and one bad minute becomes up to nine calls. Either lower the client's setting (init_chat_model(..., max_retries=0)) or keep the outer retry small.

The last fallback in that chain is worth copying: a non-LLM path. Truncating the text is a bad summary, but it is a response, and when the whole provider is down a bad response beats a 500.

Testing a pipeline

This is the payoff for all the structure, and the part most teams skip until an incident forces it.

Test the deterministic parts for real

Prompt templates and parsers involve no network and no randomness. Test them properly:

Python
def test_summary_prompt_includes_text():    msgs = summarise_prompt.format_messages(text="hello world")    assert "hello world" in msgs[0].contentdef test_json_parser_handles_fenced_output():    parser = JsonOutputParser()    assert parser.parse('```json\n{"people": ["Ada"]}\n```') == {"people": ["Ada"]}

That second test encodes a real production bug — a model wrapping JSON in a markdown fence — as a permanent regression check. Every incident you fix should leave a test behind like that.

Mock the model, not the pipeline

LangChain ships fakes for exactly this. They make tests fast, free and deterministic:

Python
from langchain_core.language_models.fake_chat_models import FakeListChatModeldef test_pipeline_shape():    fake = FakeListChatModel(responses=["A three sentence summary.",                                        '{"people": ["Ada"], "orgs": [], "dates": []}',                                        "positive"])    pipeline = build_pipeline(model=fake)          # model injected, not hard-coded    out = pipeline.invoke({"text": "Ada Lovelace joined the board."})    assert out["sentiment"] == "positive"    assert out["entities"]["people"] == ["Ada"]

The critical design decision is in the comment: build_pipeline(model=...) takes the model as an argument. If your module does model = init_chat_model(...) at import time, you cannot substitute a fake without patching, and every test hits the network. Inject the model. It costs one parameter and it is the difference between a test suite you run on every commit and one you run never.

Test error paths too, since those are the ones production will exercise:

Python
def test_entity_failure_does_not_kill_the_run():    class Boom(FakeListChatModel):        def _call(self, *a, **k): raise TimeoutError("upstream timeout")    out = build_pipeline(model=Boom(responses=[])).invoke({"text": "anything"})    assert out["entities"] == {"people": [], "orgs": [], "dates": []}

Integration tests: few, real, and pinned

Keep a small suite that makes genuine API calls — perhaps ten cases — and run it before releases rather than on every commit. Its job is not to check exact wording, which will change between model versions and break your build for no reason. Its job is to check properties:

Python
@pytest.mark.integrationdef test_entities_found_in_real_call():    out = production_pipeline.invoke({"text": ADA_BIO})    assert "Ada Lovelace" in out["entities"]["people"]      # property, not exact string    assert 30 <= len(out["summary"].split()) <= 120         # range, not equality    assert out["sentiment"] in {"positive", "negative", "neutral"}

Assert properties, not exact outputs. A test that fails whenever the model rephrases a sentence teaches your team to ignore failing tests.

Observability: knowing what actually happened

When a pipeline misbehaves in production, you need the input, the output, the latency and the token count for each component. Callbacks give you hooks at every stage:

Python
import timefrom langchain_core.callbacks import BaseCallbackHandlerclass PipelineMonitor(BaseCallbackHandler):    def __init__(self):        self.started = {}        self.tokens_in = self.tokens_out = 0    def on_llm_start(self, serialized, prompts, *, run_id, **kw):        self.started[run_id] = time.perf_counter()    def on_llm_end(self, response, *, run_id, **kw):        elapsed = time.perf_counter() - self.started.pop(run_id, time.perf_counter())        message = response.generations[0][0].message        usage = getattr(message, "usage_metadata", None) or {}        self.tokens_in  += usage.get("input_tokens", 0)        self.tokens_out += usage.get("output_tokens", 0)        logging.info("llm ok in %.2fs, %s in / %s out",                     elapsed, usage.get("input_tokens"), usage.get("output_tokens"))    def on_llm_error(self, error, *, run_id, **kw):        logging.error("llm failed: %s", error)monitor = PipelineMonitor()result = analyse.invoke({"text": doc}, config={"callbacks": [monitor]})print(monitor.tokens_in, monitor.tokens_out)

The token counts come from usage_metadata, the usage record LangChain attaches to every chat model reply in the same shape for every provider. Older examples read response.llm_output["token_usage"], which only OpenAI's integration fills in; on other providers it silently counts zero.

Two things make this useful rather than noisy. Callbacks are passed through config, so they apply to every nested component automatically — you attach one handler at the top and get events from the whole tree. And run_id is what lets you match a start to its end when calls overlap in a parallel branch; keying your timers on anything else gives you nonsense durations under concurrency.

Add a per-request identifier and the logs become answerable rather than merely voluminous:

Python
result = analyse.invoke(    {"text": doc},    config={"callbacks": [monitor],            "run_name": "document-analysis",            "tags": ["prod", f"tenant:{tenant_id}"],            "metadata": {"request_id": request_id, "doc_id": doc.id}},)

Assembling the production shape

Put the pieces together and the whole pipeline is legible in one screen — which is the actual test of whether the modularity worked:

Python
def build_pipeline(model, fast_model):    summarise = SUMMARY_PROMPT | model | StrOutputParser()    entities  = ENTITY_PROMPT  | model | JsonOutputParser()    sentiment = SENTIMENT_PROMPT | fast_model | StrOutputParser()    return (        RunnableLambda(validate_input)        | RunnableParallel(            summary=safe(summarise.with_retry(stop_after_attempt=3), "", "summary"),            entities=safe(entities.with_retry(stop_after_attempt=2), EMPTY, "entities"),            sentiment=safe(sentiment, "neutral", "sentiment"),            source=RunnablePassthrough(),          )        | RunnableLambda(assemble_record)    )

Read it top to bottom: validate, then three independent analyses concurrently, each retried and each degrading to a sensible default, then assemble the record. Nothing is hidden. Compare that with the 200-line function at the start, where the same information existed only as the order of statements.

Where people get this wrong

Splitting into components but keeping the fusion. Ten functions that each call the next one directly is not a pipeline — it is the same fused code with more files. The test is whether the composition is written down in one place as an expression. If finding the order of steps means reading ten function bodies, you have not modularised anything.

Over-decomposing. The opposite failure. A component per prompt line, wired through five layers of adapters, is harder to read than the fused version. The right granularity is "the smallest thing you would want to test, swap or reuse alone". If a piece has never been tested alone and never will be, it does not need to be a piece.

Retrying everything. Blanket retries on all exceptions turn a fast, clear 400 into a slow, confusing 400 — and during an outage, a synchronised retry storm from your own clients extends it. Retry transient network failures and rate limits; fail fast on everything else.

Hard-coding the model at module scope. The single most common reason an LLM codebase has no tests. model = init_chat_model(...) at the top of the file means every import needs an API key and every test costs money. Pass models in.

Treating parallel as free. It is free in cost and cheap in latency, but three concurrent calls consume three times the rate-limit budget in the same instant. A pipeline that runs comfortably at ten requests a minute sequentially can start returning 429s the day you parallelise it. Know your limits before you fan out.

What this buys you when the pager goes off

The argument for modular pipelines is usually made on elegance. The real argument is operational, and it shows up in three concrete moments.

When a prompt needs to change. Someone reports that summaries are too long. In the fused version you edit a string inside a 200-line function, cannot test the change in isolation, and deploy hoping nothing else moved. In the modular version you edit one prompt, run that component's tests in two seconds, and know exactly what you touched.

When costs jump. The monitor gives you tokens per component. You look and find entity extraction is 71% of your spend because it sends the full document while the summary sends a truncated one. That is a five-minute fix — but only because the number existed at component granularity. A single total tells you the bill went up and nothing else.

When a provider degrades. Latency triples. With retries, fallbacks and per-component degradation already in place, the pipeline gets slower and some enrichments go missing, and users still get answers. Without them, every request fails and you are writing resilience code during an incident, which is the worst possible time to write it.

Build the seams first. Retrofitting them into a function that already works is the hardest refactor in this whole field, because the function works and nobody wants to be the person who broke it.