Production-Grade LLM Annotation Pipeline: Pydantic, PostgreSQL & Kafka Patterns

The article details how a SaaS team evolved a fragile LLM ticket-classification script into a production-grade pipeline using strict Pydantic contracts, PostgreSQL state machines with atomic leasing, Kafka outbox patterns, immutable training snapshots, and metric-gated rollouts to handle model drift, replay, and human review safely.

Ray's Galactic Tech
Ray's Galactic Tech
Ray's Galactic Tech
Production-Grade LLM Annotation Pipeline: Pydantic, PostgreSQL & Kafka Patterns

Incident Retrospective: How Silent Errors Entered Training Candidates in 30 Minutes

Background and Legacy Implementation

A subscription SaaS support team receives ~12,000 tickets daily. They wanted an LLM to pre-label tickets as billing, technical, or other for routing and training-set inclusion. The initial implementation was a single script that called the model, extracted JSON via regex, and wrote to the database. It achieved ~99.6% parseability but lacked robustness against model switches, Kafka replays, and human review.

# Legacy implementation (not for production)
payload = re.search(r"\{.*\}", raw_response, re.S).group()
label = json.loads(payload)
confidence = float(label.get("confidence", 0))
category = label.get("category", "other")

This code silently coerced string confidence to float and defaulted missing categories to other, meaning "parse success" did not guarantee contract compliance or label correctness.

Timeline of the Incident

09:00 – 10% traffic canary to new model; no alerts.

09:12 – New model starts prefixing output with "Here is the classification result:"; regex still extracts JSON, monitoring sees nothing.

09:18 – confidence field becomes a string; float() silently converts.

09:30 – Raw JSON/Schema first-pass failure rate jumps from 0.4% to 6.8%; legacy dashboard still shows "parse success". New pipeline freezes auto-candidate promotion.

09:36 – Two workers restart, Kafka messages replay; legacy code causes duplicate LLM calls and double billing. New pipeline absorbs replay via unique keys and leases, isolates anomalies.

10:05 – Replay fixed golden set; cannot distinguish old vs new prompt results. New pipeline compares by pipeline_version and rolls back.

The incident exposed four gaps: no strict output contract, no distinguishable error codes, no idempotent state machine for external calls, and no "quality gate failure stops auto-promotion" control plane.

Define Identity, State, and Invariants First

Data Identities

Input Identity : (record_id, input_version) – allows re-annotation when ticket body updates.

Content Fingerprint : SHA-256 of raw UTF-8 text – prevents confusion when upstream reuses version numbers.

Annotation Identity : (record_id, input_version, pipeline_version) – ensures exactly one business result per replay.

Pipeline Version : immutable tuple of model + prompt + schema + route – enables side-by-side comparison and rollback.

Snapshot Identity : snapshot_id + member-list hash – makes training sets reproducible and immune to later review changes.

State Machine

The state machine (illustrated in the article) enforces:

Untrusted model responses must pass server-side schema before routing; no "extract JSON then best-effort repair".

Same annotation identity has exactly one current business result; does not promise LLM called only once.

Only the worker holding the current attempt_id may commit results; late workers cannot overwrite renewed leases or re-claimed results.

Snapshots record concrete members, not a mutable "all records before export time" query.

Boundary Validation: Strict Accept, Limited Correction

Model output is untrusted data from outside the network boundary. The contract rejects extra fields, string confidence, blank reasons, and NaN/Infinity. model_validate_json() requires the entire response to be JSON, so it will not accept "Here is the result: {...}".

from typing import Literal
from pydantic import BaseModel, ConfigDict, Field, ValidationError, field_validator

MAX_OUTPUT_BYTES = 8_192

class Label(BaseModel):
    model_config = ConfigDict(strict=True, extra="forbid")
    category: Literal["billing", "technical", "other"]
    confidence: float = Field(ge=0, le=1, allow_inf_nan=False)
    reason: str = Field(min_length=1, max_length=300)

    @field_validator("reason")
    @classmethod
    def non_blank(cls, value: str) -> str:
        if not value.strip():
            raise ValueError("reason is blank")
        return value.strip()

def parse_label(raw: str) -> Label:
    if len(raw.encode("utf-8")) > MAX_OUTPUT_BYTES:
        raise ValueError("response_too_large")
    return Label.model_validate_json(raw)

