Series overview
Part 4 of 1822% complete
2026-05-09•12 min read

Search architecture and keeping Elasticsearch in sync

This chapter decides how changes get from PostgreSQL into Elasticsearch without losing, duplicating, or resurrecting profiles. You add two things to the lab: a row_version column that increases on every change, and a trigger-fed outbox table that records every committed change as an identifier and a version. You also reproduce the bug that sinks most home-made outbox pollers. The Kotlin indexer that drains this outbox is built in chapter 14; this chapter settles its design.

It builds on chapter 03, which showed that external versions reject stale writes, that delete tombstones expire after index.gc_deletes, and that bulk requests fail per item. The transactional outbox pattern itself has its own series on this site; this chapter covers only what changes when the consumer is a search index.

The architecture

PostgreSQL, source of truth

Application services

one transaction

trigger, same transaction

poll and claim batch

re-read current rows

bulk index and delete, external versions

search

Write API

Search API

user_profile

profile_search_outbox

Indexer

user-profile-write alias

user-profile-read alias

Elasticsearch

PostgreSQL, source of truth

Application services

one transaction

trigger, same transaction

poll and claim batch

re-read current rows

bulk index and delete, external versions

search

Write API

Search API

user_profile

profile_search_outbox

Indexer

user-profile-write alias

user-profile-read alias

Elasticsearch

The diagram shows two paths. On the write path, the Write API changes user_profile in one PostgreSQL transaction, and a trigger adds a row to profile_search_outbox in the same transaction. The indexer polls the outbox, claims a batch, re-reads the current state of those profiles from user_profile, and sends bulk index and delete operations with external versions to the user-profile-write alias. On the read path, the Search API queries the user-profile-read alias and never touches PostgreSQL for search. The two aliases, introduced in chapter 05, let the index behind them change without either application noticing.

Two properties matter more than the boxes. Nothing writes to Elasticsearch except the indexer, and the indexer can always rebuild what it writes from PostgreSQL. Those are chapter 01’s source-of-truth rules turned into components.

Why the Write API must not also write to Elasticsearch

The obvious design writes to both systems in the same request handler: commit to PostgreSQL, then index into Elasticsearch. It fails in both orders.

ElasticsearchPostgreSQLWrite APIElasticsearchPostgreSQLWrite APIPostgreSQL has the change. Elasticsearch never will.UPDATE user_profile (commit)committedindex profile (timeout, process killed, or 429)
ElasticsearchPostgreSQLWrite APIElasticsearchPostgreSQLWrite APIPostgreSQL has the change. Elasticsearch never will.UPDATE user_profile (commit)committedindex profile (timeout, process killed, or 429)

The sequence diagram shows the Write API committing an update to PostgreSQL, then failing to index into Elasticsearch because of a timeout, a killed process, or a rejected request. PostgreSQL has the change and Elasticsearch never receives it. Reverse the order and a failed commit leaves Elasticsearch holding a change that never happened, which breaks the rule that nothing exists only in the projection.

A distributed transaction across both is not available: Elasticsearch has no prepare or commit phase to take part in one. Retrying in memory only moves the loss to the moment the process restarts. The only durable place to record “this profile changed and the index must follow” is the same PostgreSQL transaction that made the change.

Three ways to synchronise, compared

Bulk backfill plus dual writeTransactional outboxChange data capture (Debezium and Kafka)
How changes are capturedApplication code writes to bothTrigger or application writes an outbox row in the same transactionConnector reads PostgreSQL’s write-ahead log
Lost change on crashYes, between the two writesNoNo
Captures changes made outside the application (SQL fixes, other services)NoYes, with a triggerYes
Commit orderNot preservedMust be handled by the pollerPreserved by the log
New infrastructureNoneNoneKafka, Kafka Connect, Debezium, a replication slot
Load on PostgreSQLNone extraOutbox inserts, polling, and deletesLogical decoding; a stalled slot retains WAL on disk
ReplayNot possibleNot possible once rows are deleted, so rebuild from the table insteadFrom the topic’s retention
Operational skill neededLow, until the first drift incidentModerateHigh

Dual write is acceptable for a prototype and for nothing else. At 300 million profiles, drift is guaranteed, and nothing in the design detects or repairs it.

The transactional outbox keeps everything inside PostgreSQL. With the trigger used in this chapter, it captures every change to user_profile no matter which code path made it. Its cost is write amplification on the profile table and a poller that must be written carefully, as Stage 2 shows.

