The example is sectioning: it retrieves a fixed set of candidate document sections, asks the
model to answer from each one alone (never seeing the other sections or the other calls) and
combines whichever sections actually answered part of the question. Because each call already
knows which single section it saw, citations are exact by construction; nothing has to be parsed
back out of free text the way RAG and prompt chaining do.
ThreadPoolExecutor.map is what makes this parallel rather than sequential: it submits every
call to the pool at once, and the calls run concurrently, but the returned iterator still yields
results in the order the candidates were given, regardless of which call actually finishes
first. That is what keeps the example deterministic without an explicit sort: order comes from
retrieval, never from a race between threads.
Every tracer.record call happens on the main thread, after list(pool.map(...)) has already
collected every result: the trace itself is never written to from more than one thread at once.
A StubModel built from a fixed list of canned responses is not safe to call from several
threads concurrently, since it advances a shared counter with no lock; the example’s tests build
their stub from a function that reads the prompt instead, which has no shared state to race on.
The trace above shows a real cost of naive sectioning: two of the three candidate sections
happened to answer the same sub-question, so the combined text states the filter fact twice.
Nothing in a fixed combine step notices the overlap, because each section answered without
seeing what the others said.
examples/parallelization/run.py · lines 22–75
LEVEL = 3
CANDIDATES_K = 3
NO_ANSWER = "NOT IN THIS SECTION"
PER_SECTION_SYSTEM = (
"You are given exactly one source passage about Halvorsen appliances, and a question that "
"may have more than one part. If this passage answers all or part of the question, answer "
"briefly using only this passage. If it answers none of the question, reply with exactly "
f"'{NO_ANSWER}' and nothing else."
)
def _answer_from_one_section(question: str, section: Section, model: Model) -> Completion:
prompt = f"Passage [{section.cite}] {section.title}:\n{section.text}\n\nQuestion: {question}"
return model.complete([Message(role="system", content=PER_SECTION_SYSTEM), Message(role="user", content=prompt)], max_tokens=200)
def run(
question: str,
model: Model,
embedder: Embedder | None,
tracer: Tracer,
*,
corpus_dir: Path = DEFAULT_CORPUS_DIR,
k: int = CANDIDATES_K,
) -> Answer:
del embedder # candidates come from keyword search, not a vector index
sections = load_sections(corpus_dir)
candidates = [s for s, score in bm25_search(sections, question, k=k) if score > 0]
tracer.record(
kind="code", decided_by="code", title="Pick sections to answer in parallel",
detail=", ".join(s.cite for s in candidates) or "none",
)
# .map submits every call to the pool at once and yields results back in candidate order,
# so the calls run concurrently but the code below never has to sort them: determinism comes
# from retrieval order, not from whichever call happens to finish first.
with ThreadPoolExecutor(max_workers=max(1, len(candidates))) as pool:
completions = list(pool.map(lambda s: _answer_from_one_section(question, s, model), candidates))
for section, completion in zip(candidates, completions):
tracer.record(
kind="model", decided_by="code", title=f"Answer from {section.cite} alone", detail=completion.text[:200],
tokens_in=completion.tokens_in, tokens_out=completion.tokens_out, ms=completion.ms,
)
used = [(s, c) for s, c in zip(candidates, completions) if NO_ANSWER not in c.text.upper()]
tracer.record(
kind="code", decided_by="code", title="Combine the sections that answered",
detail=f"{len(used)} of {len(candidates)} sections answered part of the question",
)
if not used:
return Answer(text="None of the retrieved sections answered the question.", citations=[], retrieved_sources=[s.cite for s in candidates])
combined = " ".join(c.text.strip() for _, c in used)
return Answer(text=combined, citations=sorted({s.cite for s, _ in used}), retrieved_sources=[s.cite for s in candidates])
Run it yourself:
examples/parallelization/README.md · lines 16–16
python -m examples.parallelization --model stub:scripted
Real-time parallel calls like these are for when the answer is needed now. When it is not (a
nightly re-score of every open ticket, a one-time pass over a large document set), the batch
APIs three model makers publish do the same many-calls-one-submission idea asynchronously and
cheaper: Anthropic describes its Message Batches API as suited to tasks that do not need an
immediate response, “with most batches finishing in less than 1 hour while reducing costs by 50%
and increasing throughput”[4]; OpenAI’s Batch API gives a “50% cost discount compared to
synchronous APIs” with each batch completing “within 24 hours (and often more quickly)”[5]; Google states that its Gemini Batch API processes requests at “50% of the standard cost”
and says of the wait: “The target turnaround time is 24 hours, but in majority of cases, it is
much quicker”[6]. Each of those is the maker’s own published figure, checked on the
date in the source list below, and each is the same trade: give up the immediate response, halve
the price.