Series overview
Part 14 of 1878% complete
2026-05-26•13 min read

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.display and state.display are keywords without a normaliser, so facets can show Maharashtra.
  • No doc values on email and mobileNumber. Chapter 08’s disk-usage breakdown showed email as 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:

search-index/src/main/resources/es/user-profile-v2.json
{
"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:

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:

indexer/build.gradle.kts
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")
}
indexer/src/main/resources/application.yaml
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: 2

The 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:

db/sync/02-dead-letter.sql
-- 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()
);
Terminal window
docker compose exec -T postgres psql -v ON_ERROR_STOP=1 -U profiles -d profiles < db/sync/02-dead-letter.sql

The shared IndexAdmin gains settings updates, counts, refresh, and the alias swap. Replace the file:

search-index/src/main/kotlin/in/o612/eng/usersearch/index/IndexAdmin.kt
package `in`.o612.eng.usersearch.index
import co.elastic.clients.elasticsearch.ElasticsearchClient
import co.elastic.clients.elasticsearch._types.ElasticsearchException
import co.elastic.clients.elasticsearch.indices.CreateIndexRequest
import co.elastic.clients.elasticsearch.indices.PutIndicesSettingsRequest
import co.elastic.clients.elasticsearch.synonyms.PutSynonymRequest
import 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:

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/IndexerApplication.kt
package `in`.o612.eng.usersearch.indexer
import org.springframework.boot.autoconfigure.SpringBootApplication
import org.springframework.boot.context.properties.ConfigurationPropertiesScan
import org.springframework.boot.runApplication
import org.springframework.scheduling.annotation.EnableScheduling
@SpringBootApplication
@ConfigurationPropertiesScan
@EnableScheduling
class IndexerApplication
fun main(args: Array<String>) {
runApplication<IndexerApplication>(*args)
}
indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/IndexerProperties.kt
package `in`.o612.eng.usersearch.indexer
import org.springframework.boot.context.properties.ConfigurationProperties
import 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)
}
indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/ElasticsearchConfig.kt
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClient
import co.elastic.clients.json.JsonpMapper
import `in`.o612.eng.usersearch.index.ElasticsearchJson
import `in`.o612.eng.usersearch.index.IndexAdmin
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
@Configuration
class 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.

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/ProfileReader.kt
package `in`.o612.eng.usersearch.indexer
import `in`.o612.eng.usersearch.index.UserProfileDocument
import org.springframework.jdbc.core.simple.JdbcClient
import org.springframework.stereotype.Component
import java.sql.ResultSet
import java.time.LocalDate
import 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. */
@Component
class 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:

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/BulkOutcome.kt
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:

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/DeadLetters.kt
package `in`.o612.eng.usersearch.indexer
import org.springframework.jdbc.core.simple.JdbcClient
import org.springframework.stereotype.Component
/** Changes that failed permanently. Replaying one is inserting a new outbox row for its user_id. */
@Component
class 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.

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/BackfillJob.kt
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClient
import co.elastic.clients.elasticsearch._helpers.bulk.BulkIngester
import co.elastic.clients.elasticsearch._helpers.bulk.BulkListener
import co.elastic.clients.elasticsearch._types.VersionType
import co.elastic.clients.elasticsearch.core.BulkRequest
import co.elastic.clients.elasticsearch.core.BulkResponse
import co.elastic.clients.elasticsearch.core.bulk.BulkOperation
import `in`.o612.eng.usersearch.index.IndexAdmin
import org.slf4j.LoggerFactory
import org.springframework.stereotype.Component
import java.util.concurrent.ConcurrentLinkedQueue
import 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.
*/
@Component
class 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_deletes is 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. The finally block restores all three, even when the job fails.
  • Contexts connect results to rows. Each operation is added with its ProfileRow as context, and the listener zips the response items with the contexts to classify each one.
  • Retries re-read. A 429 or 5xx item 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_deletes by sending null, 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 at refresh_interval: -1 accepts 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.

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/OutboxRelay.kt
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClient
import co.elastic.clients.elasticsearch._types.VersionType
import co.elastic.clients.elasticsearch.core.bulk.BulkOperation
import org.slf4j.LoggerFactory
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty
import org.springframework.jdbc.core.simple.JdbcClient
import org.springframework.scheduling.annotation.Scheduled
import org.springframework.stereotype.Component
import 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:

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/ReindexValidator.kt
package `in`.o612.eng.usersearch.indexer
import co.elastic.clients.elasticsearch.ElasticsearchClient
import co.elastic.clients.elasticsearch._types.FieldValue
import co.elastic.clients.elasticsearch._types.SortOrder
import co.elastic.clients.elasticsearch._types.query_dsl.Operator
import `in`.o612.eng.usersearch.index.IndexAdmin
import org.springframework.stereotype.Component
/** Checks a new index version against PostgreSQL and against the version it replaces. */
@Component
class 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:

