分片 Postgres 查询的生命周期

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

中文
复制
题图:深蓝底的架构示意图,PRIMARY / REPLICA / SHARD 1-8 与 PARSER、PLANNER、ROUTER、EXECUTION 等方块用连线串起来,右侧写着 The Lifecycle of a sharded Postgres query 与 PlanetScale 的 BLOG // ENGINEERING 标记

如果你用过 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 的分片键。这些主键是自然的起点,虽然有点天真,后面很快就会看到。每次要存一行 customersorders,系统会取出 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 其实有两种发送查询的方式:SimpleExtended。你可能也不知道,不同的 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" }
 }
 }
 }
 }
 }
}

这个配置告诉路由器:customersorders 两张表都要分散到 main_shards 组的四个分片(shard)上,并且各自按 id 列的哈希值进行分片。存储新行时,路由器会:

  • 提取该行的 id
  • 将其输入一个确定性的 xxhash 函数。
  • 根据给定的范围,用得到的哈希值选择分片。

我们的四个分片各自负责一段不同的哈希值区间。前面说过,这样选择分片键会让事情变得更复杂。第一个后果就在这里:一个客户的订单可能散落在全部四个分片上,而客户行本身只存在于其中一个分片。要连接它们,就必须在路由器里把匹配的行凑到一起。不过我们的 planner 很聪明,它开始更新 JoinCluster 节点,把它变成一个用于检索并连接这些行的计划。对于这个查询,它需要把 customersorders 的工作分开,找出每个分片能应用的过滤条件,并决定路由器如何匹配结果。它首先为每张表创建一个 Route 计划节点。每个 route 内部的工作会变成一条交给 Postgres 执行的 SQL 查询。route 还会指定如何选择目标分片,因此一个 route 可以把查询发给一个分片,也可以发给多个。现在我们知道,要取回 customersorders 两张表所需的行,就得发送 scatter-gather 查询。不过,这个连接有好几种做法:是先取客户还是先取订单?订单应该在分片层过滤,还是在路由器过滤?连接应该用嵌套循环完成,还是用哈希表完成?这就是规划过程的下一步。

跨多个 Postgres 实例进行连接

为了实现 JOIN orders ON orders.customer_id = customers.id,路由器必须把符合条件的订单和对应的客户匹配起来。做法有好几种。嵌套循环连接(nested-loop join)这个概念,熟悉数据库内部实现的人可能不陌生。对其他人来说,思路很简单:从一侧输入取行,然后为每一行在另一侧输入中查找匹配项。路由器会先向每个分片发送查询,收集所有 customers 行。随着行陆续返回,路由器逐个遍历每个客户。对每个客户,它都要向所有分片发一次 scatter-gather 查询,收集该客户的所有订单。哪怕一个客户没有任何符合条件的订单,我们也要花四次分片查询才能确认这一点。router

| id | name | |---|---|---| | ··· | ··· | | ··· | ··· | | ··· | ··· | | ··· | ··· |

SELECT id, nameFROM customers;

nametotalcreated_at
·········
·········
·········
·········
·········
·········

分片 1customers

2Alex

orders

idcust.
320732
320932

分片 2customers

6Jeff

orders

idcust.
320332
32082

分片 3customers

32Maja

orders

idcust.
32022
321

分片 4customers

1Leah

orders

idcust.
32016
16

首先,从全部四个分片中收集每一位客户。

这样做效率不高,还有更好的办法。

哈希连接是另一种连接技术,更适合我们的需求。它的思路是:先根据一份输入构建一张查找表,再用它在另一份输入的行陆续到达时找出匹配项。同样地,路由器会先发一个请求,一次性取回全部 customers 行,但这次会在路由器内存中构建一张以客户 ID 为键的哈希映射。然后我们就能做一次大规模的 scatter-gather,请求所有创建于 9 月 9 日及之后的 orders 行。

router

idcust.totaldate
············
············
············
············
············
············
nametotalcreated_at
·········
·········
·········
·········
·········
·········

SELECT id, nameFROM customers;

keyname
······
······
······
······

Shard 1customers

2Alex

orders

idcust.
320732
320932

Shard 2customers

6Jeff

orders

idcust.
320332
32082

Shard 3customers

32Maja

orders

idcust.
32022
321

Shard 4customers

1Leah

orders

idcust.
32016
16

一种做法是先取一次 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 协议消息:ParseBindExecuteSync。这一次,Parse 里装的是我们的客户查询,而不是最初的 join。它把这些消息打包进一个 ExecuteRequestraw 字段,同时带上目标地址和会话信息。请求使用 protobuf,直接把 Postgres 消息以字节形式携带。

接下来需要一个发送请求的通道。路由为每个 sidecar 维护一个由长连接、双向 gRPC 流组成的池。它借用一个空闲流,没有就新开一个,然后发送 ExecuteRequest。sidecar 会在同一条流上返回响应分块。路由

原始 Postgres 消息

  1. Parse

SELECT customers. id, customers.name
FROM public .customers
  1. 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_idorders 分片,使用与 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。

来源: PlanetScale← 返回首页