def classify_parse_error(exc: Exception) -> str:
    if isinstance(exc, ValueError) and str(exc) == "response_too_large":
        return "response_too_large"
    if isinstance(exc, ValidationError):
        types = {e["type"] for e in exc.errors()}
        if "json_invalid" in types:
            return "json_invalid"
        if "extra_forbidden" in types:
            return "schema_extra_field"
        if any(t.endswith("_type") for t in types):
            return "schema_type_error"
        return "schema_validation_error"
    return "unexpected_parse_error"

On format failure, allow exactly one targeted regeneration : prompt the model to return only the specified three-field JSON object. Do not stuff full exceptions, ticket bodies, or other user data into correction prompts or logs. If the second attempt still fails, move to quarantine; do not fabricate other. Network timeouts, 429, 5xx, and auth errors are a separate exception class that must be handled explicitly on the client side, not mixed into "JSON format error" retries.

Production Implementation: PostgreSQL Atomic Claim and Final Write

SQLite is fine for demos, but multi-worker Kafka consumption requires PostgreSQL. The schema (indexes, permissions, partitioning omitted) defines an annotation_status enum and an annotation_tasks table with a composite primary key (record_id, input_version, pipeline_version). The encrypted_text_ref points to controlled object storage or an encrypted column; keeping the content hash and version in the task table enables consistency checks while avoiding plaintext exposure in routine queries.

Claim Is Not "Select Then Act"

A "SELECT, if missing call model, then INSERT" race lets two consumers both pass the check and invoke the model twice. The correct flow is claim first, then call. Two short DB transactions (no model call inside the transaction):

-- Step 1: first message tries to create a processing task.
INSERT INTO annotation_tasks (
  record_id, input_version, pipeline_version, input_sha256, source,
  encrypted_text_ref, status, lease_until, attempt_id, attempts
) VALUES (
  :record_id, :input_version, :pipeline_version, :input_sha256, :source,
  :text_ref, 'processing', now() + interval '90 seconds', :attempt_id, 1
)
ON CONFLICT (record_id, input_version, pipeline_version) DO NOTHING
RETURNING attempt_id;

-- Step 2: take over only when old lease expired. RETURNING a row permits model call.
UPDATE annotation_tasks
SET lease_until = now() + interval '90 seconds',
    attempt_id = :attempt_id,
    attempts = attempts + 1,
    updated_at = now()
WHERE record_id = :record_id
  AND input_version = :input_version
  AND pipeline_version = :pipeline_version
  AND status = 'processing'
  AND lease_until <= now()
RETURNING attempt_id;

If both INSERT and UPDATE return no rows, read the task: input_sha256 differs from current input → isolate as input_version_conflicts_with_content, alert upstream to fix version strategy; do not overwrite.

Status processing and lease not expired → another worker is running, skip model call.

Status is a business terminal state → Kafka replay, skip model call.

Final write must include the attempt_id condition. If the call exceeds 90 seconds without lease renewal, another worker may have taken over; the late worker's result is discarded and logged as lease_lost.

UPDATE annotation_tasks
SET status = :status,
    label_json = :label_json::jsonb,
    error_code = :error_code,
    review_reason = :review_reason,
    lease_until = NULL,
    updated_at = now()
WHERE record_id = :record_id
  AND input_version = :input_version
  AND pipeline_version = :pipeline_version
  AND status = 'processing'
  AND attempt_id = :attempt_id;

This still does not guarantee "model called exactly once": the process may exit after model returns but before final UPDATE. If the vendor supports idempotency keys, send a stable business key or attempt_id and store the vendor request ID to reduce duplicate billing; otherwise accept at-least-once and constrain with cost monitoring and retry budgets.

Kafka and Outbox: Make Non-Atomic Boundaries Explicit

Kafka offsets, external LLM calls, and PostgreSQL are not a single atomic transaction. Reliable consumer pseudocode:

for message in poll_records():
    try:
        record = InputRecord.model_validate_json(message.value)
    except ValidationError as exc:
        save_invalid_input(message, error_code="input_schema_invalid")
        mark_offset_done(message)  # only after isolation table commit
        continue

    claim = claim_task(record)  # short DB tx; returns claimed / terminal / in_progress
    if claim.state == "terminal":
        mark_offset_done(message)
        continue
    if claim.state == "in_progress":
        defer_without_committing_gap(message)
        continue

    outcome = call_llm_with_deadline(record, idempotency_key=claim.attempt_id)
    with db_transaction():
        finish_only_if_current_attempt(claim, outcome)
        insert_outbox_event_if_finished(claim, outcome)
        mark_offset_done(message)

commit_only_highest_contiguous_done_offset_per_partition()
mark_offset_done

