Go · Kafka · Clean Architecture

Data Indexing Pipeline

An event-driven Go consumer that retired slow Postgres views and materialized views in favor of a Kafka-driven pipeline, keeping Elasticsearch in sync with Postgres for a legal-document platform — no HTTP API, just three topics in and a search index kept honest on the way out.

Why we built it?

Search indexing used to run off a set of Postgres views and materialized views — some of them large, CTE-based queries that had accumulated years of special cases nobody fully trusted touching. Assembling a single document type could mean dozens of separate database round trips per record: fast individually, expensive multiplied across tens of thousands of them. And a multi-hour bulk reindex against a search index serving live traffic meant choosing between degraded search for hours or an outage — there wasn't a third option.

How we built it?

The rewrite moved one document type at a time, behind its own feature flag defaulting to off. A small diffing tool ran the same record through both the legacy view and the new Go processor and reported every field-level difference; legacy logic stayed the source of truth until that tool reported zero discrepancies against real data, not a handful of hand-picked examples. The clean-architecture layering shown above kept the migration's complexity contained — swapping the read path for one document type never touched the other ten.

Infrastructure: adapters to everything external Application: handler → processor → repository Domain: pure interfaces, no dependencies 01 · TRIGGER 03 · DESTINATION document.post.scraping workers × 2 document.indexing core path · workers × 3 upsert · bulk reindex · delete document.importing workers × 2 Event · Repository interfaces Handler per-topic dispatch Processor routes by type suffix Repository Postgres ↔ Elasticsearch upsert bulk delete Postgres views (legacy) diff ✓ Internal Tool · PHP API scraping-ingest only PostgreSQL source of truth · read Elasticsearch write target 2 sub-indices → Redis locking · checkpoints Slack best-effort alerts main index (shared) quick-search index shadow index (building) alias → content ❄ merge write ✓ full rewrite
Three Kafka topics feed a clean-architecture Go consumer that reads PostgreSQL and writes Elasticsearch, with side connections to Redis for checkpointing, Slack for alerts, and an internal PHP API for the scraping path.

What changed after everything above?

  • Query count, not query speed, turned out to be the real design smell — batching duplicate lookups and running independent queries concurrently cut round trips by two-thirds for one document type, turning a multi-hour job into well under an hour.
  • Reindexing no longer means an outage: a shadow index builds in the background, then an atomic alias swap flips traffic to it instantly.
  • A metadata-only update skips the expensive content-assembly step entirely — a merge write, not a full rewrite.
  • Since the migration landed, the pipeline has scaled exponentially from processing 200K to 2M documents in under an hour.
← Back to Casket