Scaling Shopify's Search: Enhancing Elasticsearch Resilience With Kubernetes and KE... Leila Vayghan

Leila Vayghan

KubeCon + CloudNativeCon Europe 2025 · Session

Overview

In this insightful talk from KubeCon EU, Leila Vayghan, a Site Reliability Engineer in Shopify’s resiliency organization, detailed a critical project aimed at dramatically improving Elasticsearch indexing performance and infrastructure reliability. The core of the problem addressed was the contention between high-volume, bursty reindexing operations and latency-sensitive real-time indexing, both vying for the same compute and storage resources. This contention led to performance degradation, stale search results for merchants, and significant operational overhead.

Watch on YouTube

Visual summary for Scaling Shopify's Search: Enhancing Elasticsearch Resilience With Kubernetes and KE... Leila Vayghan by Leila Vayghan
Visual summary for Scaling Shopify's Search: Enhancing Elasticsearch Resilience With Kubernetes and KE... Leila Vayghan by Leila Vayghan

Key moments

  1. 0:00 Introduction and talk agenda overview
  2. 2:00 Shopify's massive scale and BFCM statistics
  3. 2:40 Elasticsearch: distributed search engine at Shopify
  4. 4:00 Elasticsearch clusters hosted on Kubernetes in GCP
  5. 5:00 General data ingest pipeline and indexing definition
  6. 6:40 Two distinct indexing pipelines: real-time and reindex
  7. 8:30 Kafka's role in managing real-time and reindex pipelines

Scaling Shopify's Search: Enhancing Elasticsearch Resilience With Kubernetes and KEDA

Speakers: Leila Vayghan, Site Reliability Engineer, Shopify

Conference: KubeCon EU

YouTube: https://www.youtube.com/watch?v=r59IfCSmUBQ

Overview

In this insightful talk from KubeCon EU, Leila Vayghan, a Site Reliability Engineer in Shopify’s resiliency organization, detailed a critical project aimed at dramatically improving Elasticsearch indexing performance and infrastructure reliability. The core of the problem addressed was the contention between high-volume, bursty reindexing operations and latency-sensitive real-time indexing, both vying for the same compute and storage resources. This contention led to performance degradation, stale search results for merchants, and significant operational overhead.

Vayghan presented Shopify’s innovative solution: the strategic separation of real-time and reindex workloads onto dedicated Kubernetes node pools. Crucially, these reindex node pools are dynamically scaled using KEDA (Kubernetes Event-driven Autoscaling), responding to the actual demand of reindexing operations rather than remaining statically provisioned. This approach not only resolved the performance bottlenecks and improved system resilience but also yielded substantial infrastructure cost savings, showcasing a practical application of advanced Kubernetes and cloud-native autoscaling techniques in a massive e-commerce environment.

The talk is highly relevant for SREs, platform engineers, and architects managing large-scale distributed search systems, particularly those built on Elasticsearch and Kubernetes. It provides a blueprint for mitigating resource contention, optimizing cloud infrastructure costs, and enhancing the overall reliability and responsiveness of critical data pipelines within high-traffic, data-intensive applications like Shopify, which processes trillions of requests and billions in sales during peak events.

Background

▶ Watch: Introduction and talk agenda overview (0:00)

Shopify operates at an immense scale, serving over 3 million businesses across 175 countries and processing a significant portion of global e-commerce. As highlighted by Vayghan, during the Black Friday Cyber Monday (BFCM) week in 2024, Shopify managed 58 petabytes of data, served over a trillion edge requests, and processed more than 10 trillion database queries, culminating in $12 billion in global sales. Search is a fundamental component of this platform, enabling buyers to find products and merchants to manage orders and inventory efficiently.

At the heart of Shopify's search infrastructure is Elasticsearch, a distributed text search and analytics engine built on Lucene. Elasticsearch is chosen for its full-text search capabilities, scalability, and fault tolerance, making it well-suited for the e-commerce domain. Data in Elasticsearch is organized into indices, which are logically partitioned into shards (primary and replica) for distribution and resilience. Shopify runs its extensive fleet of Elasticsearch clusters—over 800 distinct clusters, some as large as 260 nodes, storing more than 3 petabytes of data—on managed Kubernetes clusters (GKE) within Google Cloud Platform (GCP), using a custom Kubernetes controller for deployment and maintenance.

