Elasticsearch Out of Sync With Database: Index From the WAL

Distributed Data · Advanced · 7 min read · published

This article was written by Claude (Anthropic) and published automatically.

What this solves: Your search index shows stale prices, deleted items, or missing rows because the app writes to the database and the search index separately. Here's the pipeline that makes them converge.

The Forces at Play

If your Elasticsearch index is out of sync with your database, the cause is almost always a dual write: the request handler commits to Postgres, then calls Elasticsearch. Those are two systems with no shared transaction, so every failure between them becomes permanent drift. Search shows a product you deleted, a price from last week, or no row at all for something that exists.

The obvious fixes don't close the gap:

The tension is that the database is the source of truth, but search needs a differently shaped, denormalized copy, and you can't make both writes atomic. The way out is to stop treating indexing as a side effect of a request. Instead, derive the index from the database's own commit log.

The Shape

Change Data Capture (CDC) reads committed changes from the write-ahead log, so only things that actually committed are ever indexed, in commit order.

flowchart LR
  App[App / API] -->|single write, one txn| PG[(Postgres primary)]
  PG -->|WAL via logical replication slot| DBZ[Debezium connector]
  DBZ -->|key = product_id, payload includes LSN| K[[Kafka topic: db.public.products]]
  K -->|partition 0..N, ordered per key| IDX[Indexer consumers]
  IDX -->|re-read current row + joins| RR[(Read replica)]
  IDX -->|bulk index, version = LSN, version_type=external_gte| ALIAS{{alias: products}}
  ALIAS --> V2[(index products_v2)]
  BF[Backfill job: snapshot scan] -.->|same message contract| K
  ALIAS -.->|atomic swap after rebuild| V3[(index products_v3)]
  Search[Search API] --> ALIAS

Three properties come from this structure. The app makes exactly one write. Ordering is guaranteed per document because the Kafka key is the primary key. Stale writes are rejected by Elasticsearch itself, because the LSN only goes up and is used as an external version.

How Data Flows Through It

Follow one price change for product 8812:

  1. The API runs UPDATE products SET price = 19.99 WHERE id = 8812 and commits. That's the only thing the request does.
  2. Postgres appends the change to the WAL at LSN 0/3A4F1C8. The logical replication slot holds the WAL until Debezium confirms it.
  3. Debezium decodes the row change and publishes it to db.public.products with key {"id":8812}. Every change to 8812 hashes to the same partition.
  4. The indexer consumer for that partition receives the event. It uses the event as a trigger rather than a payload. It re-reads product 8812 with its category and inventory joins from a replica, then builds the search document. Re-reading means one event can safely stand in for several coalesced ones, and the document always reflects a consistent row.
  5. It sends a bulk request versioned with the event's LSN:
{ "index": { "_index": "products", "_id": "8812",
             "version": 61134216, "version_type": "external_gte" } }
{ "name": "Trail Shoe", "price": 19.99, "category": "running", "in_stock": true }
  1. If a slower, older event for 8812 arrives later, perhaps from a replay, ES rejects it with a version conflict. The indexer treats that as success and commits the Kafka offset.
  2. Deletes arrive as a delete event plus a Kafka tombstone. The indexer issues a versioned delete.

If the replica you re-read from lags behind the event's LSN, check pg_last_wal_replay_lsn() against it and wait or retry. Otherwise you'll index pre-change state under a post-change version, and that state will stick.

What Each Piece Owns

Where It Breaks Down

When This Is Overkill

For a single service with modest write volume, a simpler design is usually right:

  1. Index after commit, best-effort.
  2. Run a reconciliation job every few minutes that re-indexes rows WHERE updated_at > last_run - interval '5 min', plus a nightly diff of IDs to catch deletes.

That bounds drift to minutes with no new infrastructure. If your search needs are modest, Postgres full-text search with tsvector and a GIN index removes the second store entirely.

You've outgrown the simple design when any of these become true:

At that point, move indexing onto the log.

Key takeaway: Never dual-write to your database and search index from app code: stream committed changes from the WAL, key them by document ID, and version every write with the LSN so late events can't overwrite newer data.

Real-world challenge

You run Debezium against Postgres to keep Elasticsearch in sync. On Monday morning the primary database alerts at 94% disk usage. Nothing about write traffic changed. The Kafka Connect cluster was down over the weekend after a failed deploy, and nobody noticed because search still served results, just slightly stale ones. Diagnose and fix it, and make sure it can't silently happen again.

Diagnosis. A logical replication slot pins WAL until its consumer confirms it. With Connect down, Postgres kept every WAL segment since Friday on the primary's disk.

SELECT slot_name, active,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained
FROM pg_replication_slots;

You'll see active = f and tens of GB retained.

Immediate fix. Bring Connect back first. The slot resumes from confirmed_flush_lsn and WAL is freed as the indexer catches up. If disk will run out before that happens, drop the slot (pg_drop_replication_slot). That loses your position, so you must then run a full snapshot plus alias-swap reindex.

Make it fail safe.