DEV Community

Anusha Renangi
Anusha Renangi

Posted on

OpenSearch Under the Hood: How It Actually Works

I have been using OpenSearch for querying and aggregating data at work. Initially, I was just maintaining the code, but after working with it for some time, I started wondering what actually happens behind those queries.

When I search for something, does OpenSearch search the original JSON documents?

And when an aggregation runs across multiple shards, how does OpenSearch get one final result?

So I decided to understand OpenSearch from the inside rather than treating it as just another API that returns search results.

First, why do we need OpenSearch?

Let’s say an application has millions of documents.

We may want to do things like:

  • Search text
  • Filter documents
  • Sort by fields
  • Group documents
  • Calculate aggregations

Doing this efficiently over a large amount of data is not the same as simply storing JSON and scanning every document whenever a query comes in.

OpenSearch is designed for search and analytics over large datasets.

To understand how it does this, let’s start with how the data is organized.

Cluster → Nodes → Indexes → Shards → Documents

These terms are easy to mix up when we first start working with OpenSearch.

Here is the simplest way I understand them now.

Cluster

A cluster is a collection of OpenSearch nodes that work together.

A cluster can contain multiple indexes, and each index can be divided into multiple shards that are distributed across the cluster’s nodes.

You can loosely compare an OpenSearch cluster to a database system in the relational-database world, but they are not exactly the same thing.

Node

A node is an OpenSearch server that is part of the cluster.

Nodes provide the compute and storage where shards are placed.

A cluster can have multiple nodes, and a node can hold multiple shards from one or more indexes.

Index

An index is a logical collection of related documents.

It is somewhat similar to a table in a relational database.

For example:
orders
products
users

An index can contain a large number of documents, so OpenSearch divides the index into shards.

Shard

A shard is a partition of an index.

These shards can be distributed across different nodes in the cluster.

Each shard contains a subset of the documents belonging to the index.

Document

A document is the actual piece of data stored in the index.

For example:
{
“name”: “Laptop”,
“region”: “India”,
“amount”: 75000
}

Opensearch High level Storage Structure

What actually lives inside a shard?

A shard is backed by Apache Lucene.

And OpenSearch doesn’t simply keep our JSON documents somewhere and scan them whenever we search.

It maintains different structures that are useful for different operations.

Shard Internal Structure

Inverted Index

Let’s start with search.

Suppose we have these documents:

Document 1: “I love OpenSearch”
Document 2: “OpenSearch is fast”
Document 3: “I love search”

If we stored only the documents in their original form, one simple way to search would be to look through the documents and check which ones contain the word we are searching for.

An inverted index changes the way we look at the data.

Instead of:
Document → Words
we have something more like:
Word → Documents containing that word
For example:
I → 1, 3
love → 1, 3
OpenSearch → 1, 2
is → 2
fast → 2
search → 3

Now if the query is: “OpenSearch”

the inverted index can tell us: OpenSearch → Document 1, Document 2

We don’t have to think of search as scanning every JSON document anymore.

Doc Values

Search is only one type of operation.

Suppose I want to run: SUM(amount)

or sort documents by: amount

For these operations, we need efficient access to field values.

A simplified way to think about them is:

Document amount

Doc 1 75000
Doc 2 45000
Doc 3 90000
Doc 4 30000

_source

Then there is the original JSON.

Suppose we indexed:

{
“name”: “Laptop”,
“region”: “India”,
“amount”: 75000
}

When we retrieve the document, we generally want to get something close to what we originally sent.

That’s what _source is for.

“_source”: {
“name”: “Laptop”,
“region”: “India”,
“amount”: 75000
}
It isn’t simply storing JSON documents and searching through them.

Different data structures are maintained because different operations need different ways of accessing the data.

What happens when we actually send a query?

Now let’s say our orders index has three shards.

Suppose we send a search query.

The request reaches a node, which coordinates the request.

The query is sent to the relevant shards.

Each shard searches its own local data.

Then the results come back.

The coordinating node merges the shard results and produces the final response.

Opensearch Query Process

But does OpenSearch always give us the exact global result?

Suppose we run an aggregation:

Group documents by region
Return the top 2 regions by document count

The documents are distributed across multiple shards.

Each shard does its own local processing and sends results back to the coordinating node.

The coordinating node then combines those results.

A value that isn’t among the top candidates from one shard can still become important when the results from all shards are combined.

For example, a region might have a lower count on each individual shard, but when those counts are combined across all shards, that region could have a high global count.

So how can OpenSearch confidently say, “These are the top results across the entire cluster”?

That’s where distributed aggregation gets tricky and sometimes, approximate.

In the next article, we’ll see exactly why this happens and how OpenSearch improves the accuracy of these results.

Top comments (0)