The data pipeline feeding Elasticsearch is critical. Shopify Core, a large Rails monolith, underpins merchant shops. Updates from SQL databases, the primary data store, are streamed through Kafka topics and consumed by Kafka consumers, which then write these updates to the appropriate Elasticsearch indices in real time. This process is termed "indexing." Shopify employs two distinct indexing pipelines to handle different write profiles:

  1. Real-time pipeline: Indexes immediate changes, such as a merchant adding a new product, making it instantly searchable for buyers. This pipeline is highly sensitive to latency.
  2. Reindex pipeline: Handles structural changes to indices, such as adding or removing fields, or modifying analyzers. When such changes occur, an entirely new version of the index (.new index) must be rebuilt from SQL records. This is a bursty, high-volume operation that can involve rebuilding indices with hundreds of terabytes of data. For instance, the real-time indexing rate peaks at 90,000 documents per second, while the reindex rate can reach 500,000 documents per second.

For resilience, Shopify deploys its infrastructure across multiple GCP regions and availability zones. Inter-region data replication ensures that if one Elasticsearch cluster fails, query traffic can be failed over to another. Within a region, Elasticsearch's zone awareness feature, combined with Kubernetes node affinity rules and taints and tolerations, ensures that primary and replica shards are distributed across different availability zones and that Elasticsearch pods are scheduled only on dedicated GKE nodes, enhancing fault tolerance and enabling faster maintenance.

Key Findings

▶ Watch: Elasticsearch: distributed search engine at Shopify (2:40)

The central problem identified by Shopify was the resource contention arising from running both real-time and reindex workloads on the same Elasticsearch clusters and underlying GKE nodes. While an Elasticsearch index is split into many shards, and each shard has a replica, both the production index (receiving real-time writes and queries) and a newly created .new index (receiving heavy reindex writes) would share the same compute and storage resources. This led to several critical issues:

  • Performance Degradation: The heavy, bursty nature of reindex writes would significantly impact real-time writes, slowing them down. This resulted in stale or inaccurate search results, directly affecting merchant revenue.
  • Operational Overhead and Pager Fatigue: Site Reliability Engineers (SREs) were frequently paged due to real-time indexing delays, requiring manual intervention such as stopping and restarting reindexing processes, which was inefficient and disruptive.
  • Slowed Feature Rollouts: The need to pause or manage reindexes carefully to avoid impacting production meant that developers’ new features, which often necessitated index changes, could not be rolled out as quickly.
  • Excessive Costs: Shopify's initial approach to mitigate the problem was to overprovision Elasticsearch clusters to always be ready for peak reindex load. As the platform grew, this became prohibitively expensive and unsustainable.

To address these challenges, Shopify implemented a comprehensive solution centered on workload isolation and intelligent autoscaling. The key findings and results of this initiative were:

  • Workload Isolation: By separating real-time and reindex workloads onto dedicated GKE node pools, resource contention was eliminated. This ensured that heavy reindex operations no longer impacted the performance of latency-sensitive real-time indexing and production queries.
  • Significant Performance Improvements:
  • Reindexing of large indices (e.g., the orders index with 100 terabytes of data) became 40% faster.
  • Queued byte threads on the real-time nodes, a critical indicator of indexing backlog, dropped by 98% during reindexes.
  • This directly translated to improved real-time write performance and the elimination of pager fatigue related to indexing delays.
  • Developers were able to ship features requiring index changes much faster.
  • Substantial Cost Savings: Leveraging node autoscaling for the reindex node pools, Shopify achieved remarkable infrastructure cost reductions:
  • A 58% reduction in total CPU cores used by the real-time node pool.
  • A 15% reduction in memory used by the real-time node pool.
  • An overall 43% saving in infrastructure costs. This was achieved by no longer needing to overprovision the real-time nodes for reindex load and by scaling down the dedicated reindex nodes when not in use.

These findings underscore the critical importance of designing distributed systems with workload isolation and dynamic resource management in mind, especially for platforms operating at Shopify's scale and handling diverse data processing demands.

