Beyond LangChain: Production AI Orchestration Patterns

Why LangChain struggles in production and the orchestration patterns that work better for reliable, observable, and maintainable AI workflows at scale.

#ai#orchestration#langchain#production
Cover image for the article: Beyond LangChain: Production AI Orchestration Patterns

LangChain got us started, but production AI systems need orchestration patterns that prioritize reliability, observability, and debuggability over developer convenience. After migrating three production systems away from LangChain and building custom orchestration for 15+ AI workflows, here's why the framework approach fails at scale and what to build instead.

The Problem: Framework Complexity vs Production Reality

LangChain and similar frameworks provide high-level abstractions that make prototyping fast. But in production, those abstractions become liabilities:

  • Debugging opacity: When a chain fails, the stack trace traverses dozens of framework-internal methods. Finding the actual failure point requires deep framework knowledge.
  • Version fragility: Framework updates frequently break existing chains. Pinning versions means missing security patches and model provider updates.
  • Hidden retries and fallbacks: Framework-level retry logic conflicts with application-level retry policies, causing cascading failures.
  • Observability gaps: Generic tracing captures chain execution but not the semantic meaning of each step. "Chain step 4 failed" tells you nothing about business impact.
  • Testing difficulty: Mocking framework internals for unit tests requires intimate knowledge of implementation details that change between versions.

The alternative isn't "no orchestration" — it's orchestration designed for production from the start.

Architecture: Composable Step-Based Orchestration

The pattern that works: explicit steps with typed inputs/outputs, built-in observability, and deterministic control flow. No magic, no hidden state, no framework coupling.

AI Orchestration Architecture

Type-Safe Workflow Engine

The workflow engine enforces type safety between steps, provides built-in tracing, and keeps control flow explicit.

import asyncio
import time
import uuid
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import TypeVar, Generic, Any, Optional
from enum import Enum

T = TypeVar("T")
U = TypeVar("U")

class StepStatus(Enum):
    PENDING = "pending"
    RUNNING = "running"
    COMPLETED = "completed"
    FAILED = "failed"
    SKIPPED = "skipped"
    RETRYING = "retrying"

@dataclass
class StepResult(Generic[T]):
    status: StepStatus
    output: Optional[T]
    error: Optional[str]
    duration_ms: float
    retries: int
    metadata: dict = field(default_factory=dict)

@dataclass
class WorkflowContext:
    workflow_id: str
    trace_id: str
    step_results: dict[str, StepResult] = field(default_factory=dict)
    metadata: dict = field(default_factory=dict)
    started_at: float = field(default_factory=time.time)

class Step(ABC, Generic[T, U]):
    def __init__(self, name: str, max_retries: int = 2, timeout_ms: float = 30000):
        self.name = name
        self.max_retries = max_retries
        self.timeout_ms = timeout_ms

    @abstractmethod
    async def execute(self, input_data: T, context: WorkflowContext) -> U:
        pass

    async def run(self, input_data: T, context: WorkflowContext) -> StepResult[U]:
        retries = 0
        last_error = None

        while retries <= self.max_retries:
            start = time.time()
            try:
                result = await asyncio.wait_for(
                    self.execute(input_data, context),
                    timeout=self.timeout_ms / 1000,
                )
                duration = (time.time() - start) * 1000
                step_result = StepResult(
                    status=StepStatus.COMPLETED,
                    output=result,
                    error=None,
                    duration_ms=duration,
                    retries=retries,
                )
                context.step_results[self.name] = step_result
                return step_result

            except asyncio.TimeoutError:
                last_error = f"Timeout after {self.timeout_ms}ms"
                retries += 1
            except Exception as e:
                last_error = str(e)
                retries += 1
                if retries <= self.max_retries:
                    await asyncio.sleep(2 ** retries * 0.1)

        duration = (time.time() - start) * 1000
        step_result = StepResult(
            status=StepStatus.FAILED,
            output=None,
            error=last_error,
            duration_ms=duration,
            retries=retries - 1,
        )
        context.step_results[self.name] = step_result
        return step_result


