本文面向具备一定分布式系统和算法基础的中高级工程师,旨在深度剖析“多对多”批量撮合(Batch Matching)的核心机制。我们将从金融市场(如股票集合竞价)和数字资产交易的真实需求出发,回归到微观经济学的供需平衡原理,推导并实现最大成交量算法。文章将覆盖从基础数据结构到CPU缓存友好的算法优化,再到基于消息队列的高可用、可扩展的分布式撮合架构,最终提供一套从简单到复杂的架构演进路径,帮助技术负责人应对不同阶段的业务挑战。
现象与问题背景
在高性能交易系统中,撮合引擎是心脏。我们常见的“连续撮合”(Continuous Matching),即订单簿(Order Book)模型,遵循“价格优先、时间优先”原则,逐笔处理新进入的订单。这种模式适用于流动性好的市场,能提供即时成交的体验。然而,在某些特定场景下,连续撮合会暴露出明显的问题:
- 开盘/收盘定价: 股票市场每日开盘和收盘时,需要在瞬间处理大量隔夜订单或收盘前涌入的订单,以形成一个公允的开盘价/收盘价。连续撮合无法实现这一目标,因为它会按毫秒级的时间顺序处理,导致价格剧烈波动且不公平。
- 新资产发行(ICO/IEO/IPO): 在一个新资产首次发行的抢购阶段,瞬时流量极大,订单堆积如山。如果采用连续撮合,最先到达的订单(通常是低延迟网络的“科学家”)会占据绝对优势,引发“抢跑”(Front-running)和公平性问题。
- 低流动性市场: 对于一些交易不活跃的资产,订单簿稀疏,买卖价差(Spread)巨大。连续撮合可能导致一笔市价单就能造成巨大的价格滑点,市场容易被操控。
为了解决这些问题,批量撮合,也称为集合竞价(Call Auction),应运而生。其核心思想是在一个指定的时间窗口内(例如开盘前的5分钟),只接收订单而不进行撮合。时间窗口关闭后,系统将所有累积的买单和卖单视为一个整体,通过特定算法计算出一个唯一的“均衡价格”(Equilibrium Price),使得在此价格下能够成交的数量最大化。所有符合条件的订单都将以这个统一价格成交。这不仅解决了公平性问题,也有效地为市场“发现”了价格。
核心挑战在于:如何设计一个高效、确定且可扩展的算法与系统,在海量订单中,快速找到那个能实现最大成交量的价格平衡点?
关键原理拆解
在深入工程实现之前,我们必须回归到计算机科学与经济学的基础原理。批量撮合的本质,是一个基于离散订单数据的优化问题,其理论根基是微观经济学中的供需曲线模型。
想象一下,我们将所有买单按价格从高到低排序,所有卖单按价格从低到高排序。对于任何一个可能的价格 P,我们可以计算出两个关键指标:
- 累计买方需求量 (Cumulative Buy Quantity): 所有出价
>= P的买单数量之和。因为任何愿意出更高价的买家,在P这个价格上更愿意购买。 - 累计卖方供给量 (Cumulative Sell Quantity): 所有出价
<= P的卖单数量之和。因为任何愿意以更低价出售的卖家,在P这个价格上更愿意出售。
在价格 P 上,理论上可以撮合的成交量为 V(P) = min(累计买方需求量, 累计卖方供给量)。我们的目标就是找到一个价格 P*,使得 V(P*) 最大。这个 P* 就是我们寻找的均衡价格,V(P*) 就是最大成交量。
从算法角度看,一个朴素的想法是遍历所有可能的价格点,计算每个点的 V(P),然后找到最大值。但这存在两个问题:价格是连续的还是离散的?遍历范围是多大?
关键的洞察在于:累计需求量和供给量的阶跃变化只会发生在订单簿中真实存在的那些价格点上。 在两个相邻订单价格之间,累计量是恒定的。因此,我们只需要在所有订单出现过的价格点(以及它们的邻近价格)上进行计算即可。这瞬间将一个看似无限的问题,转化为了一个有限离散点的寻优问题。
基于此,我们可以设计出核心算法的框架:
- 数据聚合: 将时间窗口内的所有买单和卖单收集起来。
- 价格水平构建: 对所有买单和卖单,按价格进行分组,累加每个价格水平(Price Level)上的总数量。例如,买单 {P:100, Q:50}, {P:100, Q:30} 会被合并为 {P:100, Q:80}。
- 累计曲线生成:
- 对买单价格水平,从高到低遍历,生成每个价格点的累计买量曲线。
- 对卖单价格水平,从低到高遍历,生成每个价格点的累计卖量曲线。
- 均衡点查找: 遍历所有出现过的价格点,计算每个点的
min(累计买量, 累计卖量),记录下使该值最大的价格P*和成交量V*。
这里还有一个重要的约束:成交价格的确定性规则。如果多个价格点都能达成同样的最大成交量,我们必须有一个明确的规则来选择最终的均衡价格。常见的规则有:
- 高价优先: 选择成交价格最高的那个。
- 低价优先: 选择成交价格最低的那个。
- 参考价优先: 选择最接近上一个收盘价或某个基准价的价格。
- 买卖压力平衡: 选择使得未成交的买卖双方订单不平衡量最小的价格。
交易所通常会明确定义这些规则。例如,选择的价格 P* 不仅要满足最大成交量,还要使得所有价格高于 P* 的买单和低于 P* 的卖单都能完全成交。
系统架构总览
一个生产级的批量撮合系统,远不止于算法本身。它是一个涉及高可用、数据一致性和低延迟通信的复杂分布式系统。我们可以用文字勾勒出其核心架构:
整个系统围绕一个核心的、不可变的消息日志(通常由 Kafka 或类似组件扮演)构建,分为几个关键服务:
- 1. 订单网关 (Order Gateway):
- 作为系统入口,负责接收来自客户端的订单请求。
- 进行基础的协议解析、参数校验、用户认证和风控检查(如账户余额)。
- 通过验证的订单被序列化成标准格式,并被原子地、有序地发布到消息队列(如 Kafka 的一个特定 Topic Partition)中。发布成功后才向客户端确认。这确保了订单的持久化和顺序性。
- 2. 消息队列/定序器 (Sequencer - Kafka):
- 所有交易指令的唯一真相来源(Single Source of Truth)。
- 为整个系统提供了削峰填谷、异步解耦和故障恢复的能力。一个交易对(Symbol)通常对应一个 Partition,保证了单个交易对的订单处理是严格有序的。
- 3. 批量撮合引擎 (Batch Matching Engine):
- 系统的核心计算单元。它订阅特定交易对的 Kafka Topic。
- 在批量撮合窗口期,它只消费消息并构建内存中的订单集合。
- 当触发信号(如预设的定时器或一个特殊的控制指令消息)到达时,引擎停止接收新订单,对内存中的订单集合执行上文所述的最大成交量算法。
- 计算出均衡价格和所有成交明细(Trades/Fills)。
- 4. 成交结果发布器 (Trade Publisher):
- 撮合引擎计算出的成交结果(是一个原子性的集合),被打包成消息,发布到另一个“成交结果”的 Kafka Topic 中。
- 5. 下游消费者 (Downstream Consumers):
- 清算服务 (Clearing Service): 订阅成交结果 Topic,进行账户资金和持仓的变更。
- 行情服务 (Market Data Service): 订阅成交结果,生成K线、最新成交价等市场行情数据,并推送给客户端。
- 持久化服务 (Persistence Service): 将订单状态和成交记录异步地写入数据库(如 MySQL/PostgreSQL)以供历史查询和审计。
这种架构的优点是显而易见的:通过 Kafka 将核心的撮合计算与外围的IO、持久化、行情推送等任务完全解耦,使得撮合引擎可以专注于其核心使命——计算,从而达到极高的性能和吞吐量。同时,基于日志的架构天然支持故障恢复和水平扩展。
核心模块设计与实现
我们聚焦于撮合引擎的内部实现。这部分,我们得像个极客一样思考内存、CPU缓存和代码效率。
模块一:订单聚合与价格水平构建
当订单从 Kafka 流入时,我们需要一个高效的数据结构来存储它们。一个简单的 List 是不够的,因为它在后续计算中需要反复遍历。更优的做法是直接构建价格水平视图。
我们可以使用两个 Map(或哈希表)来分别存储买方和卖方的价格水平:`buyLevels map[int64]int64` 和 `sellLevels map[int64]int64`。其中 key 是价格(为避免浮点数精度问题,通常会将价格乘以一个固定的倍数转换为整数),value 是该价格上的累计订单数量。
// Order represents a simplified order structure from Kafka
type Order struct {
Side string // "BUY" or "SELL"
Price int64 // Price scaled to integer
Quantity int64
}
// buildPriceLevels consumes orders and aggregates them into price levels.
func buildPriceLevels(orders []Order) (buyLevels, sellLevels map[int64]int64) {
buyLevels = make(map[int64]int64)
sellLevels = make(map[int64]int64)
for _, order := range orders {
if order.Side == "BUY" {
buyLevels[order.Price] += order.Quantity
} else {
sellLevels[order.Price] += order.Quantity
}
}
return
}
这个步骤的时间复杂度是 O(N),其中 N 是订单数量。这是非常高效的预处理。
模块二:最大成交量算法实现
预处理完成后,我们得到了聚合后的价格水平。接下来的步骤是算法的核心。
第一步:提取并排序所有唯一价格点。
我们需要一个包含所有买卖价格的、无重复且有序的列表。这是后续计算的基础。
第二步:构建累计数量曲线。
这里有一个工程上的技巧。我们不需要为每个价格点都存储一个完整的累计值。我们可以先将价格水平转换为切片并排序,然后在计算过程中动态累加。买单价格从高到低排序,卖单价格从低到高排序。
第三步:寻找均衡点。
我们将买卖双方的有序价格水平进行一次“合并扫描”,即可找到均衡点。这比遍历所有价格点更高效。
import "sort"
type PriceLevel struct {
Price int64
Quantity int64
}
// FindMatchPrice calculates the equilibrium price and max volume.
func FindMatchPrice(buyLevels, sellLevels map[int64]int64) (matchPrice, maxVolume int64) {
// 1. Convert maps to sorted slices
sortedBuys := make([]PriceLevel, 0, len(buyLevels))
for p, q := range buyLevels {
sortedBuys = append(sortedBuys, PriceLevel{p, q})
}
sort.Slice(sortedBuys, func(i, j int) bool {
return sortedBuys[i].Price > sortedBuys[j].Price // DESC for buys
})
sortedSells := make([]PriceLevel, 0, len(sellLevels))
for p, q := range sellLevels {
sortedSells = append(sortedSells, PriceLevel{p, q})
}
sort.Slice(sortedSells, func(i, j int) bool {
return sortedSells[i].Price < sortedSells[j].Price // ASC for sells
})
// 2. Create cumulative curves
cumulativeBuys := make([]PriceLevel, len(sortedBuys))
var totalBuyQty int64 = 0
for i, level := range sortedBuys {
totalBuyQty += level.Quantity
cumulativeBuys[i] = PriceLevel{level.Price, totalBuyQty}
}
cumulativeSells := make([]PriceLevel, len(sortedSells))
var totalSellQty int64 = 0
for i, level := range sortedSells {
totalSellQty += level.Quantity
cumulativeSells[i] = PriceLevel{level.Price, totalSellQty}
}
// 3. Find the equilibrium point by iterating through potential prices
// We use two pointers, one for the buy curve and one for the sell curve.
buyIdx, sellIdx := 0, 0
maxVolume = 0
matchPrice = 0
for buyIdx < len(cumulativeBuys) && sellIdx < len(cumulativeSells) {
buyLevel := cumulativeBuys[buyIdx]
sellLevel := cumulativeSells[sellIdx]
// A match is only possible if the highest bid is >= the lowest ask
if buyLevel.Price < sellLevel.Price {
break
}
var currentPrice int64
var currentVolume int64
// We check the interval defined by the current buy and sell prices
if buyLevel.Price >= sellLevel.Price {
// Potential match exists. The volume is min of cumulative quantities.
vol := min(buyLevel.Quantity, sellLevel.Quantity)
// This volume is achievable at any price between sellLevel.Price and buyLevel.Price
// Let's check the buy price and sell price as candidates
if vol > maxVolume {
maxVolume = vol
// Tie-breaking rule: choose higher price. Can be configured.
matchPrice = buyLevel.Price
} else if vol == maxVolume {
// Apply tie-breaking, e.g., choose price that is more "balanced" or just higher/lower.
if buyLevel.Price > matchPrice {
matchPrice = buyLevel.Price
}
}
if vol > maxVolume {
maxVolume = vol
matchPrice = sellLevel.Price
} else if vol == maxVolume {
if sellLevel.Price > matchPrice {
matchPrice = sellLevel.Price
}
}
}
// Advance the pointer that corresponds to the "smaller" price range
// to explore the next price interval.
if buyIdx+1 < len(cumulativeBuys) && sellIdx+1 < len(cumulativeSells) {
if cumulativeBuys[buyIdx+1].Price > cumulativeSells[sellIdx+1].Price {
buyIdx++
} else {
sellIdx++
}
} else if buyIdx+1 < len(cumulativeBuys) {
buyIdx++
} else {
sellIdx++
}
}
return matchPrice, maxVolume
}
func min(a, b int64) int64 {
if a < b {
return a
}
return b
}
复杂度分析: 假设有 N 个买单和 M 个卖单,形成了 P_buy 和 P_sell 个独特的价格水平。
- 订单聚合:O(N+M)
- 价格水平排序:O(P_buy * log(P_buy) + P_sell * log(P_sell))
- 累计曲线生成和均衡点查找:O(P_buy + P_sell)
在大多数市场中,价格水平的数量 P 远小于订单总数 N,所以性能瓶颈在于排序。此实现已经相当高效。
性能优化与高可用设计
性能优化
- CPU Cache 优化: 上述代码中将 Map 转换为 Slice 再进行处理,这是一种对 CPU 缓存更友好的做法。连续的内存访问(遍历 Slice)比离散的内存访问(遍历 Map)要快得多。在性能极致的场景中,可以考虑使用基数排序(Radix Sort)代替比较排序,如果价格是分布在有限范围内的整数。
- 内存管理: 对于 Go 语言,撮合引擎这种需要长时间运行且对延迟敏感的服务,要特别注意内存分配。可以为订单对象、价格水平对象等建立对象池(sync.Pool),避免在高并发时给 GC 带来巨大压力,从而减少 STW(Stop-The-World)暂停时间。
- 并发与并行: 撮合算法的核心计算部分(寻找均衡点)为了保证确定性和简单性,通常是单线程执行的。但是,准备数据的阶段,如从 Kafka 消费、反序列化、聚合到价格水平 Map,完全可以并行化,为每个 CPU核心分配一个协程来处理一部分订单。
高可用设计
撮合引擎是单点,一旦宕机,整个交易市场就会停摆。高可用是必须的。
- 主备(Active-Passive)模式: 这是最经典可靠的模式。部署两台完全相同的撮合引擎实例,一个为主(Active),一个为备(Passive)。它们都订阅同一个 Kafka Topic Partition。
- 主节点正常处理订单,并将自己处理到的 Kafka Offset 定期汇报给一个协调服务(如 ZooKeeper 或 etcd)。
- 备用节点也在消费数据,构建内存状态,但它不向外发布任何成交结果。它只是默默地“跟随”主节点。
- 当主节点心跳丢失时,协调服务会触发主备切换。备用节点确认自己已经追赶到主节点最后汇报的 Offset,然后切换为 Active 状态,开始对外发布成交结果。
- 状态快照与恢复: 一个长时间运行的撮合引擎,内存中的订单状态会非常庞大。如果重启,从 Kafka 的最开始位置回放所有订单会非常耗时。因此,引擎需要定期将内存中的订单簿状态(即 `buyLevels` 和 `sellLevels`)进行快照,并与当时的 Kafka Offset 一同持久化到磁盘或分布式存储。当引擎重启时,它可以直接加载最新的快照,然后从快照对应的 Offset 开始消费 Kafka 消息,大大缩短恢复时间(RTO)。
架构演进与落地路径
一个复杂的系统不是一蹴而就的。根据业务发展阶段,可以规划清晰的演进路径。
第一阶段:单体 MVP (Minimum Viable Product)
- 架构: 一个单体应用,包含了订单接收、内存撮合、结果持久化到数据库(如 PostgreSQL)的所有逻辑。
- 部署: 单机部署。
- 适用场景: 业务初期,交易对少,并发量低。快速验证商业模式是首要目标。
- 缺点: 无高可用,性能瓶颈明显,所有模块紧耦合。
第二阶段:服务化与高可用
- 架构: 引入 Kafka,将系统拆分为订单网关、撮合引擎、清算服务等微服务。实现上文描述的基于消息队列的解耦架构。
- 部署: 为撮合引擎部署 Active-Passive 主备集群。其他服务可以进行常规的无状态水平扩展。
- 适用场景: 业务进入成长期,对系统的稳定性和可用性提出高要求。并发量显著增加。这是大多数生产系统的标准形态。
- 优点: 高可用,各模块可独立演进和扩缩容,系统健壮性大幅提升。
第三阶段:多交易对水平扩展(Sharding)
- 架构: 当交易对数量达到成百上千,单个撮合引擎集群(即使是主备)也可能成为瓶颈。此时需要按交易对进行分片(Sharding)。
- 实现: 在 Kafka 中,创建多个 Topic,或一个 Topic 的多个 Partition。例如,将哈希(交易对名称) % N 的结果,决定该交易对的订单路由到第 N 个 Partition。部署 N 组独立的撮合引擎主备集群,每组集群只负责处理一部分 Partition 的数据。
- 部署: 订单网关需要增加路由逻辑。整个撮合层可以近似线性地水平扩展。
- 适用场景: 大型交易所,需要支持海量交易对和极高的并发订单处理能力。
通过这样的演进路径,团队可以在不同阶段聚焦于最核心的矛盾,用合适的架构成本支撑业务发展,避免过度设计,也为未来的大规模扩展预留了清晰的路径。
延伸阅读与相关资源
-
想系统性规划股票、期货、外汇或数字币等多资产的交易系统建设,可以参考我们的
交易系统整体解决方案。 -
如果你正在评估撮合引擎、风控系统、清结算、账户体系等模块的落地方式,可以浏览
产品与服务
中关于交易系统搭建与定制开发的介绍。 -
需要针对现有架构做评估、重构或从零规划,可以通过
联系我们
和架构顾问沟通细节,获取定制化的技术方案建议。