Technical Deep Dive

▶ Watch: Elasticsearch clusters hosted on Kubernetes in GCP (4:00)

The technical solution implemented by Shopify to achieve workload isolation and cost efficiency involved several key components and configurations within their Kubernetes and Elasticsearch environment.

The fundamental change was to create a dedicated infrastructure for reindex operations. This began with the introduction of a separate GKE node pool specifically for running reindex Elasticsearch pods. To enforce this separation, Kubernetes taints and tolerations were leveraged:

  • Taints were applied to the GKE nodes designated for reindexes. A taint essentially "repels" pods unless they have a matching toleration.
  • Tolerations were added to the reindex Elasticsearch pods, allowing them to be scheduled exclusively on the tainted reindex GKE nodes.
  • Similarly, real-time Elasticsearch pods were configured with tolerations to be scheduled on their respective real-time node pools, ensuring complete isolation.

Beyond node-level segregation, Elasticsearch's built-in shard allocation settings were crucial. These settings were configured to ensure that:

  • All production indices were hosted on Elasticsearch pods running on the real-time node pool.
  • Indices whose aliases ended with .new (the temporary indices created during a reindex) were exclusively hosted on Elasticsearch pods running on the reindex node pool.

Once a reindex was complete, and the .new index was ready to go live, an alias switch was performed. This action automatically triggered Elasticsearch to relocate the shards from the now-production index (which was previously the .new index) from the reindex node pool to the real-time node pool, making them searchable by clients and freeing up the reindex resources.

The dedicated reindex node pool, while providing isolation, presented a new cost challenge: it would be expensive to maintain a large, constantly provisioned pool for intermittent, bursty reindex operations. This is where KEDA (Kubernetes Event-driven Autoscaling) became indispensable.

KEDA is an open-source project designed to extend the capabilities of the Kubernetes Horizontal Pod Autoscaler (HPA). Instead of scaling based solely on CPU or memory utilization, KEDA allows workloads to scale based on metrics from a wide variety of external event sources, such as queue lengths, message counts, or custom metrics.

KEDA is implemented as a Kubernetes operator and uses Custom Resources (CRs). Its two main components are:

  • KEDA Controller: Watches for KEDA custom resources, specifically ScaledObjects, and manages the lifecycle of corresponding HPA resources.
  • Metrics Adapter: Collects metrics from configured event sources (called scalers) and provides them to the Kubernetes HPA, which then makes scaling decisions for application pods.

For Shopify's use case, the Prometheus scaler was chosen. A ScaledObject was configured to target the Elasticsearch StatefulSet running the reindex pods. Key settings included:

  • scaleTargetRef: Pointing to the reindex StatefulSet.
  • minReplicaCount: Set as low as 3, ensuring a minimal footprint when no reindexes are active.
  • maxReplicaCount: Set to around 200, allowing for significant scaling during large reindexes.
  • Prometheus Trigger: This trigger periodically queries a specified Prometheus server. The query designed for this solution was: number_of_shards / number_of_nodes.
  • Threshold: Set to 1. The ideal state for peak reindex performance was determined to be one shard per node. KEDA compares the query result to this threshold to decide whether to scale up or down. If the ratio is above 1, it indicates that nodes are overburdened, prompting a scale-up. If it's below 1 (or 0 when no shards are present), it signals a scale-down.

This setup allowed the reindex GKE node pool to dynamically expand when heavy reindexing started and contract back to its minimum size once the reindex was completed and shards were evacuated, leading to significant cost savings. The scaling behavior, including the rate of scaling, could also be defined within the ScaledObject configuration.

Demo / Proof of Concept

▶ Watch: Two distinct indexing pipelines: real-time and reindex (6:40)

While the talk did not feature a live demo, Leila Vayghan walked through a clear conceptual example demonstrating how KEDA-driven autoscaling for the reindex node pool operates. This illustrative walkthrough highlighted the system's dynamic response to reindexing events.

