分片 Postgres 查询的生命周期
Postgres 分片到千台服务器后,一条 SELECT 如何穿越认证、wire protocol 与分布式查询规划器,Neki 路由器用 Go 复刻 SCRAM-SHA-256 并同时支持 Simple 与 Extended 协议。
中文
复制

如果你用过 Postgres,并且觉得“哇,这软件真简单,每一部分我都完全看得懂”,那你还没见过查询规划器。现在再想想,把这个数据库分片到一千台服务器上。你要怎么处理一个查询?你需要一套系统,能复刻 Postgres 的认证机制、wire protocol 和解析器,还要有能感知分片的分布式查询规划器、对所有服务器故障场景的优雅处理,以及能绕开 Postgres 每连接一进程架构的连接池。很简单,对吧?让我们跟着一个 Postgres 查询,走一遍这个优雅又精巧复杂的分布式系统的每一层。这样我们就能理解,要让横跨数千台服务器的大规模 Postgres 分片部署看起来像一台 Postgres 服务器,究竟需要做些什么。真实数据库可能有几百张表,但这个例子里我们把 schema 保持简单:两张表,分布在四个分片上。
CREATE TABLE customers (
id BIGINT PRIMARY KEY,
name TEXT NOT NULL,
email TEXT NOT NULL,
country TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE orders (
id BIGINT PRIMARY KEY,
customer_id BIGINT NOT NULL,
total NUMERIC(12,2) NOT NULL,
status TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
我们要跟踪的这个查询也很简单,一条 SELECT 语句,查询所有近期订单的客户名、订单金额和下单日期。
SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1;
在执行这个查询之前,有一件关于分片的重要事情需要理解:分片的意义在于让数据库突破单台服务器的限制。这意味着我们现在身处一个分布式系统。行数据可以分布在不同的服务器上,所以我们需要规则来决定数据存在哪里、怎么找到它。在这个例子里,我们选择 customers.id 作为 customers 的分片键,选择 orders.id 作为 orders 的分片键。这些主键是自然的起点,虽然有点天真,后面很快就会看到。每次要存一行 customers 或 orders,系统会取出 id,算一个哈希,再用这个哈希决定存到哪里。注意客户和订单是如何分散在许多服务器上的。这让分片数据库几乎可以无限扩展,但也带来了不少有意思的工程难题。用 Follow Maja 按钮看看一个客户和他们的订单是如何定位分片的。眼尖的人会发现,它们并不在同一分片上。先记住这一点,很快就用得上。让我们开始这段穿越 Neki 分片数据库的旅程。
认证
应用要发送查询,必须先通过认证。有些分片数据库把路由交给应用层处理:由应用自己挑选合适的 Postgres 实例并建立连接。我们的路由器把这份复杂度藏了起来。对你的应用来说,它_就是_ Postgres,负责认证、连接以及把查询分发到各个服务器。为此,我们用 Go 完整实现了 Postgres 的认证流程,包括 SCRAM-SHA-256。TCP 建立、SSL 协商和 TLS 握手之后,客户端发来一条 startup 消息,其中包含 Postgres 角色、数据库名以及请求的会话设置。路由器检查自己的认证规则,随后 SCRAM 交换开始。
路由器发出 challenge。驱动用密码和 challenge 算出 proof,路由器拿它和存储的 verifier 比对。接着路由器返回一个 signature 供驱动验证。这样路由器就能确认客户端知道正确的密码,而密码本身无需在网络上传输。
最后,路由器检查数据库访问权限、应用会话设置,并发送 startup 响应。ReadyForQuery 告诉客户端可以发送查询了。
协议
安全连接建立之后,应用就可以把查询交给它的 Postgres 驱动了。很多人不知道,Postgres 的 wire protocol 其实有两种发送查询的方式:Simple 和 Extended。你可能也不知道,不同的 Postgres 驱动,甚至同一个驱动里的不同设置,发送同一条查询的方式都可能不同。路由器的目标是与任何 Postgres 客户端兼容,所以两种都支持。
Postgres Simple 与 Postgres Extended 协议
使用 simple 协议时,客户端把整条 SQL 语句连同时间戳一起放在一条 Query 消息里发送。
Query
SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= '2026-09-09 00:00:00+00';
在 extended 协议中,我们的查询分 5 个部分传递。
Parse
statement: ""
query: SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1;
Bind
statement: ""
portal: ""
parameters: ["2026-09-09 00:00:00+00"]
Describe
portal: ""
Execute
portal: ""
Sync
- Parse 要求路由器准备 SQL,此时
$1仍然留在原位。这会创建一个未命名的 prepared statement,尽管应用并没有显式执行过PREPARE。 - Bind 提供时间戳,并为这次执行创建一个 portal。portal 跟踪一条查询的执行状态,把 prepared statement 与它的参数值和请求的结果格式绑定在一起。
- Describe 询问这个 portal 会返回哪些列、以什么格式返回。
- Execute 让路由器执行该 portal。
- Sync 标记这一批的结束。
对于临时查询(比如来自 psql 这类命令行),简单协议提供的是一种 简单 的交互方式:发送完整的 SQL,拿回结果。它还能在一次请求中接受多条 SQL 语句。而扩展协议让 Postgres 驱动可以分别控制准备、参数绑定和执行。当应用程序反复用不同的值执行同一条语句时,这种分离就很有用——复用已准备好的语句,就不必反复解析 SQL、反复解析其中的名称和类型。我们的 router 并不是真正的 Postgres 进程,所以我们在实现这套协议行为时加入了一些自己的优化。不过那是另一篇文章的话题了;从应用程序的角度看,我们的 router 就是 Postgres。
解析
无论是简单协议还是扩展协议,最终拿到的都是一串必须处理的 SQL。直接与 Postgres 通信时,Postgres 的解析器会先对查询做词法分析,构建抽象语法树(AST),再交给 Postgres 查询规划器。AST 表示查询的各个组成部分:引用了哪些表、如何连接、施加了什么过滤条件、要返回哪些列。在 router 这边,我们必须做同样的解析和建树工作,并且与 Postgres 的语法精确一致——支持的语法要一致,拒绝的语法也要一致。我们用大约 18,000 行 Go 实现了它,并配以严格的要求和测试,确保我们的版本性能不逊于、甚至优于 Postgres 自身的 C 实现。它就像编译器的解析器一样,把源文本变成抽象语法树,让 router 在推敲如何执行查询时,有一个可以检视和改写的结构。有些情况下 router 根本不需要解析 SQL。为了性能,router 的目标是尽可能避免做无用功。它会先检查是否有可复用的缓存计划,命中缓存就能跳过解析和规划。但这一次没有命中。router 必须解析,随后这棵 AST 被交给查询规划器。
规划
到目前为止,我们已经做了大量工作,而完成的不过是把一条 SQL 查询变成了抽象语法树。分片还没考虑,Postgres 还没碰,那个问题也还没回答:“怎么才能快速执行这条查询?”
一个分片数据库系统好不好,取决于查询规划器做得有多扎实。坐稳了,你马上就会知道,为了让 Neki 尽可能少干活,我们花了多少功夫。它懒到家了。
回头看看我们一路跟进的那条原始查询:
SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1;
如果只有单个 Postgres 节点,执行它需要做什么?没有分片、没有代理、没有查询路由器(暂时)。Postgres 有自己的查询规划器,从解析器拿到 AST 之后,它得决定这些事:
- 用内部统计信息估算有多少订单匹配
created_at >= $1。 - 决定每张表的读取方式:顺序扫描、索引扫描、位图扫描还是仅索引扫描。
- 选择连接方向,以及用嵌套循环、哈希连接还是归并连接。
- 检查
orders.created_at上的索引能否帮助过滤订单。
这还只是我们这条简单查询所需的部分决策。查询越复杂,规划器要考虑的东西越多,包括更多的连接、聚合、子查询和排序。
别忘了,这是分片数据库,每张表的行都分散在各个分片上。这里我们有 4 个分片,但本节讲的原理,无论 4 个分片存 1 TB,还是 400 个分片存 1 PB,都成立。
要决定每一行落到哪个分片,我们指定某一列作为分片键。前面我们给 customers 选了 customers.id,给 orders 选了 orders.id。再看一眼之前的示意图。这个决定很快就会变成麻烦。
按这种布局,一个客户的行落在一个分片上,而这个客户的订单却散落在许多分片上。现在该怎么连接?
分片查询规划器必须做到:接收任意 SQL 查询,为跨多个 Postgres 节点的执行构造出最优计划。
构建计划
查询规划器的产物同样是一棵树。规划器的职责,是把描述查询请求的 AST 转换成查询计划,而查询计划描述的是集群为算出结果必须完成的全部操作序列。要做到这一点,路由器里的规划器必须同时掌握两样东西:(a) 完整的 Postgres schema,(b) 数据在各分片间的分布规则。它从数据库的_权威分片_(authoritative shard)获取 (a)——该分片被指定为这个数据库 schema 元数据的来源。路由器从这里拿到表名、列和类型。它不会为每次查询都向权威分片请求这些信息,而是维护一份缓存,并定期与权威分片通信以获取更新。它从数据拓扑获取 (b),这份拓扑集中存放在 etcd 中,并在每个路由器节点上缓存。借助这些信息,路由器解析出查询中的列及其类型,这一步称为语义分析。它还会寻找可能绕开分布式规划的捷径,可惜我们的查询没这个运气。
怎么 join?
接下来:把 AST 变成查询计划树。第一步是先转换成一棵粗略的计划树,之后再分几个阶段细化。由原始查询生成的基础计划大致是这样:这很简单。JoinCluster 把我们的两张表以及它们的 inner join 条件归到一起。Inner join 可以重排,这给了规划器不同的方式来应对同一条查询。OutputClauses 描述拿到行之后要做什么。这里我们只需要返回三列。其他查询可能需要排序、分组、聚合、LIMIT,甚至 DISTINCT。规划器还得决定如何 join 这些行,这才是更难的部分。
跨分片 join
规划器用我们的数据拓扑来判断查询的哪些部分可以在分片上执行。前面没有展示,但我们是在一个数据拓扑 JSON 文件里指定这张路由映射的。对于我们的两张表,各自按自身 id 的哈希分片,映射大致如下(数据库名用 postgres):
{
"authoritative_shard_group": "metadata",
"shard_indexes": {
"xxhash_id": {
"type": "xxhash",
"columns": ["id"]
}
},
"shard_groups": [
{
"uid": "metadata",
"key_ranges": [{ "shard_uid": "shard1" }]
},
{
"uid": "main_shards",
"default_shard_index": "xxhash_id",
"key_ranges": [
{ "shard_uid": "shard1", "end": "40" },
{ "shard_uid": "shard2", "start": "40", "end": "80" },
{ "shard_uid": "shard3", "start": "80", "end": "c0" },
{ "shard_uid": "shard4", "start": "c0" }
]
}
],
"databases": {
"postgres": {
"schemas": {
"public": {
"tables": {
"customers": { "shard_group": "main_shards" },
"orders": { "shard_group": "main_shards" }
}
}
}
}
}
}
这个配置告诉路由器:customers 和 orders 两张表都要分散到 main_shards 组的四个分片(shard)上,并且各自按 id 列的哈希值进行分片。存储新行时,路由器会:
- 提取该行的
id。 - 将其输入一个确定性的
xxhash函数。 - 根据给定的范围,用得到的哈希值选择分片。
我们的四个分片各自负责一段不同的哈希值区间。前面说过,这样选择分片键会让事情变得更复杂。第一个后果就在这里:一个客户的订单可能散落在全部四个分片上,而客户行本身只存在于其中一个分片。要连接它们,就必须在路由器里把匹配的行凑到一起。不过我们的 planner 很聪明,它开始更新 JoinCluster 节点,把它变成一个用于检索并连接这些行的计划。对于这个查询,它需要把 customers 和 orders 的工作分开,找出每个分片能应用的过滤条件,并决定路由器如何匹配结果。它首先为每张表创建一个 Route 计划节点。每个 route 内部的工作会变成一条交给 Postgres 执行的 SQL 查询。route 还会指定如何选择目标分片,因此一个 route 可以把查询发给一个分片,也可以发给多个。现在我们知道,要取回 customers 和 orders 两张表所需的行,就得发送 scatter-gather 查询。不过,这个连接有好几种做法:是先取客户还是先取订单?订单应该在分片层过滤,还是在路由器过滤?连接应该用嵌套循环完成,还是用哈希表完成?这就是规划过程的下一步。
跨多个 Postgres 实例进行连接
为了实现 JOIN orders ON orders.customer_id = customers.id,路由器必须把符合条件的订单和对应的客户匹配起来。做法有好几种。嵌套循环连接(nested-loop join)这个概念,熟悉数据库内部实现的人可能不陌生。对其他人来说,思路很简单:从一侧输入取行,然后为每一行在另一侧输入中查找匹配项。路由器会先向每个分片发送查询,收集所有 customers 行。随着行陆续返回,路由器逐个遍历每个客户。对每个客户,它都要向所有分片发一次 scatter-gather 查询,收集该客户的所有订单。哪怕一个客户没有任何符合条件的订单,我们也要花四次分片查询才能确认这一点。router
| id | name | |---|---|---| | ··· | ··· | | ··· | ··· | | ··· | ··· | | ··· | ··· |
SELECT id, nameFROM customers;
| name | total | created_at |
|---|---|---|
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
分片 1customers
| 2 | Alex |
|---|
orders
| id | cust. |
|---|---|
| 3207 | 32 |
| 3209 | 32 |
分片 2customers
| 6 | Jeff |
|---|
orders
| id | cust. |
|---|---|
| 3203 | 32 |
| 3208 | 2 |
分片 3customers
| 32 | Maja |
|---|
orders
| id | cust. |
|---|---|
| 3202 | 2 |
| 32 | 1 |
分片 4customers
| 1 | Leah |
|---|
orders
| id | cust. |
|---|---|
| 3201 | 6 |
| 1 | 6 |
首先,从全部四个分片中收集每一位客户。
这样做效率不高,还有更好的办法。
哈希连接是另一种连接技术,更适合我们的需求。它的思路是:先根据一份输入构建一张查找表,再用它在另一份输入的行陆续到达时找出匹配项。同样地,路由器会先发一个请求,一次性取回全部 customers 行,但这次会在路由器内存中构建一张以客户 ID 为键的哈希映射。然后我们就能做一次大规模的 scatter-gather,请求所有创建于 9 月 9 日及之后的 orders 行。
router
| id | cust. | total | date |
|---|---|---|---|
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| name | total | created_at |
|---|---|---|
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
SELECT id, nameFROM customers;
| key | name |
|---|---|
| ··· | ··· |
| ··· | ··· |
| ··· | ··· |
| ··· | ··· |
Shard 1customers
| 2 | Alex |
|---|
orders
| id | cust. |
|---|---|
| 3207 | 32 |
| 3209 | 32 |
Shard 2customers
| 6 | Jeff |
|---|
orders
| id | cust. |
|---|---|
| 3203 | 32 |
| 3208 | 2 |
Shard 3customers
| 32 | Maja |
|---|
orders
| id | cust. |
|---|---|
| 3202 | 2 |
| 32 | 1 |
Shard 4customers
| 1 | Leah |
|---|
orders
| id | cust. |
|---|---|
| 3201 | 6 |
| 1 | 6 |
一种做法是先取一次 customers,之后每来一条 order,就在 router 里查它对应的 customer。
把 customer 哈希表放在内存里会占用资源,但订单到达时可以直接匹配,不必为每个 customer 再发一次查询。join 规则会比较这两种做法的预估代价。两者都在 router 里完成匹配,用的是 Postgres 在各 shard 上提供的行。
planner 也会考虑把输入顺序反过来。先取 customers 并不是这两种 join 的硬性要求。
这些并不是 planner 仅有的选项。它还支持 merge join、批量 nested-loop join 等等。它的任务是在已有信息下,为查询找出最高效的执行计划。
针对这条查询的选择
planner 先取哪张表、用哪种 join 方式,很大程度上取决于表的大小。假设四个 shard 上一共有 100,000 个 customers 和 1,000,000 条 orders。按它对日期不等式的默认估算,大约三分之一的订单会命中,也就是约 300,000 条。下面的 Picasso 图展示了 planner 在哪些情形下会选择不同的 join 算法和顺序。改动行数,就能看到 planner 的选择如何变化。Orders1101001,00010,000100,0001M10M 1101001,00010,000100,0001M10MCustomers__Nested loop · orders first__Nested loop · customers first__Hash join · build customers__Hash join · build orders__Batched nested loop 我们来看看,为什么在这些行数下它会选择 hash join。用刚才说的 nested loop 方式,就意味着要做 100,000 次 customer 查找(全部都要查),然后对每一个 customer 在 orders 上做一次 scatter-gather。这样大约要取 400,000 行(100,000 个 customers 加 300,000 条 orders),但分散在 router 与 shard 之间成千上万次单独的查询里。hash join 则用四次 shard 查询取 customers,再用四次取命中的 orders。这八次查询同样会把大约 400,000 行传进 router,其中还包括可能永远匹配不上的 customers。这种办法还多了一步:在内存里构建哈希表。planner 还会统计索引列上不同值的数量,用来估算等值过滤或 join 会匹配多少行。根据这些估算,它会权衡 shard 请求数、传输行数、处理行数和驻留内存的行数。在这些行数下,先把 customer 输入取来、留在本地供查找使用,预估代价远低于反复发请求。这里 planner 还有一个选择要做:哈希表可以用任意一个输入来构建。planner 估算 100,000 行 customers 的内存需求小于 300,000 条命中的 orders,所以它选 customers。做完这些选择,就可以把最初的计划变成最终的查询计划。customers 这条路径负责构建哈希表,orders 这条路径提供用于探测的行。两条都是 scatter 路径,也就是说各自把查询发给全部四个 shard。两者都没有把请求收窄到特定 shard-key 值的条件。可如果哈希表放不进内存呢?内存又快又好用,前提是你没用完。一旦哈希表超出内存预算,router 就会溢写到磁盘。它按 join key 把 customers 和 orders 分区,然后逐个 join 对应的分区。
*示意性行数据与一个较小的内存预算,展示的是溢写,而不是文章示例所需的内存。分片位于路由器下方,应用位于其上方,临时磁盘存储位于路由器内部 RAM 的右侧。蓝色的 customer 行从分片向上流动,填充路由器的哈希表。达到内存预算后,它们溢写到三个临时磁盘分区。黄色的 orders 从分片到达,经过 RAM 中的写缓冲区,然后按匹配的连接键刷写到磁盘分区。两个输入都完成分区后,路由器将一个 customer 分区加载到 RAM,用对应的 orders 探测它,返回匹配结果,然后处理下一个分区。磁盘副本保留,而内存中的哈希表被复用。示意性内存计量条显示蓝色 customer 内存、黄色 order 写/读缓冲区以及处理内存。orders 以流式方式通过保留的 customer 哈希表;整个 order 分区不会加载进去。该计量条不是实测字节数,也不是精确的内存预算核算计数器,并且不包含路由器的其他工作。每个内存中的 customer 分区在其连接完成后清除,然后才加载下一个分区。每个黄色 order 都携带 customer_id、total 和 created_at 经过 RAM 和磁盘。连接结果包含 customers.name、orders.total 和 orders.created_at。时间戳以紧凑日期显示:September 9, 2026;示例时间范围为 12:00 到 17:00 UTC。全部六个匹配结果向上传递到应用。暂停之后,循环重复。
无论哈希表能放进内存还是溢写到磁盘,连接仍然在路由器中发生。现在我们有了一个计划。来执行它。
执行
我们仍然没有离开路由器,但很快会推进到集群中的另一个节点。如果这一切都发生在单个实例上,Postgres 可以在本地连接这些表。然而,我们的计划把连接放在路由器中,而路由器并不持有这两张表中的任何一张。_多亏了分片。_计划所需的行分散在四个独立实例上,所以路由器现在必须向分片发送单独的 customer 和 order 查询,然后连接这些查询返回的行。首先,四个分片中的每一个都需要执行 customer 查询:
SELECT customers.id, customers.name
FROM public.customers;
路由构建好客户哈希表后,每个分片都要执行并返回以下查询的结果:
SELECT orders.customer_id, orders.total, orders.created_at
FROM public.orders
WHERE orders.created_at >= $1;
订单查询在全部四个分片上使用同一个 $1 时间戳参数 '2026-09-09 00:00:00+00'。前四个请求已经可以从路由发出了。
连接 Postgres
我们的客户请求要发往全部四个分片。这里只看 2 号分片上发生了什么。路由并不直接连接 Postgres,而是把请求发给 sidecar——一个与 Postgres 并行运行的独立进程。这让我们离 Postgres 之间多了一步,但也把连接管理集中到了每个实例旁边,不同路由发来的请求可以共用同一个 Postgres 连接池。sidecar 负责保证这些连接可以安全复用。如果当初选择直连 Postgres,查询会在一个预先建立的会话里执行。但在我们的架构里,那个会话属于路由。分片请求需要把数据库、认证角色和会话设置随查询一起带上,这样每一部分才能以应用期望的权限和设置运行。路由会准备一串扩展 Postgres 协议消息:Parse、Bind、Execute 和 Sync。这一次,Parse 里装的是我们的客户查询,而不是最初的 join。它把这些消息打包进一个 ExecuteRequest 的 raw 字段,同时带上目标地址和会话信息。请求使用 protobuf,直接把 Postgres 消息以字节形式携带。
接下来需要一个发送请求的通道。路由为每个 sidecar 维护一个由长连接、双向 gRPC 流组成的池。它借用一个空闲流,没有就新开一个,然后发送 ExecuteRequest。sidecar 会在同一条流上返回响应分块。路由
原始 Postgres 消息
-
Parse
SELECT customers. id, customers.name
FROM public .customers
-
Bind
parameters none 3. #### Execute
max_rows = 0 4. #### Sync
targetshard 2 · primary
databaseapplication database
session_userauthenticated role
options.settingssession settings
SidecarShard 2 · primary
我们的客户端请求终于到达了 shard 2,但还没进入 Postgres。
Sidecar 内部
我们已经到了 shard 2 的 sidecar,但还没有运行在 Postgres 里。Postgres 的每连接一进程模型又一次挡在了前面。别忘了,每个 Postgres 连接都需要一个独立的服务器后端进程,即使空闲也占着内存。如果每个 router 都各自维护一套到每个 shard 的连接池,随着规模扩大,这些进程会迅速堆积。sidecar 让这些 router 共享一个热 Postgres 连接池,按需借用。我们的客户端查询可以用上其中一个连接,但这个连接可能还带着上一次使用留下的设置和身份。这时,随 SQL 一起发送的会话信息就派上用场了。Router 借用一条可用的 gRPC stream 连到 shard 2 的 sidecar。我们的请求携带 SQL、Settings 和身份,在 sidecar 界面上一起显示。一个已有的 Postgres 连接从 Sidecar 的独立池中升起。它此前的设置和身份让位给请求中的状态。先对齐设置,再对齐身份。两者都就绪之后,SQL 才被送往 Postgres,也就是 shard 2 的 primary。已经匹配的状态可以直接复用,无需改动。Stream 和 Postgres 连接并非永久绑定。每个池都有自己的并行通道。池的占用情况和时序仅为示意。升起和对接表现的是取出连接和准备会话的过程,不是新建物理连接,也不是转移连接池本身。循环重复这段说明,不表现响应或释放。
SELECT customers.id, customers.name FROM public.customers
SettingsSession settingsIdentityAuthenticated role
SettingsPrevious useSession settingsIdentityPrevious useAuthenticated role
连接池首先寻找设置与我们的请求匹配的可用连接。找到了,就可以原样使用,不必改动。否则,sidecar 会把此前的设置重置为默认值,再应用我们的设置。
sidecar 还会检查会话的角色并确认其正确,防止角色污染。无论上一个借用该连接的是谁,我们的查询都必须以我们的权限运行。
sidecar 把 SQL 转发出去,我们的客户查询终于抵达 Postgres。
终于到 Postgres 了
客户查询现在抵达 Postgres。
SELECT customers.id, customers.name
FROM public.customers;
我们在 router 层面花了那么多时间做规划,现在 Postgres 也得自己出一份计划。这是一个真正的 Postgres 实例,所以它用自己的 planner 和统计信息来决定怎么读 customers 表,然后针对 shard 2 上的 customers 执行这份计划。行一旦就绪,Postgres 就通过那条借来的连接把它们发回 sidecar。我们的哈希表有了雏形!
回到 router
sidecar 把这些行通过发送请求时用的那条 gRPC 流中继给 router。行是分批到达的,所以 shard 还在返回数据时,router 就能开始构建 customer 哈希表。四个 customer 请求全部完成、哈希表建好之后,哈希连接启动 orders Route 计划节点。又有四个请求沿着我们刚刚走过的传输和连接池路径出发。这一次,Postgres 在每个 shard 上应用 9 月 9 日这个过滤条件,把符合条件的订单发回来。返回的行产生了匹配!router 用每个订单的 customer_id 在哈希表里找到对应的 customer,然后构造出一行结果,包含客户姓名、订单总额和创建时间戳。这里有个相当巧妙的做法,能帮我们省下一些计算周期。我们向 Postgres 请求输出字段时,用的是客户端为原始查询选定的同样的文本或二进制格式。这些值已经编码好了。router 直接在内存中从匹配的行里复制已编码的姓名、总额和时间戳来构造结果,省掉了先解码成 Go 值、再重新编码的功夫。customersid 32 · Maja440000001A0002000000080000000000000020000000044D616A61
orderscustomer_id 32440000003500030000000800000000000000200000000533302E303000000016323032362D30392D30392031323A30303A30302B3030
DataRow → driver_7 新头部字节_ + 43 已复制44000000310003000000044D616A610000000533302E303000000016323032362D30392D30392031323A30303A30302B3030
000000044D616A610000000533302E303000000016323032362D30392D30392031323A30303A30302B3030
customer 行包含 ID 32 和 Maja。order 行包含 customer ID 32、total 30.00 以及时间戳 2026-09-09 12:00:00+00。两个二进制 ID 被读取用于匹配,不会进入输出。编码后的 name、total 和时间戳的副本,连同各自的字段长度前缀,被移动到新消息中的确切位置。这些输出值使用文本格式。新头部是 44 00 00 00 31 00 03。其长度为 49,不包括 D 标签,列数为 3。完整的 DataRow 为 50 字节。源字节保持不变。
我们也不必等待每个 order 请求都完成。应用的 Postgres 客户端可以在分片仍在返回 order 时就开始接收结果。
当四个分片全部返回完 order 后,我们这条查询的旅程就结束了。应用拿到了结果,连接也准备好迎接下一条查询。
求值引擎
有些情况下,我们不能使用这种直接复制字节的捷径。在我们一直跟踪的这条查询中,输出里的每个值都已经存在于分片返回的行中。路由器可以把它们编码后的字节直接复制到结果里。但如果我们改为请求每个 customer 的平均 order total,会发生什么?
SELECT customers.name, AVG(orders.total) AS average_total
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1
GROUP BY customers.id, customers.name;
这需要路由器中某个组件做额外的工作,你可能认得它,它是 Postgres 的一个内部组件:Eval Engine。记住,单个 customer 的 order 分散在全部四个分片上。单个分片无法独自算出完整的平均值。我们也不能对分片平均值再求平均。一个分片可能有某个 customer 的两条 order,而另一个分片有二十条。于是路由器把平均值改写为求和与计数,由每个分片在本地计算。结果返回后,路由器按 customer 合并它们。路由器的求值引擎再用合并后的总和除以合并后的计数,得出该 customer 的平均值。
路由器把 AVG(orders.total) 改写成 SUM(total) 和 COUNT(total),再把这项工作发给四个分片。每个分片为一个示例客户组计算总和与计数。对于客户 32,也就是 Maja,分片 1 返回总和 $100、计数 2;分片 2 返回 $90 和 3;分片 3 返回 $80 和 1;分片 4 返回 $50 和 2。计数只包含非空的订单总额。路由器把它们合并为 $320 和 8。它的求值引擎用 $320 除以 8,得出平均值 $40,连同 Maja 的名字一起返回给应用。若对四个分片的平均值再取平均,会错误地得到 $46.25。这些面板区分的是概念步骤,不是物理计划算子。这个特写省略了客户 join 和其他分组;每个订单只匹配一个客户。数值、到达顺序和时序均为示意。开启“减弱动态效果”时,计算完成后的结果会保持可见。
不过,对于我们最初的查询,求值引擎并没有多少工作要做。
更好的拓扑
我们已经看到,路由器为了跨分片 join 和计算结果要做多少工作。我们选择的分片键让这些工作变得比必要的更繁重。按主键分片本身并不是问题。我们可以继续用 customers.id 来分布客户。复杂性来自按 orders.id 分布订单,而这与我们需要的客户是彼此独立的,join 时又必须把它们对应起来。如果我们按 customer_id 对 orders 分片,使用与 customers 相同的哈希和分片映射,那么每个客户就会和他们的订单待在一起。Postgres 随后就能在本地完成 join。我们的查询仍然会到达全部四个分片,因为它要的是所有客户的近期订单。但现在我们把完整的 join 发给每个 Postgres 实例,由它在本地 join 自己的行并返回结果。四次分片请求而不是八次,路由器也不需要再构建客户哈希表。客户仍然按 customers.id 分片。这些就是文章前面用过的同一批示例行。订单现在使用 customer_id,并采用与 customers.id 相同的分片映射。应用把同一个查询发给路由器。路由器把一条 join 后的查询发给四个分片中的每一个。Postgres 在本地 join 匹配的客户和符合条件的订单。分片 4 也会被查询,但它没有符合条件的订单,因此不返回任何 join 后的行。其他分片总共返回六行 join 结果,路由器在它们到达时随即转发。两种布局从路由器返回给应用的都是同样的六个结果,然后在完成后的结果上暂停,再重复播放。移动这些示例行是在对比两种布局,不是真实的重新分片过程。动画速度不是基准测试。
- 分片 1 保存客户 2,以及为客户 2 选出的订单 3202、3208。
- 分片 2 保存客户 6,以及为客户 6 选出的订单 3201。
- 分片 3 保存客户 32,以及为客户 32 选出的订单 3207、3203、3209。
- 分片 4 保存客户 1,选出的订单为空。
我们仍然需要刚才讲过的认证、传输和连接管理。但把分别取回的客户和订单在 router 里拼到一起这件事,没有了。最快的 router join,就是根本不用跑的那个。这话其实几千字之前就该说了。
看,很简单吧?
欢迎来到 Neki。