is not an immediate commit(): with concurrent per-partition processing, only the highest contiguously completed offset may be committed. Later messages finishing first cannot leapfrog earlier unfinished ones. On DB failure, pause the affected partition with jitter backoff; do not commit offsets, avoiding infinite hot loops.

Downstream publishing uses a transactional outbox. Task terminal state and outbound events are written in one PostgreSQL transaction; the publisher can redeliver, and consumers deduplicate by event_id.

CREATE TABLE outbox_events (
  event_id uuid PRIMARY KEY,
  aggregate_key text NOT NULL,
  event_type text NOT NULL,
  payload jsonb NOT NULL,
  created_at timestamptz NOT NULL DEFAULT now(),
  published_at timestamptz,
  publish_attempts integer NOT NULL DEFAULT 0
);

-- Publisher batch-claims to avoid duplicate sends.
SELECT event_id, payload
FROM outbox_events
WHERE published_at IS NULL
ORDER BY created_at
FOR UPDATE SKIP LOCKED
LIMIT 100;

Set published_at after successful send. If a restart occurs between "send succeeded" and "mark published", the event is resent, so consumers must still deduplicate by event_id. DLQ/isolation records store topic, partition, offset, record ID, error code, replay count, and replay policy; never put raw sensitive bodies into the DLQ.

Testing and Acceptance: Verify Bad Paths, Not Just Happy Path

The following minimum automated acceptance suite should run before release. It does not replace load testing, chaos drills, or human quality evaluation, but prevents the most common regressions.

String confidence : input "confidence":"0.92" → first pass schema_type_error; must not become auto-candidate.

JSON prefix : input "Explanation: {...}" → json_invalid; no brace extraction.

Extra field : input includes debug_trace → schema_extra_field.

Two workers concurrent : same business key claimed simultaneously → at most one gets claimed.

Lease expiry : worker A does not commit, lease expires → B can take over; A's final write affects 0 rows.

Version conflict : same ID/version, different content hash → reject processing and alert.

Review rejection : review_pending → rejected → record absent from snapshot members.

Snapshot immutability : modify review label after snapshot → old snapshot's members and manifest unchanged.

Pydantic layer can be covered directly with pytest:

import pytest
from pydantic import ValidationError

def test_does_not_coerce_string_confidence():
    with pytest.raises(ValidationError):
        parse_label('{"category":"billing","confidence":"0.92","reason":"duplicate charge"}')

def test_does_not_extract_json_from_prose():
    with pytest.raises(ValidationError):
        parse_label('以下是结果:{"category":"billing","confidence":0.92,"reason":"duplicate charge"}')

def test_rejects_unknown_field():
    with pytest.raises(ValidationError):
        parse_label('{"category":"billing","confidence":0.92,"reason":"duplicate charge","debug":true}')

Concurrency and lease tests must run against a real PostgreSQL test database; mocks cannot replace transaction semantics. In practice, spin up two independent connections, use a barrier to make them execute the claim SQL simultaneously, assert only one attempt_id is returned; then use an injectable clock or very short test lease to verify takeover and late updates. Snapshot tests should read snapshot_members, not re-run "currently eligible tasks" queries.

Freeze Training Snapshots: Save Members, Not a Time Condition

Saving only a cutoff_at and later querying "all reviewed data up to that time" is unreliable: backfills, backdated corrections, and review amendments change results. When creating a snapshot, write snapshot metadata and member list in a single transaction.

CREATE TABLE dataset_snapshots (
  snapshot_id text PRIMARY KEY,
  pipeline_version text NOT NULL,
  policy_version text NOT NULL,
  manifest_sha256 text NOT NULL,
  created_at timestamptz NOT NULL DEFAULT now()
);

CREATE TABLE snapshot_members (
  snapshot_id text NOT NULL REFERENCES dataset_snapshots(snapshot_id),
  record_id text NOT NULL,
  input_version integer NOT NULL,
  pipeline_version text NOT NULL,
  final_category text NOT NULL,
  PRIMARY KEY (snapshot_id, record_id, input_version, pipeline_version)
);

-- Run in one transaction: select only auto-candidate and human-reviewed-passed data.
INSERT INTO snapshot_members (snapshot_id, record_id, input_version, pipeline_version, final_category)
SELECT :snapshot_id, record_id, input_version, pipeline_version,
CASE WHEN status = 'reviewed' THEN reviewer_label
ELSE label_json->>'category' END
FROM annotation_tasks
WHERE pipeline_version = :pipeline_version
  AND status IN ('auto_candidate', 'reviewed');