Scaling Up Scenario:

  1. Initial State: The reindex node pool is at its minimum size, typically three nodes, hosting no indices. The Prometheus query (number_of_shards / number_of_nodes) returns 0. KEDA, comparing this to its threshold of 1, takes no action as the workload is already at its minimum.
  2. Reindex Starts: A reindex operation begins for an index, for instance, one with four primary shards and four replica shards, totaling eight shards. Due to the configured Elasticsearch shard allocation rules, these eight shards of the new index are created on the pods running within the reindex node pool.
  3. KEDA Detects Load: At this point, the Prometheus query calculates 8 shards / 3 nodes, resulting in a value greater than 1. This signals to KEDA that the current node capacity is insufficient for optimal performance.
  4. Workload Scales Up: KEDA initiates a scale-up of the reindex StatefulSet. This action, in turn, causes the underlying GKE node pool to expand. As new nodes become available, Elasticsearch automatically rebalances the shards across them. The scaling continues until the Prometheus query result approaches the threshold of 1 (e.g., 8 shards / 8 nodes = 1).
  5. Steady State: Once the ratio is 1, KEDA recognizes optimal resource allocation and ceases further scaling actions, maintaining the expanded node pool until the reindex completes. For very large indices, such as those with thousands of shards, scaling up to the maximum replica count (e.g., 200 nodes) might still mean each node hosts multiple shards (e.g., 10 shards/node), but this is still a vast improvement over contention. The speaker noted that scaling up a large index took approximately 30 minutes, which was acceptable for a maintenance operation.

Scaling Down Scenario:

  1. Reindex Completes: After the reindex is finished, the alias for the .new index is flipped, promoting it to the production index. Elasticsearch then begins the process of evacuating all shards from the reindex node pool, relocating them to the real-time node pool.
  2. KEDA Detects Idleness: As shards are moved off the reindex nodes, the Prometheus query (number_of_shards / number_of_nodes) eventually drops below the threshold of 1, eventually returning 0 when all shards have been evacuated.
  3. Workload Scales Down: KEDA detects this state of idleness and triggers a scale-down of the reindex workload. Consequently, the underlying GKE node pool shrinks back to its minReplicaCount (e.g., three nodes).
  4. Data Integrity During Scale Down: A question from the audience clarified how data loss is prevented during scale-down. Vayghan explained that the Prometheus query is the safeguard: it will not return 0 until all shards have been successfully moved off the nodes. KEDA will only scale down the node pool to its minimum once there are zero shards remaining on the reindex nodes. For the largest 100TB index, the full evacuation and scale-down process took about two hours.

This detailed conceptual walk-through effectively served as a proof of concept, illustrating the system's ability to dynamically adapt to varying reindex loads, optimize resource utilization, and ensure data integrity throughout the process.

Defensive Implications

▶ Watch: Kafka's role in managing real-time and reindex pipelines (8:30)

The solution presented by Shopify offers several critical defensive implications for organizations managing large-scale distributed systems, particularly those relying on Elasticsearch and Kubernetes:

  • Workload Isolation as a Primary Defense: The most significant takeaway is the importance of isolating resource-intensive, bursty workloads (like reindexing or batch processing) from latency-sensitive, critical production workloads. This architectural separation prevents "noisy neighbor" issues and safeguards the performance and reliability of core services. Defenders should actively identify such conflicting workloads within their environments and design dedicated resource pools for them.
  • Leveraging Kubernetes Taints and Tolerations for Segregation: Kubernetes' built-in mechanisms like taints and tolerations are powerful tools for enforcing strict workload separation at the node level. By applying taints to nodes and corresponding tolerations to pods, organizations can ensure that specific types of workloads are scheduled only on designated infrastructure, providing a robust layer of isolation.
  • Embracing Event-Driven Autoscaling (KEDA): Traditional autoscaling based solely on CPU/memory often falls short for workloads driven by external events or custom metrics. KEDA provides a flexible framework for event-driven autoscaling, enabling cost optimization by scaling resources up only when needed and down when idle. Defenders should explore KEDA for intermittent high-load applications, queue-based systems, or any scenario where scaling decisions can be tied to specific business or application events.
  • Strategic Use of Elasticsearch Shard Allocation: Elasticsearch's native shard allocation settings are crucial for managing data distribution and resource utilization within a cluster. By configuring these settings to direct specific indices or index types to particular node groups, organizations can optimize performance, ensure resilience (e.g., zone awareness), and facilitate maintenance operations.
  • Designing Resilient Data Pipelines: The talk highlighted the necessity of designing data pipelines with distinct paths for different write profiles (real-time vs. batch/reindex). This forethought in pipeline architecture can prevent bottlenecks and ensure that diverse data processing requirements are met without impacting each other.
  • Proactive Monitoring with Custom Metrics: The success of KEDA's implementation hinged on a well-crafted Prometheus query (number_of_shards / number_of_nodes). Defenders should develop custom metrics that accurately reflect the health and load of their critical components, especially for workloads that don't fit standard CPU/memory scaling patterns. These custom metrics can then feed into advanced autoscaling solutions.
  • Cost Optimization through Dynamic Provisioning: Overprovisioning infrastructure as a default solution for peak loads is unsustainable. By adopting dynamic provisioning strategies like KEDA-based autoscaling, organizations can significantly reduce cloud infrastructure costs while maintaining or improving performance and reliability. This shifts the defensive posture from "always ready" to "ready when needed."

