Reducing Join Latency in YugabyteDB with Colocated Tables
Stop wasting milliseconds on network shuffles. Learn how YugabyteDB Colocated Tables keep related data on the same node to accelerate joins in multi‑tenant apps.
08 Jul 2025, 23:49 UTC

The Distributed Join Problem
\nIn a distributed SQL database like YugabyteDB, data is split into shards called tablets and spread across multiple nodes. While this allows for massive scaling, it introduces a performance tax when you perform joins. By default, if a parent table and a child table are sharded independently, the database must perform a \\"network shuffle.\\" This means the node coordinating the query must pull data from various nodes across the cluster to assemble the final result, adding significant network latency to every request.
\nThe takeaway: If your application frequently joins a parent table to its children using a shared ID, you can eliminate this network overhead by using Colocated Tables. This ensures that related rows from different tables are physically stored on the same tablet on the same node.
\n\nHow Colocation Works
\nColocation is a physical storage strategy. Normally, YugabyteDB shards tables based on their own primary keys. With colocation, you designate a parent table and one or more child tables. The child tables are then sharded based on the parent's primary key rather than their own.
\nThis is particularly powerful for multi-tenant architectures. If you have a tenants table and a orders table, colocation ensures that all orders for Tenant A live on the same node as the Tenant A profile. When you join these tables, the database performs a local join, avoiding the network entirely.
Implementing Colocated Tables
\nTo implement colocation, the child table must include the parent's primary key as the first part of its own primary key. This creates the physical link required for the database to group the data together.
\n\nExample Configuration
\nAssume you are running YugabyteDB 2.x. Run these commands via ysqlsh (the YSQL shell) as a user with CREATE permissions.
-- 1. Create the parent table
CREATE TABLE customers (
customer_id INT PRIMARY KEY,
name TEXT
);
-- 2. Create the child table as a colocated table
-- The child table MUST include the parent's PK as the first column of its PK
CREATE TABLE orders (
customer_id INT,
order_id INT,
amount DECIMAL,
PRIMARY KEY (customer_id, order_id),
-- This clause tells YugabyteDB to store data with the parent
COLOCATED WITH (customers)
) ;
\n\nVerifying the Performance Gain
\nTo verify that colocation is working, use the EXPLAIN ANALYZE command. This reveals whether the database is performing a distributed scan or a local one.
EXPLAIN ANALYZE \nSELECT c.name, o.amount \nFROM customers c \nJOIN orders o ON c.customer_id = o.customer_id \nWHERE c.customer_id = 123;
\n\nWhat to look for: In a non-colocated setup, you will see references to remote scans or network shuffles. In a colocated setup, the plan should indicate that the join is happening locally on the node where the data resides.
\n\nTrade-offs and Limitations
\nColocation is not a \\"silver bullet\\" for all performance issues. There are three primary risks to consider:
\n- \n
- Hot Tablets: If one parent key (e.g., a single massive enterprise customer) has millions of child rows and generates huge traffic, that specific node will become a bottleneck. You cannot spread a single parent's children across multiple nodes. \n
- Query Constraints: Colocation only helps if your query includes the sharding key. If you join
ordersandcustomersusing a column other thancustomer_id, the database must still perform a full distributed scan. \n - Schema Rigidity: Once a table is created as colocated, you cannot simply \\"toggle\\" this off. Changing the sharding strategy requires creating a new table and migrating the data. \n
Decision Summary
\nUse colocated tables when you have a clear one-to-many relationship and your most frequent, latency-sensitive queries always filter by the parent ID. If your data access patterns are unpredictable or you have extreme data skew (some parents are 1,000x larger than others), stick with standard sharding to maintain better load balancing across your cluster.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.