class Workflow:
    def __init__(self, name: str):
        self.name = name
        self.steps: list[tuple[Step, Optional[str]]] = []  # (step, depends_on)

    def add_step(self, step: Step, depends_on: Optional[str] = None) -> "Workflow":
        self.steps.append((step, depends_on))
        return self

    async def run(self, initial_input: Any) -> WorkflowContext:
        context = WorkflowContext(
            workflow_id=f"{self.name}-{uuid.uuid4().hex[:8]}",
            trace_id=uuid.uuid4().hex,
        )

        current_input = initial_input
        for step, depends_on in self.steps:
            if depends_on and depends_on in context.step_results:
                dep_result = context.step_results[depends_on]
                if dep_result.status == StepStatus.FAILED:
                    # Skip steps that depend on failed steps
                    context.step_results[step.name] = StepResult(
                        status=StepStatus.SKIPPED,
                        output=None,
                        error=f"Dependency {depends_on} failed",
                        duration_ms=0,
                        retries=0,
                    )
                    continue
                current_input = dep_result.output

            result = await step.run(current_input, context)
            if result.status == StepStatus.FAILED:
                break
            current_input = result.output

        return context


# Example: RAG workflow with explicit steps
class RetrieveContextStep(Step[str, list[dict]]):
    def __init__(self, retriever):
        super().__init__("retrieve_context", max_retries=2, timeout_ms=5000)
        self.retriever = retriever

    async def execute(self, query: str, context: WorkflowContext) -> list[dict]:
        results = await self.retriever.search(query, top_k=5)
        context.metadata["retrieval_count"] = len(results)
        return results


class GenerateResponseStep(Step[list[dict], str]):
    def __init__(self, llm_client):
        super().__init__("generate_response", max_retries=1, timeout_ms=15000)
        self.llm = llm_client

    async def execute(self, context_docs: list[dict], context: WorkflowContext) -> str:
        query = context.metadata.get("original_query", "")
        prompt = self._build_prompt(query, context_docs)
        response = await self.llm.generate(prompt)
        return response

    def _build_prompt(self, query: str, docs: list[dict]) -> str:
        context_text = "\n\n".join(d.get("content", "") for d in docs)
        return f"Context:\n{context_text}\n\nQuestion: {query}\n\nAnswer:"

Observable Workflow with Structured Tracing

Every production AI workflow needs tracing that captures not just timing but semantic context: what was retrieved, what prompt was sent, what the model returned, and why decisions were made.

interface TraceSpan {
  spanId: string;
  parentSpanId?: string;
  name: string;
  startTime: number;
  endTime?: number;
  status: 'ok' | 'error';
  attributes: Record<string, string | number | boolean>;
  events: Array<{ name: string; timestamp: number; attributes: Record<string, unknown> }>;
}

interface WorkflowTrace {
  traceId: string;
  workflowName: string;
  spans: TraceSpan[];
  input: unknown;
  output: unknown;
  totalDurationMs: number;
  stepCount: number;
  errorCount: number;
}

interface StepConfig {
  name: string;
  maxRetries: number;
  timeoutMs: number;
  circuitBreaker?: {
    failureThreshold: number;
    resetTimeMs: number;
  };
}

class ObservableWorkflowEngine {
  private traces: Map<string, WorkflowTrace> = new Map();
  private circuitStates: Map<string, { failures: number; lastFailure: number; open: boolean }> = new Map();

