Backfill, the outbox relay, and zero-downtime reindexing
Until now, the index was filled by a shell script, and the outbox from chapter 04 has been collecting events that nothing consumed. This chapter builds the component that owns every write to Elasticsearch: a Kotlin indexer application with two jobs. The backfill copies every profile from PostgreSQL into a named index, reading with keyset pagination and writing through BulkIngester with backpressure and per-item error handling. The relay drains the outbox continuously, re-reading current rows as chapter 04 designed.
Then you use both to move the live system from user-profile-v1 to a new user-profile-v2, with changes made in PostgreSQL while the backfill runs, and without the Search API noticing: create, backfill, catch up, validate, swap both aliases atomically, and keep v1 current for a rollback.
You need the lab and the user-search project from chapter 13. The chapter takes about 90 minutes.
Why mapping changes mean a new index
Chapter 02 showed that an existing field’s type cannot change, and chapter 06 that analysis settings apply to terms already written. The same holds for most mapping decisions that matter: turning off doc values, changing an analyser, splitting a field into sub-fields for existing documents, or changing the primary shard count. Adding a brand-new field is allowed in place, but existing documents do not gain values for it until they are indexed again. In practice, any change to how existing profiles are indexed means writing every document again, into a new index. The versioned names and separate aliases from chapter 05 exist for this moment.
The v2 index in this chapter makes two changes found in earlier chapters:
- Display sub-fields for facets. Chapter 12’s facets returned normalised keys such as
maharashtra.city.displayandstate.displayare keywords without a normaliser, so facets can showMaharashtra. - No doc values on
emailandmobileNumber. Chapter 08’s disk-usage breakdown showedemailas the largest indexed field, partly because of doc values that no query shape uses. Exact lookups still work: they use the inverted index, not doc values.
Create es/user-profile-v2.json in the lab directory, and copy it to search-index/src/main/resources/es/ in the project. It is v1 with those changes, and without an aliases block, because the aliases move to v2 only at the swap:
{ "settings": { "index": { "number_of_shards": 1, "number_of_replicas": 0, "refresh_interval": "1s" }, "analysis": { "filter": { "name_edge_ngram": { "type": "edge_ngram", "min_gram": 2, "max_gram": 15 }, "name_synonyms": { "type": "synonym_graph", "synonyms_set": "profile-name-synonyms", "updateable": true } }, "analyzer": { "name_index": { "type": "custom", "tokenizer": "standard", "filter": ["lowercase", "asciifolding"] }, "name_search": { "type": "custom", "tokenizer": "standard", "filter": ["lowercase", "asciifolding", "name_synonyms"] }, "name_prefix_index": { "type": "custom", "tokenizer": "standard", "filter": ["lowercase", "asciifolding", "name_edge_ngram"] } }, "normalizer": { "lowercase_ascii": { "type": "custom", "filter": ["lowercase", "asciifolding"] }, "lowercase_only": { "type": "custom", "filter": ["lowercase"] } } } }, "mappings": { "dynamic": "strict", "properties": { "userId": { "type": "keyword" }, "fullName": { "type": "text", "analyzer": "name_index", "search_analyzer": "name_search", "fields": { "prefix": { "type": "text", "analyzer": "name_prefix_index", "search_analyzer": "name_index" }, "keyword": { "type": "keyword", "normalizer": "lowercase_ascii", "ignore_above": 256 } } }, "firstName": { "type": "text", "analyzer": "name_index", "search_analyzer": "name_search" }, "lastName": { "type": "text", "analyzer": "name_index", "search_analyzer": "name_search" }, "email": { "type": "keyword", "normalizer": "lowercase_only", "doc_values": false }, "mobileNumber": { "type": "keyword", "doc_values": false }, "city": { "type": "keyword", "normalizer": "lowercase_ascii", "fields": { "text": { "type": "text", "analyzer": "name_index" }, "display": { "type": "keyword" } } }, "state": { "type": "keyword", "normalizer": "lowercase_ascii", "fields": { "text": { "type": "text", "analyzer": "name_index" }, "display": { "type": "keyword" } } }, "country": { "type": "keyword" }, "pincode": { "type": "keyword" }, "dateOfBirth": { "type": "date", "index": false, "doc_values": false }, "gender": { "type": "keyword" }, "accountStatus": { "type": "keyword" }, "createdAt": { "type": "date" }, "updatedAt": { "type": "date" } } }}Compared with v1, only three blocks differ: email and mobileNumber gain "doc_values": false, and city and state each gain a display sub-field.
Stage 1 — The indexer module
Add the module to settings.gradle.kts:
rootProject.name = "user-search"
include("search-index", "client-tour", "search-api", "indexer")The indexer needs JDBC for PostgreSQL and the Elasticsearch client, and nothing for the web:
import org.springframework.boot.gradle.plugin.SpringBootPlugin
plugins { kotlin("jvm") kotlin("plugin.spring") id("org.springframework.boot")}
kotlin { jvmToolchain(21) compilerOptions { freeCompilerArgs.add("-Xjsr305=strict") }}
dependencies { implementation(platform(SpringBootPlugin.BOM_COORDINATES)) implementation(project(":search-index")) implementation("org.springframework.boot:spring-boot-starter-elasticsearch") implementation("org.springframework.boot:spring-boot-starter-jdbc") implementation("tools.jackson.module:jackson-module-kotlin") implementation("org.jetbrains.kotlin:kotlin-reflect") runtimeOnly("org.postgresql:postgresql")}spring: application: name: indexer main: web-application-type: none datasource: url: ${DATABASE_URL} username: ${DATABASE_USERNAME} password: ${DATABASE_PASSWORD} elasticsearch: uris: ${ELASTICSEARCH_URIS} username: ${ELASTICSEARCH_USERNAME} password: ${ELASTICSEARCH_PASSWORD} connection-timeout: 2s socket-timeout: 60s restclient: ssl: bundle: elasticsearch ssl: bundle: pem: elasticsearch: truststore: certificate: ${ELASTICSEARCH_CA_CERT}
indexer: # relay (default, runs until stopped), create, backfill, validate, or swap command: ${INDEXER_COMMAND:relay} source-index: ${INDEXER_SOURCE_INDEX:} target-index: ${INDEXER_TARGET_INDEX:} relay: poll-interval: 1s batch-size: 500 # Indices or aliases every change is written to. Add the new index version during a reindex. targets: ${INDEXER_RELAY_TARGETS:user-profile-write} backfill: page-size: 2000 bulk-operations: 1000 concurrent-requests: 2The socket timeout is longer than the Search API’s, because a bulk request with a thousand documents legitimately takes longer than a search. indexer.command selects what the application does; relay, the default, runs until stopped, and the others run once and exit.
Add a table for permanent failures to the lab database:
-- Changes the indexer could not apply and will not retry. Replay: fix the cause, then-- INSERT INTO profile_search_outbox (user_id, row_version, operation) for the affected user_id.CREATE TABLE profile_search_dead_letter ( id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, user_id bigint NOT NULL, row_version bigint NOT NULL, target text NOT NULL, error_type text NOT NULL, error_reason text, created_at timestamptz NOT NULL DEFAULT now());docker compose exec -T postgres psql -v ON_ERROR_STOP=1 -U profiles -d profiles < db/sync/02-dead-letter.sqlThe shared IndexAdmin gains settings updates, counts, refresh, and the alias swap. Replace the file:
package `in`.o612.eng.usersearch.index
import co.elastic.clients.elasticsearch.ElasticsearchClientimport co.elastic.clients.elasticsearch._types.ElasticsearchExceptionimport co.elastic.clients.elasticsearch.indices.CreateIndexRequestimport co.elastic.clients.elasticsearch.indices.PutIndicesSettingsRequestimport co.elastic.clients.elasticsearch.synonyms.PutSynonymRequestimport java.io.StringReader
/** Index and alias management for the user-profile index. Used by applications and jobs, never by search requests. */class IndexAdmin(private val client: ElasticsearchClient) {
/** Indices an alias currently points to; empty if the alias does not exist. */ fun aliasTargets(alias: String): Set<String> = if (client.indices().existsAlias { it.name(alias) }.value()) { client.indices().getAlias { it.name(alias) }.aliases().keys } else { emptySet() }
/** * Creates `user-profile-v<version>` from its reviewed definition, aliases included. * Returns false, and changes nothing, if the index already exists. */ fun createIndex(version: Int): Boolean { val name = UserProfileIndex.indexName(version) if (client.indices().exists { it.index(name) }.value()) return false ensureSynonymsSet() val request = CreateIndexRequest.of { it.index(name).withJson(StringReader(UserProfileIndex.definition(version))) } client.indices().create(request) return true }
/** The name analysers reference the synonyms set, so it must exist before any index version. */ fun ensureSynonymsSet() { val exists = try { client.synonyms().getSynonym { it.id(UserProfileIndex.SYNONYMS_SET) } true } catch (e: ElasticsearchException) { if (e.status() != 404) throw e false } if (!exists) { client.synonyms().putSynonym( PutSynonymRequest.of { it.id(UserProfileIndex.SYNONYMS_SET).withJson(StringReader(UserProfileIndex.synonymsDefinition())) }, ) } }
/** Applies dynamic index settings given as JSON, for example `{"index": {"refresh_interval": "-1"}}`. */ fun putSettings(index: String, settingsJson: String) { client.indices().putSettings( PutIndicesSettingsRequest.of { it.index(index).withJson(StringReader(settingsJson)) }, ) }
/** The number of replicas an index is configured with. */ fun replicaCount(index: String): String = client.indices().getSettings { it.index(index) }[index]?.settings()?.index()?.numberOfReplicas() ?: error("No settings for $index")
fun count(index: String): Long = client.count { it.index(index) }.count()
fun refresh(index: String) { client.indices().refresh { it.index(index) } }
/** * Moves both aliases from [fromIndex] to [toIndex] in one atomic request. * Searches and writes switch together; there is no moment when an alias points nowhere. */ fun swapAliases(fromIndex: String, toIndex: String) { client.indices().updateAliases { u -> u.actions { a -> a.remove { r -> r.index(fromIndex).alias(UserProfileIndex.READ_ALIAS) } } .actions { a -> a.remove { r -> r.index(fromIndex).alias(UserProfileIndex.WRITE_ALIAS) } } .actions { a -> a.add { ad -> ad.index(toIndex).alias(UserProfileIndex.READ_ALIAS) } } .actions { a -> a.add { ad -> ad.index(toIndex).alias(UserProfileIndex.WRITE_ALIAS).isWriteIndex(true) } } } }}swapAliases sends four actions in one update aliases request: remove both aliases from the old index and add both to the new one. Elasticsearch applies the actions atomically, so there is no moment when user-profile-read points at nothing, or at both indices.
The application class, its typed configuration, and the client configuration from chapter 10:
package `in`.o612.eng.usersearch.indexer
import org.springframework.boot.autoconfigure.SpringBootApplicationimport org.springframework.boot.context.properties.ConfigurationPropertiesScanimport org.springframework.boot.runApplicationimport org.springframework.scheduling.annotation.EnableScheduling
@SpringBootApplication@ConfigurationPropertiesScan@EnableSchedulingclass IndexerApplication
fun main(args: Array<String>) { runApplication<IndexerApplication>(*args)}package `in`.o612.eng.usersearch.indexer
import org.springframework.boot.context.properties.ConfigurationPropertiesimport java.time.Duration
@ConfigurationProperties("indexer")data class IndexerProperties( val command: String, val sourceIndex: String, val targetIndex: String, val relay: Relay, val backfill: Backfill,) { data class Relay(val pollInterval: Duration, val batchSize: Int, val targets: List<String>)
data class Backfill(val pageSize: Int, val bulkOperations: Int, val concurrentRequests: Int)}package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClientimport co.elastic.clients.json.JsonpMapperimport `in`.o612.eng.usersearch.index.ElasticsearchJsonimport `in`.o612.eng.usersearch.index.IndexAdminimport org.springframework.context.annotation.Beanimport org.springframework.context.annotation.Configuration
@Configurationclass ElasticsearchConfig { @Bean fun jsonpMapper(): JsonpMapper = ElasticsearchJson.jsonpMapper()
@Bean fun indexAdmin(client: ElasticsearchClient) = IndexAdmin(client)}Stage 2 — Reading PostgreSQL with keyset pagination
What it is. Keyset pagination reads the next page as “rows with a key greater than the last one I saw”, WHERE user_id > :after ORDER BY user_id LIMIT :limit, instead of OFFSET.
Why it matters at 300 million profiles. OFFSET 250000000 makes PostgreSQL walk and discard 250 million rows before returning the page, so every page is slower than the last and the backfill slows to a crawl near the end. A keyset page is a range scan on the primary key, and it costs the same whether it is the first page or the last.
Example. ProfileReader has three queries: a keyset page for the backfill, a lookup by IDs for the relay and for retries, and a count for validation.
package `in`.o612.eng.usersearch.indexer
import `in`.o612.eng.usersearch.index.UserProfileDocumentimport org.springframework.jdbc.core.simple.JdbcClientimport org.springframework.stereotype.Componentimport java.sql.ResultSetimport java.time.LocalDateimport java.time.OffsetDateTime
/** A profile row, with the version the index uses as its external version. */data class ProfileRow(val document: UserProfileDocument, val rowVersion: Long) { val userId: Long get() = document.userId.toLong()}
/** Reads profiles from PostgreSQL, the source of truth. */@Componentclass ProfileReader(private val jdbc: JdbcClient) {
/** * Keyset pagination: the next [limit] profiles after [afterUserId], by primary key. * Every page is an index range scan, however deep the backfill is. OFFSET would re-read every skipped row. */ fun page(afterUserId: Long, limit: Int): List<ProfileRow> = jdbc.sql("$SELECT WHERE user_id > :after ORDER BY user_id LIMIT :limit") .param("after", afterUserId) .param("limit", limit) .query { rs, _ -> rs.toProfileRow() } .list()
/** The current state of specific profiles. Missing ids have been deleted. */ fun byIds(userIds: Collection<Long>): Map<Long, ProfileRow> = if (userIds.isEmpty()) { emptyMap() } else { jdbc.sql("$SELECT WHERE user_id = ANY(:ids)") .param("ids", userIds.toTypedArray()) .query { rs, _ -> rs.toProfileRow() } .list() .associateBy { it.userId } }
fun count(): Long = jdbc.sql("SELECT count(*) FROM user_profile").query(Long::class.java).single()
private fun ResultSet.toProfileRow() = ProfileRow( document = UserProfileDocument( userId = getLong("user_id").toString(), fullName = getString("full_name"), firstName = getString("first_name"), lastName = getString("last_name"), email = getString("email"), mobileNumber = getString("mobile_number"), city = getString("city"), state = getString("state"), country = getString("country"), pincode = getString("pincode"), dateOfBirth = getObject("date_of_birth", LocalDate::class.java), gender = getString("gender"), accountStatus = getString("account_status"), createdAt = getObject("created_at", OffsetDateTime::class.java).toInstant(), updatedAt = getObject("updated_at", OffsetDateTime::class.java).toInstant(), ), rowVersion = getLong("row_version"), )
private companion object { const val SELECT = """ SELECT user_id, full_name, first_name, last_name, email, mobile_number, city, state, country, pincode, date_of_birth, gender, account_status, created_at, updated_at, row_version FROM user_profile """ }}ProfileRow carries row_version next to the document, because every write to Elasticsearch uses it as the external version.
Common mistake. Paging by updated_at for a backfill. Many profiles share a timestamp, so a page boundary inside a group of equal timestamps skips or repeats rows, and rows updated during the backfill move between pages. The primary key never changes and never repeats.
Stage 3 — Classifying every bulk item
Chapter 03 established the rules for bulk item results. They become one small enum used by both jobs:
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.core.bulk.BulkResponseItem
/** What to do with one bulk item's result (chapter 03's table). */enum class BulkOutcome { /** Written. */ APPLIED,
/** 409 on an external-version write: Elasticsearch already has this version or a newer one. */ ALREADY_CURRENT,
/** 429 or 5xx: temporary. Retry later with backoff. */ RETRY,
/** Any other failure: the request or document is wrong. Retrying will not help. */ DEAD_LETTER, ;
companion object { fun of(item: BulkResponseItem): BulkOutcome = when { item.error() == null -> APPLIED item.status() == 409 -> ALREADY_CURRENT item.status() == 429 || item.status() >= 500 -> RETRY else -> DEAD_LETTER } }}Permanent failures go to the dead-letter table with the target and the error, so they can be investigated and replayed:
package `in`.o612.eng.usersearch.indexer
import org.springframework.jdbc.core.simple.JdbcClientimport org.springframework.stereotype.Component
/** Changes that failed permanently. Replaying one is inserting a new outbox row for its user_id. */@Componentclass DeadLetters(private val jdbc: JdbcClient) {
fun record(userId: Long, rowVersion: Long, target: String, errorType: String, errorReason: String?) { jdbc.sql( """ INSERT INTO profile_search_dead_letter (user_id, row_version, target, error_type, error_reason) VALUES (:userId, :rowVersion, :target, :errorType, :errorReason) """, ) .param("userId", userId) .param("rowVersion", rowVersion) .param("target", target) .param("errorType", errorType) .param("errorReason", errorReason) .update() }}Stage 4 — The backfill job
What it is. The backfill writes every profile into one named index, never an alias. It uses the Java client’s BulkIngester, which collects operations into bulk requests and sends them when a request reaches a number of operations or a size in bytes, 1,000 operations and 5 MiB by default.
Why it matters at 300 million profiles. A backfill is the largest write load the cluster sees. It must go as fast as the cluster allows, and no faster: an indexer that sends requests faster than the cluster can process them fills the write thread pool’s queue and gets 429 rejections. BulkIngester provides the brake. The client documentation states that when the maximum number of concurrent requests is in flight and the next request is full, adding an operation blocks. The reading loop therefore slows to the cluster’s pace. That is backpressure.
Example.
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClientimport co.elastic.clients.elasticsearch._helpers.bulk.BulkIngesterimport co.elastic.clients.elasticsearch._helpers.bulk.BulkListenerimport co.elastic.clients.elasticsearch._types.VersionTypeimport co.elastic.clients.elasticsearch.core.BulkRequestimport co.elastic.clients.elasticsearch.core.BulkResponseimport co.elastic.clients.elasticsearch.core.bulk.BulkOperationimport `in`.o612.eng.usersearch.index.IndexAdminimport org.slf4j.LoggerFactoryimport org.springframework.stereotype.Componentimport java.util.concurrent.ConcurrentLinkedQueueimport java.util.concurrent.atomic.AtomicLong
data class BackfillReport(val read: Long, val applied: Long, val alreadyCurrent: Long, val deadLettered: Long, val retried: Long)
/** * Copies every profile from PostgreSQL into one index, by name. Never targets an alias: * a backfill fills a new index version before anything reads it. */@Componentclass BackfillJob( private val client: ElasticsearchClient, private val reader: ProfileReader, private val indexAdmin: IndexAdmin, private val deadLetters: DeadLetters, private val properties: IndexerProperties,) { private val log = LoggerFactory.getLogger(javaClass)
fun run(targetIndex: String): BackfillReport { val settings = properties.backfill val replicas = indexAdmin.replicaCount(targetIndex) // No refresh and no replicas while loading; tombstones kept long enough to outlive the load (chapter 04). indexAdmin.putSettings(targetIndex, """{"index":{"refresh_interval":"-1","number_of_replicas":0,"gc_deletes":"1h"}}""")
val applied = AtomicLong() val alreadyCurrent = AtomicLong() val deadLettered = AtomicLong() val retry = ConcurrentLinkedQueue<Long>() var read = 0L
val listener = object : BulkListener<ProfileRow> { override fun beforeBulk(executionId: Long, request: BulkRequest, contexts: List<ProfileRow>) = Unit
override fun afterBulk(executionId: Long, request: BulkRequest, contexts: List<ProfileRow>, response: BulkResponse) { response.items().zip(contexts).forEach { (item, row) -> when (BulkOutcome.of(item)) { BulkOutcome.APPLIED -> applied.incrementAndGet() BulkOutcome.ALREADY_CURRENT -> alreadyCurrent.incrementAndGet() BulkOutcome.RETRY -> retry.add(row.userId) BulkOutcome.DEAD_LETTER -> { deadLetters.record(row.userId, row.rowVersion, targetIndex, item.error()?.type() ?: "unknown", item.error()?.reason()) deadLettered.incrementAndGet() } } } }
override fun afterBulk(executionId: Long, request: BulkRequest, contexts: List<ProfileRow>, failure: Throwable) { log.warn("Bulk request {} failed as a whole; {} items will be retried: {}", executionId, contexts.size, failure.message) contexts.forEach { retry.add(it.userId) } } }
try { // add() blocks while the maximum number of requests is in flight: that is the backpressure. BulkIngester.of<ProfileRow> { it.client(client) .maxOperations(settings.bulkOperations) .maxConcurrentRequests(settings.concurrentRequests) .listener(listener) }.use { ingester -> var after = 0L while (true) { val page = reader.page(after, settings.pageSize) if (page.isEmpty()) break page.forEach { row -> ingester.add(indexOperation(targetIndex, row), row) } read += page.size after = page.last().userId if (read % 100_000 < settings.pageSize) log.info("Backfill read {} profiles", read) } } // close() flushes the last batch and waits for every request val retried = retryFailed(targetIndex, retry.toSet(), applied, alreadyCurrent) return BackfillReport(read, applied.get(), alreadyCurrent.get(), deadLettered.get(), retried) } finally { // Always restore, even after a failure: an index left at refresh -1 never shows new writes. // gc_deletes goes back to its documented default, 60s; the typed settings API cannot send null. indexAdmin.putSettings(targetIndex, """{"index":{"refresh_interval":"1s","number_of_replicas":$replicas,"gc_deletes":"60s"}}""") indexAdmin.refresh(targetIndex) } }
/** Re-reads each retryable profile, so a retry writes its current state, and tries again with backoff. */ private fun retryFailed(targetIndex: String, userIds: Set<Long>, applied: AtomicLong, alreadyCurrent: AtomicLong): Long { var pending = userIds var delayMillis = 1_000L repeat(MAX_RETRIES) { if (pending.isEmpty()) return userIds.size.toLong() Thread.sleep(delayMillis) delayMillis *= 2 val rows = reader.byIds(pending) // deleted since: the outbox relay handles them if (rows.isEmpty()) return userIds.size.toLong() val response = client.bulk { b -> b.operations(rows.values.map { indexOperation(targetIndex, it) }) } pending = response.items().filter { item -> when (BulkOutcome.of(item)) { BulkOutcome.APPLIED -> { applied.incrementAndGet(); false } BulkOutcome.ALREADY_CURRENT -> { alreadyCurrent.incrementAndGet(); false } else -> true } }.map { it.id()!!.toLong() }.toSet() } check(pending.isEmpty()) { "${pending.size} profiles still failing after $MAX_RETRIES retries" } return userIds.size.toLong() }
private fun indexOperation(targetIndex: String, row: ProfileRow): BulkOperation = BulkOperation.of { op -> op.index { i -> i.index(targetIndex) .id(row.document.userId) .document(row.document) .version(row.rowVersion) .versionType(VersionType.External) } }
private companion object { const val MAX_RETRIES = 3 }}The job applies every rule from earlier chapters:
- Settings for the load, and a guaranteed restore. Refresh is off and replicas are zero during the load, as chapter 08 recommends, and
gc_deletesis raised to one hour, which closes the resurrection gap chapter 04 left open: a delete that happens during the backfill keeps its tombstone longer than the backfill’s read-to-write delay. Thefinallyblock restores all three, even when the job fails. - Contexts connect results to rows. Each operation is added with its
ProfileRowas context, and the listener zips the response items with the contexts to classify each one. - Retries re-read. A
429or5xxitem is not re-sent as it was. After the load, retryable profiles are read again from PostgreSQL, so a retry writes their current state, with exponential backoff. A profile deleted in the meantime is skipped; the relay deletes it. - External versions make it safe to repeat. Run the backfill twice and the second run applies nothing: every item returns
409, counted as already current.
Production note — The first version of this job restored
gc_deletesby sendingnull, the usual way to reset a setting to its default. The typed settings API in client 9.4.5 rejected it while parsing the JSON (Unexpected JSON event 'VALUE_NULL'), after the load had finished, and the index was left with refresh disabled. The job now restores the documented default,60s, explicitly. Whatever your restore mechanism, test the failure path: an index left atrefresh_interval: -1accepts writes and never shows them.
Trade-off. This backfill rebuilds from PostgreSQL. The _reindex API can copy v1 into v2 inside the cluster, which is faster and does not load PostgreSQL. It also copies whatever drift v1 has accumulated, and it cannot apply changes to how documents are built from rows. Rebuilding from the source of truth doubles as a repair; _reindex is the better choice when only the mapping changes and v1 is trusted.
Stage 5 — The outbox relay
The relay implements chapter 04’s design exactly: claim a batch with SKIP LOCKED, collapse it to one change per profile, re-read the current rows, and write each change with its row version as the external version. It commits the claim only when every item is applied or dead-lettered.
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClientimport co.elastic.clients.elasticsearch._types.VersionTypeimport co.elastic.clients.elasticsearch.core.bulk.BulkOperationimport org.slf4j.LoggerFactoryimport org.springframework.boot.autoconfigure.condition.ConditionalOnPropertyimport org.springframework.jdbc.core.simple.JdbcClientimport org.springframework.scheduling.annotation.Scheduledimport org.springframework.stereotype.Componentimport org.springframework.transaction.support.TransactionTemplate
data class OutboxEvent(val id: Long, val userId: Long, val rowVersion: Long, val operation: String)
/** * Drains profile_search_outbox (chapter 04). Events are notifications: the relay re-reads each profile's * current row and writes that, with its row version as the external version, to every configured target. */@Component@ConditionalOnProperty("indexer.command", havingValue = "relay", matchIfMissing = true)class OutboxRelay( private val jdbc: JdbcClient, private val transactions: TransactionTemplate, private val reader: ProfileReader, private val client: ElasticsearchClient, private val deadLetters: DeadLetters, private val properties: IndexerProperties,) { private val log = LoggerFactory.getLogger(javaClass)
@Scheduled(fixedDelayString = "\${indexer.relay.poll-interval}") fun poll() { do { val processed = transactions.execute { relayOneBatch() } ?: 0 } while (processed == properties.relay.batchSize) }
/** Claims a batch, writes it to every target, and commits only when every item is applied or dead-lettered. */ fun relayOneBatch(): Int { val events = claimBatch(properties.relay.batchSize) if (events.isEmpty()) return 0
val rows = reader.byIds(events.map { it.userId }.toSet()) val changes = events.groupBy { it.userId }.mapNotNull { (userId, userEvents) -> val row = rows[userId] val delete = userEvents.filter { it.operation == "DELETE" }.maxByOrNull { it.rowVersion } when { row != null -> Change(userId, row.rowVersion, row) delete != null -> Change(userId, delete.rowVersion, null) else -> null // deleted after these events; its DELETE event is still in the queue } }
for (target in properties.relay.targets) { write(target, changes) } log.info("Relayed {} events as {} changes to {}", events.size, changes.size, properties.relay.targets) return events.size }
private fun write(target: String, changes: List<Change>) { var pending = changes var delayMillis = 500L repeat(MAX_ATTEMPTS) { attempt -> if (pending.isEmpty()) return if (attempt > 0) Thread.sleep(delayMillis).also { delayMillis *= 2 } val response = client.bulk { b -> b.operations(pending.map { it.operation(target) }) } pending = response.items().zip(pending).mapNotNull { (item, change) -> when (BulkOutcome.of(item)) { BulkOutcome.APPLIED, BulkOutcome.ALREADY_CURRENT -> null BulkOutcome.RETRY -> change BulkOutcome.DEAD_LETTER -> { deadLetters.record(change.userId, change.version, target, item.error()?.type() ?: "unknown", item.error()?.reason()) null } } } } // Still failing: roll back the transaction so the events return to the queue. check(pending.isEmpty()) { "${pending.size} changes to $target still failing after $MAX_ATTEMPTS attempts" } }
private fun claimBatch(limit: Int): List<OutboxEvent> = jdbc.sql( """ WITH batch AS ( SELECT id FROM profile_search_outbox ORDER BY id LIMIT :limit 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 """, ) .param("limit", limit) .query { rs, _ -> OutboxEvent(rs.getLong("id"), rs.getLong("user_id"), rs.getLong("row_version"), rs.getString("operation")) } .list()
/** One profile's net change: index its current row, or delete it. */ private data class Change(val userId: Long, val version: Long, val row: ProfileRow?) { fun operation(target: String): BulkOperation = BulkOperation.of { op -> if (row != null) { op.index { i -> i.index(target).id(userId.toString()).document(row.document) .version(version).versionType(VersionType.External) } } else { op.delete { d -> d.index(target).id(userId.toString()).version(version).versionType(VersionType.External) } } } }
private companion object { const val MAX_ATTEMPTS = 4 }}indexer.relay.targets is the list of indices or aliases every change goes to. Normally it is just user-profile-write. During a reindex it also names the new version, and after the swap, the old one. Stage 7 uses both.
Stage 6 — Commands for the reindex
Validation compares the new index with PostgreSQL and with the index it replaces:
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClientimport co.elastic.clients.elasticsearch._types.FieldValueimport co.elastic.clients.elasticsearch._types.SortOrderimport co.elastic.clients.elasticsearch._types.query_dsl.Operatorimport `in`.o612.eng.usersearch.index.IndexAdminimport org.springframework.stereotype.Component
/** Checks a new index version against PostgreSQL and against the version it replaces. */@Componentclass ReindexValidator( private val client: ElasticsearchClient, private val indexAdmin: IndexAdmin, private val reader: ProfileReader,) { /** Sample searches that must return the same profiles, in the same order, from both versions. */ private val samples = listOf("prashant kumar", "mohd", "jose fernandes", "priya", "lakshmi iyer")
fun validate(sourceIndex: String, targetIndex: String): Boolean { indexAdmin.refresh(targetIndex) val expected = reader.count() val actual = indexAdmin.count(targetIndex) println("count postgres=$expected $targetIndex=$actual ${if (expected == actual) "OK" else "MISMATCH"}") var ok = expected == actual
for (name in samples) { val before = topIds(sourceIndex, name) val after = topIds(targetIndex, name) val same = before == after println("search '$name' total ${before.first}/${after.first} top-10 ${if (same) "identical" else "DIFFERENT"}") ok = ok && same } return ok }
private fun topIds(index: String, name: String): Pair<Long, List<String>> { val response = client.search({ s -> s.index(index).size(10).trackTotalHits { it.enabled(true) }.source { it.fetch(false) } .query { q -> q.bool { b -> b.must { m -> m.match { mt -> mt.field("fullName").query(name).operator(Operator.And) } } .filter { f -> f.term { t -> t.field("accountStatus").value(FieldValue.of("ACTIVE")) } } } } .sort { it.score { sc -> sc.order(SortOrder.Desc) } } .sort { it.field { f -> f.field("updatedAt").order(SortOrder.Desc) } } .sort { it.field { f -> f.field("userId").order(SortOrder.Asc) } } }, Void::class.java) return (response.hits().total()?.value() ?: 0) to response.hits().hits().map { it.id()!! } }}The sample searches must return the same totals and the same top ten, in the same order, from both versions. That checks that nothing in the mapping change altered search results for the shapes that matter. Extend the list with your own recorded queries.
The command runner dispatches create, backfill, validate, and swap, and exits with a non-zero status when a command fails, so a deployment pipeline can stop on it:
package `in`.o612.eng.usersearch.indexer
import `in`.o612.eng.usersearch.index.IndexAdminimport `in`.o612.eng.usersearch.index.UserProfileIndeximport org.springframework.boot.ApplicationArgumentsimport org.springframework.boot.ApplicationRunnerimport org.springframework.boot.SpringApplicationimport org.springframework.context.ApplicationContextimport org.springframework.stereotype.Componentimport kotlin.system.exitProcess
/** Runs a one-off command (create, backfill, validate, swap) and exits. The default, relay, keeps running. */@Componentclass CommandRunner( private val properties: IndexerProperties, private val backfill: BackfillJob, private val validator: ReindexValidator, private val indexAdmin: IndexAdmin, private val context: ApplicationContext,) : ApplicationRunner {
override fun run(args: ApplicationArguments) { val succeeded = when (properties.command) { "relay" -> return "create" -> { val version = requireTarget().substringAfterLast("-v").toInt() println("created ${requireTarget()}: ${indexAdmin.createIndex(version)}") true } "backfill" -> { val report = backfill.run(requireTarget()) println("backfill $report") report.deadLettered == 0L } "validate" -> validator.validate(requireSource(), requireTarget()) "swap" -> { indexAdmin.swapAliases(requireSource(), requireTarget()) println("read -> ${indexAdmin.aliasTargets(UserProfileIndex.READ_ALIAS)}, write -> ${indexAdmin.aliasTargets(UserProfileIndex.WRITE_ALIAS)}") true } else -> error("Unknown indexer.command '${properties.command}'") } exitProcess(SpringApplication.exit(context, { if (succeeded) 0 else 1 })) }
private fun requireSource() = properties.sourceIndex.ifBlank { error("Set indexer.source-index") }
private fun requireTarget() = properties.targetIndex.ifBlank { error("Set indexer.target-index") }}Stage 7 — Start the relay
Every command below runs from the user-search directory. Export the connection settings once; the database password comes from the lab’s .env:
set -a; source ../user-search-lab/.env; set +aexport DATABASE_URL=jdbc:postgresql://localhost:5432/profilesexport DATABASE_USERNAME=profiles DATABASE_PASSWORD="$POSTGRES_PASSWORD"export ELASTICSEARCH_URIS=https://localhost:9200export ELASTICSEARCH_USERNAME=elastic ELASTICSEARCH_PASSWORD="$ELASTIC_PASSWORD"export ELASTICSEARCH_CA_CERT="file:$(cd ../user-search-lab && pwd)/certs/ca/ca.crt"./gradlew :indexer:bootRunThe outbox has held seven events since chapter 04. Within a second:
INFO ... c.e.usersearch.indexer.OutboxRelay : Relayed 7 events as 5 changes to [user-profile-write]Seven events became five changes because profile 42’s three updates collapsed into one. Now change PostgreSQL in another terminal:
UPDATE user_profile SET account_status = 'INACTIVE', updated_at = now() WHERE user_id = 1;DELETE FROM user_profile WHERE user_id = 2;INFO ... c.e.usersearch.indexer.OutboxRelay : Relayed 2 events as 2 changes to [user-profile-write]GET user-profile-read/_doc/1 returns "_version":2 and "accountStatus":"INACTIVE", matching the row’s row_version of 2; GET user-profile-read/_doc/2 returns "found":false. The search projection now follows PostgreSQL within about a second, with no application code writing to Elasticsearch.
Stage 8 — Reindex from v1 to v2 without downtime
The procedure has six steps. Searches and writes continue throughout.
The diagram shows the six steps in order: create v2, backfill it while the relay writes every change to both versions, wait until the outbox is drained, validate, swap both aliases in one request, and keep v1 current for a rollback window before deleting it. A failed validation loops back to fixing the cause and running the backfill again.
a. Create v2. From its reviewed definition:
./gradlew :indexer:bootRun --args='--indexer.command=create --indexer.target-index=user-profile-v2'created user-profile-v2: trueb. Backfill v2 while the relay writes to both. Stop the relay and restart it with v2 as a second target, so every change from now on reaches both versions:
INDEXER_RELAY_TARGETS=user-profile-write,user-profile-v2 ./gradlew :indexer:bootRunIn a second terminal, start the backfill:
./gradlew :indexer:bootRun --args='--indexer.command=backfill --indexer.target-index=user-profile-v2'While it runs, change three profiles in PostgreSQL: one that the backfill has already passed, one it has not reached, and one deletion:
UPDATE user_profile SET account_status = 'SUSPENDED', updated_at = now() WHERE user_id = 3;DELETE FROM user_profile WHERE user_id = 5;UPDATE user_profile SET city = 'Kochi', updated_at = now() WHERE user_id = 999999;The relay applies them to both indices while the backfill continues. The backfill ends with a report. This one comes from a run into an empty copy of v2 at the end of the lab session, when 999,997 profiles remained; your read count reflects the profiles you have created and deleted:
backfill BackfillReport(read=999997, applied=999997, alreadyCurrent=0, deadLettered=0, retried=0)In the lab, reading and indexing about a million profiles took about 96 seconds, including application startup. Your number will differ. A second run of the same command reports applied=0 and alreadyCurrent equal to every profile.
c. Catch up. Changes made during the backfill have already reached v2 through the relay. Before swapping, confirm the queue is empty:
SELECT count(*) FROM profile_search_outbox; count------- 0The three changed profiles agree in PostgreSQL and in both indices: profile 3 is SUSPENDED at version 2, profile 5 is absent, and profile 999999 is in Kochi at version 2. The backfill had read profile 999999 after its update, and profile 3 before it; external versions made the order irrelevant.
d. Validate.
./gradlew :indexer:bootRun --args='--indexer.command=validate --indexer.source-index=user-profile-v1 --indexer.target-index=user-profile-v2'count postgres=999998 user-profile-v2=999998 OKsearch 'prashant kumar' total 1191/1191 top-10 identicalsearch 'mohd' total 23333/23333 top-10 identicalsearch 'jose fernandes' total 1170/1170 top-10 identicalsearch 'priya' total 23384/23384 top-10 identicalsearch 'lakshmi iyer' total 1164/1164 top-10 identicalYour counts differ by the profiles you created and deleted; what matters is that PostgreSQL and v2 agree, and every sample search is identical.
e. Swap. One request moves both aliases:
./gradlew :indexer:bootRun --args='--indexer.command=swap --indexer.source-index=user-profile-v1 --indexer.target-index=user-profile-v2'read -> [user-profile-v2], write -> [user-profile-v2]The Search API now reads v2. It was not restarted and did not notice.
f. Keep v1 current for rollback. Restart the relay with v1 as the extra target, so the old version keeps receiving every change, including deletions, for as long as a rollback might be needed:
INDEXER_RELAY_TARGETS=user-profile-write,user-profile-v1 ./gradlew :indexer:bootRunA deletion made now reaches both indices. A rollback is the same swap command with source and target exchanged, and it is instant. When the rollback window closes, stop writing to v1 and delete it with DELETE user-profile-v1.
Security note — An old index version is a copy of personal data. If you keep v1 without updating it, a profile erased from PostgreSQL survives in v1, and a rollback would bring it back into search results. Either keep writing to the retained version, as step f does, or delete it promptly. Chapter 16 covers erasure across index versions and snapshots.
Stage 9 — Use the new fields
The Search API’s facets can now count the display sub-fields. In UserFacetsBuilder, add the field mapping after BUCKETS_PER_FACET, and use it in the terms aggregation:
/** Facets count the display sub-fields, which keep the original case (user-profile-v2, chapter 14). */ private val AGGREGATION_FIELDS = mapOf("state" to "state.display", "city" to "city.display", "gender" to "gender") .aggregations("values", Aggregation.of { t -> t.terms { tt -> tt.field(AGGREGATION_FIELDS.getValue(field)).size(BUCKETS_PER_FACET) } })Restart search-api:
curl -s 'localhost:8080/api/users/facets?state=bihar&city=patna'{"facets":{"state":[{"value":"Bihar","count":70214}],"city":[{"value":"Patna","count":70214},{"value":"Gaya","count":69711}], "gender":[{"value":"FEMALE","count":34752},{"value":"MALE","count":34707},{"value":"OTHER","count":755}]}}The counts are the same as in chapter 12; the values are now the ones stored in PostgreSQL.
Production note — Deploy this change only after the swap. Against v1, which has no
displaysub-fields, the facets would return empty lists, not an error. Code that depends on a new mapping ships after the index that provides it, and a rollback to v1 needs the previous code too.
The new index is also smaller. GET _cat/indices/user-profile-*?v&h=index,docs.count,store.size reported 150.5 MB for v2, against the 157.6 MB chapter 05 measured for v1 when freshly loaded: the doc values removed from email and mobileNumber outweigh the new display sub-fields.
Checkpoint
Both aliases point to user-profile-v2, the outbox is empty, validate reports OK and every sample identical, and the facets endpoint returns Bihar rather than bihar. If the swap command fails, nothing changed: the aliases still point to v1, because the four actions are applied together or not at all.
Common mistakes with ingestion and reindexing
OFFSETpagination for the backfill. It slows down with every page. Use the primary key.- Ignoring bulk item results. An HTTP
200can hide rejected documents. - Retrying rejected items as they were. Re-read them, so the retry writes their current state.
- Leaving load settings in place. Restore refresh, replicas, and
gc_deletesin afinallyblock, and test that it runs. - Backfilling through an alias. Write the new version by name, so nothing reads it until validation has passed.
- Swapping without catching up. Changes made during the backfill must reach the new index before it serves searches.
- Keeping the old version stale. A rollback to a stale index serves old and deleted profiles.
What you built, and what comes next
The indexer application owns every write to Elasticsearch. Its relay keeps the projection within a second or so of PostgreSQL, with re-reads, external versions, retries, and dead letters. Its backfill rebuilds any index version from the source of truth with backpressure and a guaranteed settings restore. With them you moved the live system from v1 to v2, with concurrent changes, validation, and an atomic swap, while the Search API kept serving.
The relay runs as a single process here. Several instances can share the outbox safely because of SKIP LOCKED and external versions, as chapter 04 argued, but this chapter did not run them concurrently.
Chapter 15 measures what the settings in this chapter and earlier ones actually buy: refresh and replicas during indexing, filter caching and request caching during search, index sorting, and the profile API for finding out where a slow query spends its time.