indexer/src/main/kotlin/in/o612/eng/usersearch/indexer/CommandRunner.kt
package `in`.o612.eng.usersearch.indexer
import `in`.o612.eng.usersearch.index.IndexAdmin
import `in`.o612.eng.usersearch.index.UserProfileIndex
import org.springframework.boot.ApplicationArguments
import org.springframework.boot.ApplicationRunner
import org.springframework.boot.SpringApplication
import org.springframework.context.ApplicationContext
import org.springframework.stereotype.Component
import kotlin.system.exitProcess
/** Runs a one-off command (create, backfill, validate, swap) and exits. The default, relay, keeps running. */
@Component
class 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:

Terminal window
set -a; source ../user-search-lab/.env; set +a
export DATABASE_URL=jdbc:postgresql://localhost:5432/profiles
export DATABASE_USERNAME=profiles DATABASE_PASSWORD="$POSTGRES_PASSWORD"
export ELASTICSEARCH_URIS=https://localhost:9200
export ELASTICSEARCH_USERNAME=elastic ELASTICSEARCH_PASSWORD="$ELASTIC_PASSWORD"
export ELASTICSEARCH_CA_CERT="file:$(cd ../user-search-lab && pwd)/certs/ca/ca.crt"
./gradlew :indexer:bootRun

The 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.

fails

a. Create v2

b. Backfill v2

relay writes v1 and v2

c. Outbox drained

d. Validate v2

e. Swap both aliases

f. Keep v1 current

for rollback

Delete v1

Fix and backfill again

fails

a. Create v2

b. Backfill v2

relay writes v1 and v2

c. Outbox drained

d. Validate v2

e. Swap both aliases

f. Keep v1 current

for rollback

Delete v1

Fix and backfill again

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:

Terminal window
./gradlew :indexer:bootRun --args='--indexer.command=create --indexer.target-index=user-profile-v2'
created user-profile-v2: true

b. 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:

Terminal window
INDEXER_RELAY_TARGETS=user-profile-write,user-profile-v2 ./gradlew :indexer:bootRun

In a second terminal, start the backfill:

Terminal window
./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
-------
0

The 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.

Terminal window
./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 OK
search 'prashant kumar' total 1191/1191 top-10 identical
search 'mohd' total 23333/23333 top-10 identical
search 'jose fernandes' total 1170/1170 top-10 identical
search 'priya' total 23384/23384 top-10 identical
search 'lakshmi iyer' total 1164/1164 top-10 identical

Your 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:

Terminal window
./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:

Terminal window
INDEXER_RELAY_TARGETS=user-profile-write,user-profile-v1 ./gradlew :indexer:bootRun

A 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:

search-api/src/main/kotlin/in/o612/eng/usersearch/api/search/UserFacetsBuilder.kt
/** 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")
search-api/src/main/kotlin/in/o612/eng/usersearch/api/search/UserFacetsBuilder.kt
.aggregations("values", Aggregation.of { t -> t.terms { tt -> tt.field(AGGREGATION_FIELDS.getValue(field)).size(BUCKETS_PER_FACET) } })

Restart search-api:

Terminal window
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 display sub-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

  • OFFSET pagination for the backfill. It slows down with every page. Use the primary key.
  • Ignoring bulk item results. An HTTP 200 can 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_deletes in a finally block, 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.

ElasticsearchSpring BootKotlinPostgres

Type to search the site.

↑↓ navigate⏎ openPowered by Pagefind