Q: Your e-commerce catalog contains 100 million products. When a merchant updates an item's price or stock in the primary database, the updated information must be searchable on the website within 1 second. Dual-writing from application code causes race conditions and stale search results. How do you design an asynchronous, real-time search indexing pipeline with CDC, Kafka, and OpenSearch?
Engineering a real-time document search indexing pipeline synchronizing 100 million database entities to OpenSearch with sub-second update visibility, out-of-order event reconciliation, and zero reindexing downtime.
Want to master this scenario in a live sandbox? The Linux Foundation's FinOps Certified Practitioner (FOCP) Program covers this exact problem with hands-on terminal drills.
🛠️ Production Runbook & Step-by-Step Resolution
Capture Database Changes via Debezium CDC into Apache Kafka
Extract committed mutations from the primary relational database with zero impact on queries:
- WAL Capture: Debezium reads committed transaction logs from PostgreSQL/MySQL, capturing product inserts, updates, and deletes in real time.
- Kafka Partitioning: Events publish to Kafka topic
catalog.cdc.eventspartitioned strictly byproduct_id. - Per-Entity Ordering: Partitioning by product ID guarantees that all updates for a single product are processed in strict sequential order.
Denormalize & Enrich Search Documents via Apache Flink / Kafka Streams
Transform relational normalized rows into comprehensive search documents:
- Stream Join: Flink consumes product updates and joins with merchant profiles, category trees, and review ratings stored in an in-memory RocksDB state store.
- Search Document Serialization: Emits a complete, fully denormalized JSON search document containing keywords, facets, price, and inventory status.
- Write Buffer: Emits enriched documents to
opensearch.indexing.stream.
Ingest Documents via OpenSearch Bulk API & External Versioning
Maximize indexing throughput while guaranteeing protection against out-of-order writes:
- Bulk Worker Fleet: Autoscaled indexing workers consume from Kafka and flush documents using the OpenSearch Bulk API in micro-batches (5,000 docs or 500ms).
- External Versioning: Documents are indexed with
version_type=external_gteusing the database transaction commit timestamp or LSN as the version. - Stale Write Rejection: If a delayed Kafka message arrives with an older version, OpenSearch rejects the write with HTTP 409 Conflict, preserving the newer state.
Manage Schema Migrations via Index Aliases & Reindex Pipelines
Update search analyzers and field mappings without interrupting live search traffic:
- Search Aliases: Application queries search alias
products_searchpointing to physical indexproducts_v1. - Zero-Downtime Reindex: When analyzers change, pipeline creates
products_v2, runs_reindex, catches up on Kafka CDC stream, and atomically swaps the alias in < 10 milliseconds:POST /_aliases { actions: [{ add: { index: 'products_v2', alias: 'products_search' } }, { remove: { index: 'products_v1', alias: 'products_search' } }] }. - End-to-End Latency: Product updates reflect in search results in 680ms.
- Capture database mutations via Debezium CDC partitioned by entity ID for strict ordering.
- Denormalize relational entities into rich search JSON documents using stream processing.
- Index into OpenSearch using bulk micro-batches and external versioning to eliminate race conditions.
- Use atomic index aliases to execute seamless zero-downtime full reindexing.