Then export JSONL sorted by members, compute SHA-256 of exported bytes, write hash back to dataset_snapshots, and store the file in versioned, immutable object storage. Handle "DB committed but object storage write failed" by marking snapshot materializing, a retryable materialization job writes and verifies manifest, then marks ready. Training jobs may only read ready snapshots.

Data-Driven Release Loop: Metrics, Gates, Canary, Rollback

The following gates show "how to decide from metrics". Thresholds must be calibrated from historical golden sets, business tolerance, and sample size; e.g., random-sample mislabel rate should consider confidence intervals, not just point estimates.

First-pass Schema Failure Rate : (first-pass parse failures / LLM calls). Gate: relative increase >2pp or >1% → freeze auto-candidates, investigate vendor/prompt.

Random Quality Pool Mislabel Rate : (human vs final label mismatches / random reviews completed). Gate: 95% upper confidence bound exceeds business threshold → demote auto-candidates to human review.

Duplicate Call Rate : LLM calls / unique annotation identities - 1. Gate: >0.5% → check leases, timeouts, vendor idempotency keys.

Backlog Recovery Time : time to return to normal backlog after peak. Gate: exceeds SLO → throttle, scale, or pause low-priority tasks.

Category Distribution Drift : PSI / category proportion change vs fixed baseline. Gate: exceeds preset threshold → replay fixed golden set and human spot-check.

Human Disagreement Rate : (double-review mismatches / double-review samples). Gate: rises for two consecutive windows → re-examine label guidelines and model rationales.

An executable canary process:

Offline : run old and new versions on fixed golden set, compare macro-F1, schema failure rate, cost by source and category; human-review disagreement samples.

Shadow traffic : new version processes real input but does not affect routing or enter training snapshots; record diff vs old version.

5% auto-candidate : keep 100% random quality review, observe at least one full business cycle.

25% / 50% / 100% : each step requires quality, cost, latency, backlog, duplicate calls all pass gates, with sufficient random review samples.

Rollback : any hard gate fails → immediately stop auto-candidate promotion for that pipeline_version; existing tasks move to review_pending or quarantine, old version continues serving. Training snapshots only created from gate-passed versions.

Dashboards must slice by pipeline_version, source, category, and time window: volume, claim success rate, first/second-pass schema error codes, LLM call volume and cost, P50/P95 latency, 429/5xx, lease expiries, Kafka lag, isolation volume, review turnaround, random quality mislabel rate, human disagreement. Aggregate metrics without version dimension are often why "everything looks fine" during an incident.

Privacy, Audit, and Traceability

Tickets contain PII and account data. Raw text belongs in encrypted, least-privilege data domains; application logs, DLQ, metric labels, and error summaries retain only metadata needed for location and replay. Execute de-identification before training export, and define retention, access approval, and deletion processes per data domain.

Traceability ≠ verbatim reproducibility. Record input version and content hash, code commit, model identifier, prompt version, schema, routing config, human review guidelines, request parameters, and controlled response summary; but the same model name does not guarantee identical output on re-inference. DVC or similar should track frozen snapshots and their build dependencies , not treat the continuously changing Kafka consumption stream as a reproducible stage.

Conclusion

The real danger in the incident was not that the model returned a preamble, but that the system disguised "contract violation" as "usable for training". High-quality LLM data engineering gives every failure class a deterministic destination: Pydantic guards the structural boundary, PostgreSQL state machine absorbs replays and concurrency, Kafka/outbox ensures downstream recoverability, human random samples provide evidence of label correctness, immutable snapshots make training data auditable. Thus, even when vendors change, workers crash, tickets are backfilled, or review rules evolve, the system still surfaces risk, limits blast radius, and provides a verifiable recovery path.

Original Source

Signed-in readers can open the original source through BestHub's protected redirect.

Sign in to view source
Republication Notice

This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactadmin@besthub.devand we will review it promptly.

data engineeringKafkaIdempotencyPostgreSQLannotationPydanticproduction readinessLLM data pipelinetraining snapshots
Ray's Galactic Tech
Written by

Ray's Galactic Tech

Practice together, never alone. We cover programming languages, development tools, learning methods, and pitfall notes. We simplify complex topics, guiding you from beginner to advanced. Weekly practical content—let's grow together!

0 followers
Reader feedback

How this landed with the community

Sign in to like

Rate this article

Was this worth your time?

Sign in to rate
Discussion

0 Comments

Thoughtful readers leave field notes, pushback, and hard-won operational detail here.