Change data capture with the Debezium PostgreSQL connector reads committed changes from the write-ahead log through logical decoding and publishes them to Kafka. It preserves commit order, adds no writes to the profile table, and lets several consumers replay the same stream. Its cost is a distributed system to operate, and a replication slot that retains WAL on the database server if the connector stops consuming.

Recommendation for this example. Start with the trigger-fed outbox. It needs no new infrastructure, it is correct under crashes, and it covers every write path. Move to CDC when a second consumer needs the same change stream, when the outbox’s extra writes show up in PostgreSQL’s load, or when your organisation already runs Kafka Connect. Needs validation: the point where outbox overhead matters depends on your update rate and must be measured, not assumed. The indexer design in the rest of this chapter works with either source, because it treats each event only as a notification.

Stage 1 — Add a row version and the outbox

Create db/sync/01-search-outbox.sql in the lab directory:

db/sync/01-search-outbox.sql
-- Row version: increases by one on every change to a profile row.
ALTER TABLE user_profile
ADD COLUMN row_version bigint NOT NULL DEFAULT 1;
CREATE FUNCTION user_profile_bump_version() RETURNS trigger
LANGUAGE plpgsql AS $$
BEGIN
NEW.row_version := OLD.row_version + 1;
RETURN NEW;
END;
$$;
CREATE TRIGGER user_profile_bump_version
BEFORE UPDATE ON user_profile
FOR EACH ROW EXECUTE FUNCTION user_profile_bump_version();
-- Outbox: one row per committed change. It carries identifiers only, never profile data.
CREATE TABLE profile_search_outbox (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
user_id bigint NOT NULL,
row_version bigint NOT NULL,
operation text NOT NULL CHECK (operation IN ('UPSERT', 'DELETE')),
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE FUNCTION user_profile_to_outbox() RETURNS trigger
LANGUAGE plpgsql AS $$
BEGIN
IF TG_OP = 'DELETE' THEN
INSERT INTO profile_search_outbox (user_id, row_version, operation)
VALUES (OLD.user_id, OLD.row_version + 1, 'DELETE');
RETURN OLD;
END IF;
INSERT INTO profile_search_outbox (user_id, row_version, operation)
VALUES (NEW.user_id, NEW.row_version, 'UPSERT');
RETURN NEW;
END;
$$;
CREATE TRIGGER user_profile_to_outbox
AFTER INSERT OR UPDATE OR DELETE ON user_profile
FOR EACH ROW EXECUTE FUNCTION user_profile_to_outbox();

Four decisions are built into this file.

  • The version comes from the database, not the clock. row_version is the external version the indexer sends to Elasticsearch. Chapter 03 explained why updatedAt in milliseconds is not safe for that.
  • A delete gets the next version. OLD.row_version + 1 is higher than any version of the row that was ever indexed, so the delete wins over every earlier update.
  • The outbox holds no personal data. Only user_id, a version, and an operation. When a user is erased, no copy of their profile lingers in a queue, and the outbox can be logged and inspected freely.
  • Adding the column is cheap. Since PostgreSQL 11, ADD COLUMN with a constant default does not rewrite the table, so this is fast even on a large table. The triggers add a small cost to every future write.

Apply it:

Terminal window
docker compose exec -T postgres psql -v ON_ERROR_STOP=1 -U profiles -d profiles < db/sync/01-search-outbox.sql
ALTER TABLE
CREATE FUNCTION
CREATE TRIGGER
CREATE TABLE
CREATE FUNCTION
CREATE TRIGGER

Now make some changes. Later chapters rely on this exact data, so run these statements as written:

UPDATE user_profile SET city = 'Gaya', updated_at = now() WHERE user_id = 42;
UPDATE user_profile SET account_status = 'SUSPENDED', updated_at = now() WHERE user_id = 42;
UPDATE user_profile SET account_status = 'SUSPENDED', updated_at = now() WHERE user_id = 42;
INSERT INTO user_profile (first_name, last_name, email, mobile_number, city, state, pincode,
date_of_birth, gender, account_status, created_at, updated_at)
VALUES ('Zoya', 'Khan', 'zoya.khan.new@example.com', '8000000001', 'Lucknow', 'Uttar Pradesh',
'226001', '1994-03-12', 'FEMALE', 'ACTIVE', now(), now())
RETURNING user_id, row_version;
DELETE FROM user_profile WHERE user_id = 43;
SELECT id, user_id, row_version, operation FROM profile_search_outbox ORDER BY id;
user_id | row_version
---------+-------------
1000001 | 1
id | user_id | row_version | operation
----+---------+-------------+-----------
1 | 42 | 2 | UPSERT
2 | 42 | 3 | UPSERT
3 | 42 | 4 | UPSERT
4 | 1000001 | 1 | UPSERT
5 | 43 | 2 | DELETE

The third UPDATE changed nothing but updated_at, and it still produced an event. That is deliberate: the indexer re-reads the row anyway, so a redundant event costs one extra read, not a wrong result.

Checkpoint

SELECT user_id, row_version, city, account_status FROM user_profile WHERE user_id = 42; returns version 4, city Gaya, and status SUSPENDED. The outbox holds five rows. If the outbox is empty, the AFTER trigger was not created; check \d user_profile for both triggers.

Stage 2 — Reproduce the poller bug that loses events

The natural way to read an outbox is to remember the highest id you processed and ask for id > last_id next time. That loses events, because identity values are assigned when a row is inserted, not when its transaction commits.

Open two terminals in the lab directory. In the first, start a transaction that holds its change open:

Terminal window
docker compose exec postgres psql -U profiles -d profiles
BEGIN;
UPDATE user_profile SET city = 'Pune', updated_at = now() WHERE user_id = 44;

In the second terminal, commit a change to another profile, then look at the outbox as a poller would:

Terminal window
docker compose exec postgres psql -U profiles -d profiles \
-c "UPDATE user_profile SET city = 'Kochi', updated_at = now() WHERE user_id = 45;" \
-c "SELECT id, user_id, operation FROM profile_search_outbox WHERE id > 5 ORDER BY id;"
id | user_id | operation
----+---------+-----------
7 | 45 | UPSERT

A watermark poller processes row 7 and records last_id = 7. Now commit the first transaction with COMMIT; and run the same SELECT again:

id | user_id | operation
----+---------+-----------
6 | 44 | UPSERT
7 | 45 | UPSERT

Row 6 appeared behind the watermark. The poller will never ask for it, and profile 44’s move to Pune never reaches the index. Under real write concurrency this happens continuously, and nothing reports it.

The fix is to consume by claiming and deleting, not by watermark. Each poll takes whatever rows are currently visible, locks them so that parallel indexers skip them, and deletes them in the same transaction that the indexer commits only after Elasticsearch has accepted the batch:

Claim a batch (run inside the indexer's transaction)
WITH batch AS (
SELECT id
FROM profile_search_outbox
ORDER BY id
LIMIT 500
FOR UPDATE SKIP LOCKED
)
DELETE FROM profile_search_outbox AS o
USING batch
WHERE o.id = batch.id
RETURNING o.id, o.user_id, o.row_version, o.operation;

Try it inside a transaction that you roll back, so the queue stays intact for chapter 14:

BEGIN;
-- the claim statement above
ROLLBACK;
SELECT count(*) AS still_queued FROM profile_search_outbox;
id | user_id | row_version | operation
----+---------+-------------+-----------
1 | 42 | 2 | UPSERT
...
7 | 45 | 2 | UPSERT
(7 rows)
still_queued
--------------
7

A row that commits late is simply there on the next poll. If the indexer crashes after claiming a batch, the transaction rolls back and the rows return to the queue. If the indexer writes to Elasticsearch and then fails to commit, the batch is processed again. That is at-least-once delivery, which is safe only because every write is idempotent, the subject of the next section. SKIP LOCKED is what lets several indexer instances share one queue without waiting on each other’s locks.

Production note — Rows are deleted as they are consumed, so the outbox stays small and its index stays hot. If you need an audit trail of changes, write it somewhere else; do not keep consumed rows in the queue table.

Stage 3 — Design the indexer: events are notifications, not data

The outbox rows carry no profile data, so the indexer must read the current row before it writes. That is not a limitation; it is the property that makes the pipeline safe against every ordering problem from chapter 03.

For each claimed batch, the indexer:

  1. Collapses the batch to one entry per user_id, keeping the highest row_version and noting whether any entry is a DELETE.
  2. Reads the current rows in one query: SELECT ... FROM user_profile WHERE user_id = ANY(:ids).
  3. For each profile, sends one bulk operation:
    • Row exists: index the complete document, built from the row just read, with version = row_version and version_type = external.
    • Row missing and a DELETE event is present: delete with that event’s version.
    • Row missing and only UPSERT events: skip. The row was deleted after those events, and its DELETE event is already in the queue.
  4. Classifies every bulk item with chapter 03’s table. A 409 means Elasticsearch already has this or a newer version, so it is success. 429 and 5xx items are retried with backoff. Any other failure goes to a dead-letter table with the error.
  5. Commits the outbox transaction only when every item is either applied or dead-lettered.

Trace the failure cases against those rules:

SituationWhat happens
The same event is delivered twiceThe second write carries the same version and gets 409, treated as success
Events for one profile arrive out of orderThe indexer always writes the current row at its current version, so order does not matter
A stale update arrives long after a delete, when the tombstone has expiredThe row is missing, so the event is skipped. Nothing is resurrected
Elasticsearch is downBulk fails, the transaction rolls back, and the outbox grows until Elasticsearch returns
One document is invalid, for example a mapping violationThat item goes to the dead-letter table, and the rest of the batch commits
Someone runs _update_by_query on the indexVersions drift and later events are silently rejected. Forbidden by policy; see chapter 03

The third row is the reason for this design. Chapter 03 showed that an event-carried version cannot stop a resurrection once a tombstone expires. Re-reading the source can, because a deleted row cannot be read.

Trade-off. Re-reading costs one indexed lookup per changed profile, batched by primary key. It must read from the primary, or from a replica whose lag is lower than the time an event spends in the queue; otherwise it can read an older version than the event announced. The alternative, carrying the full row in the event, avoids the read but makes ordering, duplicate handling, and tombstone expiry your problem again, and it puts personal data into the queue.

Production note — One path still writes without an outbox event: the bulk backfill that builds a new index version (chapter 14). It reads a row, then writes it seconds later. If the row is deleted in between and its tombstone expires before the backfill write arrives, the profile comes back. Chapter 14 raises index.gc_deletes on the new index for the duration of the backfill, so tombstones outlive the backfill’s read-to-write delay.

Retries, dead letters, and ordering across instances

Retries. Only 429, 5xx, and connection failures are retried, with exponential backoff and a cap. A 429 means the cluster’s write queue is full. Retrying immediately makes it worse, and it is the signal to slow the indexer down. That is backpressure, and chapter 14 implements it.

Dead letters. A permanently failing item goes to a profile_search_dead_letter table with the user_id, the version, the error type, and the reason, and the rest of the batch proceeds. Replaying a dead letter is just inserting a fresh outbox row for that user_id after fixing the cause. Because events carry no data, the replay reads the corrected row.

Ordering across instances. Several indexers can share the queue thanks to SKIP LOCKED, and two of them may process events for the same profile at the same time. External versions settle the race: whichever write carries the higher row_version wins, whatever order they arrive in.

Eventual consistency, stated as a budget. A change is searchable after: the outbox poll interval, plus batch processing time, plus the index refresh interval (chapter 02). With a one-second poll and the default one-second refresh, a typical change is searchable within a few seconds. Needs validation: measure this end to end in your environment, and treat the Search API’s results as possibly stale by that budget. Anything that must be exact, such as “is this account suspended right now?”, reads PostgreSQL.

Replay and rebuild

The outbox is a queue, not a history. Once consumed, events are gone, so “replay from the beginning” is not an operation this design supports. It does not need to: the full state is in user_profile, and rebuilding the index from the table (chapter 14) is the replay. CDC with Kafka adds genuine replay from the topic’s retention, which is valuable when several systems consume the same changes.

What you built, and what comes next

The lab’s user_profile table now has a row_version column, and every committed change is recorded in profile_search_outbox in the same transaction, with identifiers only. You reproduced the watermark bug that loses late-committing events, replaced it with a claim-and-delete batch, and designed an indexer that is idempotent, order-independent, and safe against resurrection, because it treats events as notifications and PostgreSQL as the only source of profile data.

What this chapter did not build is the indexer itself, or anything in Elasticsearch. The outbox holds seven events, waiting.

Chapter 05 designs the index those events will be written to: a versioned index name, separate read and write aliases, and the complete index creation request with its settings, analysers, and mappings.

ElasticsearchPostgresOutbox

Type to search the site.

↑↓ navigate⏎ openPowered by Pagefind