Key Takeaways

  • Workload Isolation is Paramount: Separating bursty, resource-intensive reindex operations from latency-sensitive, real-time indexing is critical for maintaining performance and reliability in large-scale distributed search systems.
  • Kubernetes Taints and Tolerations Enforce Segregation: Kubernetes provides effective mechanisms (taints on nodes, tolerations on pods) to ensure that specific workloads are scheduled only on dedicated infrastructure, preventing resource contention.
  • KEDA Enables Cost-Effective Event-Driven Autoscaling: KEDA extends the HPA to scale workloads based on custom event sources, allowing for dynamic provisioning of resources (e.g., GKE node pools) only when needed, leading to significant cost savings.
  • Custom Prometheus Metrics Drive Intelligent Scaling: A custom Prometheus query (number_of_shards / number_of_nodes) with a defined threshold (e.g., 1) can effectively inform KEDA's scaling decisions for Elasticsearch, optimizing the shard-to-node ratio for performance.
  • Significant Performance and Cost Benefits: Shopify achieved a 40% faster reindexing, a 98% reduction in queued byte threads on real-time nodes, eliminated pager fatigue, and realized 43% infrastructure cost savings through this approach.
  • Elasticsearch Features for Resilience: Leveraging Elasticsearch's built-in features like zone awareness and shard allocation settings, combined with Kubernetes node affinity, is crucial for designing resilient and efficient distributed search systems.

About the Speaker(s)

Leila Vayghan is a Site Reliability Engineer (SRE) at Shopify. She is an integral part of Shopify's resiliency organization, a team whose core responsibility is to ensure that the Shopify platform remains consistently available, robust, and reliable for its vast global network of merchants and buyers. Her work focuses on tackling complex infrastructure challenges to maintain the high performance and uptime demanded by one of the world's leading commerce platforms.

Reviews

Dr. Zero (Offensive Security Researcher) — STRONG ACCEPT

This talk from Shopify's Leila Vayghan presents a highly effective and technically sound solution to a common pain point in large-scale Elasticsearch deployments: resource contention between real-time and reindex workloads. By leveraging dedicated Kubernetes node pools, KEDA for event-driven autoscaling based on a clever custom Prometheus metric (shards per node), and Elasticsearch's native allocation rules, Shopify not only eliminated critical performance bottlenecks and pager fatigue but also achieved substantial infrastructure cost savings. It's a blueprint for any SRE or platform engineer dealing with similar challenges.

Heather Calloway (CISO) — STRONG ACCEPT

Leila Vayghan's session on scaling Shopify's search is a compelling case study in operational resilience and cost optimization, demonstrating how strategic workload isolation and event-driven autoscaling fundamentally improve critical business functions. This isn't just a technical deep dive; it's a blueprint for managing systemic risk, ensuring business continuity, and achieving significant infrastructure efficiencies in large-scale environments. It directly addresses the institutional conditions that lead to performance bottlenecks and excessive costs, offering a clear, actionable path forward for platform leaders and engineers.

→ Top-rated talks at KubeCon + CloudNativeCon Europe 2025

All talks from KubeCon + CloudNativeCon Europe 2025