LangChain Mastery

Course Content

LangChain Mastery

7 sections · 109 lessons

How do you handle large-scale data in LangChain applications?


Labelling two million complaints without restarting from zeroRead a pageof 1,000 rowsabatch withmax_concurrency 40Rate limiter atprovider quotaSave results,record progressNext page; retryfailures laterreturn_exceptions=True keeps one failure from sinking the page.
Checkpointing after every page is what makes a 14-hour job safe to crash at 90 percent.

What you need to know

"Large-scale data" shows up in two places: indexing millions of documents, and processing many inputs through a model (classifying 1 million tickets, summarising 50,000 calls).

Indexing at scale

  1. Stream — loaders have lazy_load(), which yields documents one at a time instead of building a giant list.
  2. Split and embed in batches — embed_documents on a few hundred texts per call; vector stores' add_documents batch internally.
  3. Index incrementally — the index() function with a record manager skips unchanged chunks and cleans up old ones.
  4. Run offline — a scheduled job or queue worker, never in a web request.
  5. Store in a server database — pgvector, Qdrant, Milvus, OpenSearch.

Processing many inputs

Python
from langchain_core.rate_limiters import InMemoryRateLimiterfrom langchain.chat_models import init_chat_modellimiter = InMemoryRateLimiter(requests_per_second=20, max_bucket_size=20)llm = init_chat_model(settings.chat_model, rate_limiter=limiter)classify = prompt | llm.with_structured_output(TicketLabel)async def classify_all(tickets, chunk_size=1000):    for start in range(0, len(tickets), chunk_size):        part = tickets[start:start + chunk_size]        results = await classify.abatch(            [{"text": t.text} for t in part],            config={"max_concurrency": 20},            return_exceptions=True,          # one failure doesn't kill the batch        )        save_results(part, results)          # checkpoint after every chunk
  • abatch runs inputs concurrently; max_concurrency caps parallel calls.
  • InMemoryRateLimiter smooths requests per second within one process (for many processes, use a shared limiter or the provider's batch API).
  • return_exceptions=True returns exceptions in place of results, so you can retry just the failed items.
  • Saving after each chunk means a crash at 90% restarts from 90%, not zero.

For very large offline jobs, many providers offer a batch API at a discount with results in hours; it suits jobs that aren't urgent.

Cost control

Estimate before running: inputs × (average input tokens + output tokens) × price. Run on 1% first, check quality and cost, then run the rest. Use the smallest model that meets the quality bar for simple labelling.

A real-life example

A telecom company wants to label 2 million customer complaints from the past year by issue type, to plan network fixes.

A first attempt loops chain.invoke one by one: estimated 23 days. The rebuilt job reads complaints from the database in pages of 1,000, uses abatch with max_concurrency=40 and a rate limiter matching their provider quota, a small model with structured output, and writes results after each page. A pilot on 20,000 complaints showed 94% agreement with human labels and a projected cost the team approved. The full run finished in about 14 hours; 0.3% of items failed and were retried in a second pass.

Follow-up questions to expect

  • "batch vs a loop with threads?" — batch already uses a thread pool (and abatch uses asyncio) with concurrency control, and keeps callbacks and tracing working.
  • "How do you trace a million calls?" — Sample tracing (for example LANGSMITH_TRACING_SAMPLING_RATE) or disable it for bulk jobs and log metrics instead.
  • "How do you handle documents bigger than the context window?" — Split, process chunks, then combine (map-reduce), or retrieve only the relevant parts.