This Tip is the second in a two-part series looking at the performance impact of large multi-column IN lists in YugabyteDB.
In the first Tip, Reduce YSQL Planning Time for Large Multi-Column IN Lists with UNNEST, we focused on the planning side of the problem. As more tuple pairs are added to an IN list, the SQL expression grows, the number of bind parameters increases, and the optimizer may have more predicates, index paths, and plan branches to evaluate.
We showed how passing the lookup values as arrays and expanding them with UNNEST keeps the SQL statement at a fixed shape. Instead of generating a new statement with 20, 50, 100, or more bind parameters, the application can continue using the same two array parameters regardless of batch size. That can make planning more predictable and prepared-statement reuse much easier.
But planning time is only half of the story.
The way those lookup values are expressed can also change how YugabyteDB executes the query against its distributed storage layer.
A query can use the right primary key, return only a handful of rows, and still perform more distributed work than expected.
Consider a multi-column IN lookup:
SELECT item_id,
description,
category_id,
state_code,
lookup_hash,
version,
payload
FROM lookup_item
WHERE (lookup_hash, category_id) IN (
('0000000000000000000000000000000000001001', 1),
('0000000000000000000000000000000000001002', 2),
('0000000000000000000000000000000000001003', 3),
('0000000000000000000000000000000000001004', 4),
('0000000000000000000000000000000000001005', 5),
('0000000000000000000000000000000000001006', 6),
('0000000000000000000000000000000000001007', 7),
('0000000000000000000000000000000000001008', 8)
);
All eight predicates target columns at the beginning of the primary key, so at first glance this looks like an ideal set of point lookups.
But in a distributed database, the important question isn’t only:
- How many rows are returned?
It is also:
- How many storage requests are required to retrieve them?
YugabyteDB’s EXPLAIN (ANALYZE, DIST) makes this visible by exposing distributed-storage runtime statistics such as storage read requests, read operations, rows scanned, and storage execution time.
In the eight-row test used in this Tip, the tuple IN query selected a BitmapOr plan containing eight separate Bitmap Index Scan branches. That plan ultimately required 9 storage read requests to return the eight requested rows.
By rewriting the same lookup as a join using UNNEST arrays, or a VALUES relation, YugabyteDB selected a YB Batched Nested Loop Join. In the same test, the eight logical lookups were handled with just 1 storage read request.
UNNEST. This Tip focuses on
execution overhead: using that same join-oriented query shape to give
YugabyteDB the opportunity to batch multiple point lookups into fewer distributed
storage requests.
Build the Test Table
We’ll use the same anonymized schema from the first Tip in this series:
DROP TABLE IF EXISTS lookup_item;
CREATE TABLE lookup_item (
item_id bigint NOT NULL,
category_id smallint NOT NULL,
version smallint NOT NULL,
lookup_hash varchar(40) NOT NULL,
description varchar(150),
payload text NOT NULL,
state_code char(1) NOT NULL,
PRIMARY KEY ((lookup_hash) HASH, category_id ASC, state_code ASC)
);
Populate it with 10,000 sample rows:
INSERT INTO lookup_item
(item_id, category_id, version, lookup_hash, description, payload, state_code)
SELECT
g,
(g % 100)::smallint,
1,
lpad(g::text, 40, '0'),
'Sample Item ' || g,
'{"value":"sample"}',
'A'
FROM generate_series(1, 10000) AS g;
ANALYZE lookup_item;
The reproducible test used this schema and data set.
Baseline: Multi-Column IN
Run the eight-pair lookup:
EXPLAIN (ANALYZE, DIST)
SELECT item_id,
description,
category_id,
state_code,
lookup_hash,
version,
payload
FROM lookup_item
WHERE (lookup_hash, category_id) IN (
('0000000000000000000000000000000000001001', 1),
('0000000000000000000000000000000000001002', 2),
('0000000000000000000000000000000000001003', 3),
('0000000000000000000000000000000000001004', 4),
('0000000000000000000000000000000000001005', 5),
('0000000000000000000000000000000000001006', 6),
('0000000000000000000000000000000000001007', 7),
('0000000000000000000000000000000000001008', 8)
);
The interesting portion of the plan is:
YB Bitmap Table Scan on lookup_item
-> BitmapOr
-> Bitmap Index Scan on lookup_item_pkey
-> Bitmap Index Scan on lookup_item_pkey
-> Bitmap Index Scan on lookup_item_pkey
-> Bitmap Index Scan on lookup_item_pkey
-> Bitmap Index Scan on lookup_item_pkey
-> Bitmap Index Scan on lookup_item_pkey
-> Bitmap Index Scan on lookup_item_pkey
-> Bitmap Index Scan on lookup_item_pkey
Each individual Bitmap Index Scan reported one storage table read request, and the final bitmap table scan also required a storage request. The overall query therefore reported:
Storage Read Requests: 9
Storage Read Ops: 12
Storage Rows Scanned: 16
Execution Time: 4.427 ms
Here’s the execution shape visually:
YugabyteDB’s DIST option exposes distributed-storage counters at the query layer. In particular, Storage Read Requests is the total number of table and index reads across the plan, while node-level counters show where those requests occurred.
BitmapOr plan selected for this schema and data set. Different statistics,
indexes, versions, or batch sizes can produce a different plan. Always verify the actual
behavior with EXPLAIN (ANALYZE, DIST).
Rewrite the Lookup as a Join
Instead of expressing every lookup key as a separate branch of a multi-column IN predicate, we can present the values as rows and join them to the target table.
There are two convenient ways to do this.
Option 1: UNNEST Arrays
The first option is to pass the lookup keys as arrays and expand them with UNNEST:
EXPLAIN (ANALYZE, DIST)
SELECT l.item_id,
l.description,
l.category_id,
l.state_code,
l.lookup_hash,
l.version,
l.payload
FROM lookup_item l
JOIN UNNEST(
ARRAY[
'0000000000000000000000000000000000001001',
'0000000000000000000000000000000000001002',
'0000000000000000000000000000000000001003',
'0000000000000000000000000000000000001004',
'0000000000000000000000000000000000001005',
'0000000000000000000000000000000000001006',
'0000000000000000000000000000000000001007',
'0000000000000000000000000000000000001008'
]::varchar[],
ARRAY[
1,
2,
3,
4,
5,
6,
7,
8
]::smallint[]
) AS batch(hash_val, category_val)
ON l.lookup_hash = batch.hash_val
AND l.category_id = batch.category_val;
This time, YugabyteDB selected:
YB Batched Nested Loop Join
. -> Function Scan on batch
. -> Index Scan using lookup_item_pkey
The index scan also shows the lookup keys being passed as a batch:
Index Cond:
ROW(lookup_hash, category_id) =
ANY (
ARRAY[
ROW(batch.hash_val, batch.category_val),
...
]
)
The distributed statistics changed to:
Storage Read Requests: 1
Storage Read Ops: 4
Storage Rows Scanned: 8
Execution Time: 2.293 ms
In this test, the eight requested rows were retrieved using a single storage read request.
What Changed?
The key change is the execution strategy:
YB Batched Nested Loop Join
A traditional nested loop can probe the inner side once for each tuple received from the outer side. YugabyteDB’s Batched Nested Loop Join instead groups outer-side tuples and sends a batch of lookup keys to the inner relation. YugabyteDB’s own BNL example shows the inner index condition expressed as ANY (ARRAY[...]), allowing multiple keys to be handled in one inner scan rather than repeatedly looping over the inner side.
Conceptually:
Storage Read Requests vs. Storage Read Ops
One part of the DIST output deserves special attention.
The UNNEST test reported:
Storage Read Requests: 1
Storage Read Ops: 4
Storage Rows Scanned: 8
The important point is that Storage Read Requests and Storage Read Ops are different counters.
In this plan, one storage read request was able to carry multiple read operations:
The optimization we care about here is the reduction in requests. YugabyteDB documents Storage Read Requests as the sum of table and index read requests across the plan; the node-level read-request counters correspond to RPC round trips to the local YB-TServer.
Option 2: VALUES
Arrays aren’t required to obtain the execution benefit.
The same lookup can be represented as a VALUES relation:
EXPLAIN (ANALYZE, DIST)
SELECT l.item_id,
l.description,
l.category_id,
l.state_code,
l.lookup_hash,
l.version,
l.payload
FROM lookup_item l
JOIN (
VALUES
('0000000000000000000000000000000000001001', 1::smallint),
('0000000000000000000000000000000000001002', 2),
('0000000000000000000000000000000000001003', 3),
('0000000000000000000000000000000000001004', 4),
('0000000000000000000000000000000000001005', 5),
('0000000000000000000000000000000000001006', 6),
('0000000000000000000000000000000000001007', 7),
('0000000000000000000000000000000000001008', 8)
) AS batch(hash_val, category_val)
ON l.lookup_hash = batch.hash_val
AND l.category_id = batch.category_val;
This also selected:
YB Batched Nested Loop Join
with:
Storage Read Requests: 1
Storage Read Ops: 4
Storage Rows Scanned: 8
Execution Time: 1.925 ms
For the execution-side optimization covered in this Tip, both UNNEST and VALUES gave the optimizer a join shape that selected a YB Batched Nested Loop Join.
Side-by-Side Results
| Query Shape | Plan | Read Requests | Read Ops | Rows Scanned | Execution |
|---|---|---|---|---|---|
Tuple IN
|
BitmapOr
|
9 | 12 | 16 | 4.427 ms |
UNNEST
|
YB Batched Nested Loop | 1 | 4 | 8 | 2.293 ms |
VALUES
|
YB Batched Nested Loop | 1 | 4 | 8 | 1.925 ms |
These are the actual results from the anonymized reproducible test.
Why Fewer Storage Requests Matter
The cost of a lookup isn’t determined solely by the number of rows returned.
YSQL executes against YugabyteDB’s distributed storage layer, and requests from the query layer to storage have an RPC cost. YugabyteDB’s EXPLAIN (ANALYZE, DIST) output exposes those request counts specifically because they are useful when understanding distributed execution.
A few extra requests may not be noticeable in a small local test.
But as network latency, concurrency, and batch size increase, repeatedly issuing separate storage requests can become increasingly expensive.
That is why even a query retrieving a relatively small set of rows can benefit from batching.
Batched Nested Loop Configuration
Two YSQL configuration parameters are useful when investigating BNL:
SHOW yb_enable_batchednl;
SHOW yb_bnl_batch_size;
The documented defaults are:
yb_enable_batchednl = true
yb_bnl_batch_size = 1024
yb_enable_batchednl controls whether the planner can use Batched Nested Loop Join. yb_bnl_batch_size determines the size of the tuple batch taken from the outer side; setting the batch size to 1 effectively prevents BNL from being considered as a plan candidate.
yb_bnl_batch_size is 1024.
Larger is not automatically better. Batch size can affect request size,
memory usage, and execution characteristics. Start with the default and
change it only after testing a representative workload.
Verify the Plan
Rewriting a query as a join does not guarantee that the optimizer will always select a Batched Nested Loop Join.
The optimizer still chooses the execution strategy based on the query, available indexes, statistics, and other costing decisions.
Always verify with:
EXPLAIN (ANALYZE, DIST)
Look for:
YB Batched Nested Loop Join
and compare:
Storage Read Requests
Storage Read Ops
Storage Rows Scanned
Execution Time
The YugabyteDB BNL documentation demonstrates exactly this type of verification: the batched plan shows an ANY (ARRAY[...]) index condition and fewer inner-side loops than the equivalent ordinary nested loop.
UNNEST and VALUES produced YB Batched Nested Loop plans in this test,
but the optimizer still chooses the execution plan. Use
EXPLAIN (ANALYZE, DIST) against your actual schema and representative batch sizes
to confirm that batching is occurring.
UNNEST or VALUES?
For execution performance, either approach can work well when the optimizer selects BNL.
| Approach | Advantage |
|---|---|
UNNEST
|
Fixed SQL shape, two parameters, and an excellent fit for prepared statements |
VALUES
|
Straightforward for applications or query builders that already generate row-value lists |
Because the first Tip in this series demonstrated the planning advantages of keeping the statement fixed at two array parameters, UNNEST is particularly attractive when the application can easily pass arrays.
VALUES remains a useful alternative when changing application parameter handling is difficult.
Final Takeaway
The important optimization isn’t simply changing SQL syntax.
It is changing the execution shape:
BitmapOr plan and many
Storage Read Requests, consider expressing the lookup as a join using
UNNEST arrays or a VALUES relation. If YugabyteDB selects
YB Batched Nested Loop Join, multiple logical point lookups can be carried
in fewer storage requests, reducing RPC overhead and potentially lowering execution latency.
Have Fun!
I’m not normally a big BBQ guy… but apparently I just hadn’t been eating the right BBQ. 😂
After a few days of our YugabyteDB QBR here in Dallas, the group of us went to Terry Black’s BBQ tonight … and wow.
The food was fantastic! 🔥🥩
I was already looking forward to moving to Dallas next month, but after tonight, I’m even more excited.
Pretty sure Terry Black’s is about to become a regular DoorDash order at the new house. 😂🤠