  async executeStep<TInput, TOutput>(
    config: StepConfig,
    input: TInput,
    fn: (input: TInput) => Promise<TOutput>,
    traceId: string
  ): Promise<{ output: TOutput; span: TraceSpan }> {
    // Check circuit breaker
    if (config.circuitBreaker) {
      const state = this.circuitStates.get(config.name);
      if (state?.open) {
        const timeSinceFailure = Date.now() - state.lastFailure;
        if (timeSinceFailure < config.circuitBreaker.resetTimeMs) {
          throw new Error(`Circuit breaker open for ${config.name}`);
        }
        state.open = false; // Half-open: try one request
      }
    }

    const span: TraceSpan = {
      spanId: crypto.randomUUID(),
      name: config.name,
      startTime: Date.now(),
      status: 'ok',
      attributes: {
        'step.max_retries': config.maxRetries,
        'step.timeout_ms': config.timeoutMs,
      },
      events: [],
    };

    let lastError: Error | undefined;
    for (let attempt = 0; attempt <= config.maxRetries; attempt++) {
      try {
        span.events.push({
          name: attempt === 0 ? 'attempt_start' : 'retry_start',
          timestamp: Date.now(),
          attributes: { attempt },
        });

        const output = await Promise.race([
          fn(input),
          new Promise<never>((_, reject) =>
            setTimeout(() => reject(new Error('Timeout')), config.timeoutMs)
          ),
        ]);

        span.endTime = Date.now();
        span.attributes['step.attempts'] = attempt + 1;
        span.attributes['step.duration_ms'] = span.endTime - span.startTime;

        // Reset circuit breaker on success
        if (config.circuitBreaker) {
          this.circuitStates.set(config.name, { failures: 0, lastFailure: 0, open: false });
        }

        return { output, span };
      } catch (error) {
        lastError = error as Error;
        span.events.push({
          name: 'error',
          timestamp: Date.now(),
          attributes: { error: lastError.message, attempt },
        });

        if (attempt < config.maxRetries) {
          await new Promise(resolve => setTimeout(resolve, Math.pow(2, attempt) * 100));
        }
      }
    }

    // All retries exhausted
    span.endTime = Date.now();
    span.status = 'error';
    span.attributes['step.error'] = lastError?.message ?? 'Unknown error';

    // Update circuit breaker
    if (config.circuitBreaker) {
      const state = this.circuitStates.get(config.name) || { failures: 0, lastFailure: 0, open: false };
      state.failures++;
      state.lastFailure = Date.now();
      if (state.failures >= config.circuitBreaker.failureThreshold) {
        state.open = true;
      }
      this.circuitStates.set(config.name, state);
    }

    throw lastError;
  }
}

Orchestration Patterns for Production AI

Fan-Out/Fan-In for Multi-Source Retrieval

Query multiple knowledge sources in parallel, merge results, and pass to the LLM. Each source has independent timeout and fallback behavior.

Conditional Branching Based on Classification

Classify the input first, then route to specialized sub-workflows. A customer support system routes billing questions to one workflow and technical questions to another, each with different context and prompts.

Streaming Aggregation

For multi-step workflows that generate streaming output, aggregate partial results from upstream steps and stream them to downstream consumers without buffering the full intermediate result.

Checkpoint and Resume

Long-running workflows (document processing, batch analysis) checkpoint state after each step. If a step fails, the workflow resumes from the last checkpoint rather than restarting from scratch.

Comparison: Framework vs Custom Orchestration

DimensionLangChain/FrameworkCustom Orchestration
Time to prototype2-4 hours1-2 days
Time to production2-4 weeks (fighting the framework)1-2 weeks
Debugging easeLow (opaque stack traces)High (explicit flow)
ObservabilityGeneric spansSemantic traces
TestingDifficult (mock internals)Straightforward (mock interfaces)
Upgrade riskHigh (breaking changes)Low (you control the code)
PerformanceOverhead from abstractionsMinimal overhead
Team onboardingLearn framework + AILearn AI only

Benchmarks: Production Orchestration Performance

MetricLangChain (before)Custom (after)
Avg workflow latency2,400ms1,100ms
P99 workflow latency8,200ms3,400ms
Error rate3.2%0.8%
Mean time to debug failures45 min8 min
Test coverage achievable40%92%
Deploy confidenceLowHigh

When Frameworks Still Make Sense

Frameworks are appropriate when:

  • You're prototyping and expect to rewrite before production
  • The workflow is simple (single LLM call with retrieval)
  • Your team has strong framework expertise and accepts the maintenance burden
  • You need many integrations and don't want to build connectors

Frameworks are inappropriate when:

  • Reliability and observability are primary requirements
  • You need custom retry, fallback, and circuit breaker logic
  • The workflow has complex branching and conditional logic
  • You need to unit test individual steps in isolation
  • Multiple teams contribute to the same workflow

Conclusion

Production AI orchestration needs explicit control flow, typed step boundaries, built-in observability, and standard error handling — not framework magic. The extra effort to build custom orchestration pays for itself the first time you need to debug a production failure at 2 AM. Start with a simple Step abstraction with typed inputs/outputs, add tracing from day one, and compose complex workflows from tested, observable primitives. Your future self will thank you when the workflow breaks and you can identify the exact failure point in seconds, not hours.

Comments

    No comments yet. Be the first to share your thoughts.