一、账单系统面临的挑战
蚂蚁ANTOM(支付宝国际版)的账单子系统需要对全球数千万商户的每一笔交易生成账单记录。随着业务快速扩张,系统面临着几个严峻的挑战:
- 数据量爆炸:账单明细表月增量超过 2 亿条,单表数据量在 6 个月内突破 5 亿大关。
- 查询性能劣化:商户端查询自己的账单列表,原本 200ms 的响应时间逐步恶化到 5 秒以上,高峰期甚至触发慢查询超时。
- 对账时效性不足:商户对账原本采用 T+1 日的离线批处理模式,遇到大促场景时对账延迟长达 T+3,商户投诉量激增。
- 数据倾斜:头部大商户的交易量是中小商户的万倍以上,按商户ID哈希分片后,少数分片承载了绝大部分数据量。
经过多轮架构评审,我们确定了一个分阶段的改造方案:先解决存储瓶颈(分库分表),再解决查询瓶颈(ES重构),最后提升时效性(离线改在线)。
二、分库分表规则设计
在分片策略上,我们没有采用简单的哈希取模,而是设计了一套 复合分片规则 来应对数据倾斜问题:
// 分片路由逻辑伪代码
public int route(String merchantId, String billMonth) {
long merchantIdHash = hash(merchantId);
// 判断是否为头部大商户(交易量 TOP 0.1%)
if (bigMerchantDetector.isBigMerchant(merchantId)) {
// 大商户按月份水平分片,每个月份对应独立的物理分片
int monthShard = Integer.parseInt(billMonth.substring(4)) % 12;
return BIG_MERCHANT_SHARD_OFFSET + monthShard;
} else {
// 中小商户按商户ID哈希取模,路由到通用分片池
return (int)(merchantIdHash % GENERAL_SHARD_COUNT);
}
}
这套分片规则的核心思想是 "大小商户隔离":
- 通用分片池(128 个分片):承载 99.9% 的中小商户,数据分布均匀。
- 大商户专属分片(12 个月份分片 x 8 个大商户组 = 96 个分片):每个大商户的交易数据按月份隔离,避免单个大商户的数据量压垮一个物理分片。
分库分表中间件我们选择了蚂蚁自研的 DRDS(现在叫 PolarDB-X),它原生支持自定义分片规则和跨分片聚合查询。实际落地过程中,我们还解决了一个棘手问题——跨分片的商户汇总查询:
-- 商户需要按天汇总自己的交易额,但由于数据分布在多个分片上
-- DRDS 会将查询下推到各个分片并行执行,然后在中间层做归并
SELECT bill_date, SUM(amount) as total_amount, COUNT(*) as bill_count
FROM billing_detail
WHERE merchant_id = 'M_BIG_001'
AND bill_month = '202605'
GROUP BY bill_date
ORDER BY bill_date;
三、查询重构:从 MySQL 到 Elasticsearch
分库分表解决了写入和存储的问题,但 多条件模糊查询 + 排序 + 分页 的场景在分片数据库上表现仍然不理想。商户端的账单查询页面支持按交易时间、金额范围、交易状态、交易对手等多个维度筛选,这种复杂查询在 DRDS 上需要跨多个分片执行,P99 延迟高达 5 秒。
我们引入 Elasticsearch 作为查询层,架构调整为 "MySQL 写 + ES 读" 的双写模式:
// 账单写入时同步写入 ES(通过 Canal 监听 binlog 异步同步)
@Document(indexName = "billing_detail")
public class BillingDetailES {
@Id
private String billId;
@Field(type = FieldType.Keyword)
private String merchantId;
@Field(type = FieldType.Date, format = DateFormat.date_hour_minute_second_millis)
private Date tradeTime;
@Field(type = FieldType.Double)
private BigDecimal amount;
@Field(type = FieldType.Keyword)
private String tradeStatus;
@Field(type = FieldType.Keyword)
private String counterpartyName;
// 自定义分词器支持中文模糊搜索
@Field(type = FieldType.Text, analyzer = "ik_max_word", searchAnalyzer = "ik_smart")
private String description;
}
数据同步方面,我们没有采用应用层双写(容易导致不一致),而是基于 Canal + MQ + 消费者写入 ES 的异步同步方案。为了保证最终一致性,还设计了一个定时对账任务,每小时比对 MySQL 和 ES 的记录数差异:
- MySQL 批量 COUNT:使用 DRDS 的分片并行能力,128 个分片并行统计。
- ES 批量 COUNT:利用 ES 的路由机制,按 merchantId 路由到对应 shard。
- 差异修复:对不一致的商户ID记录到修复队列,由补偿程序重推。
重构后,商户端账单查询的 P99 延迟从 5 秒降至 100 毫秒以内,用户体验有了质的飞跃。
四、离线改在线:对账时效从 T+1D 到 T+1H
原有的对账流程是每天凌晨 2 点通过 Spark 离线任务对前一天的账单进行汇总比对,整个流程需要 4-6 小时,商户最快也要到第二天中午才能看到对账结果。
改造方案采用了 "微批 + 流式"混合架构:
// 微批对账调度:每小时触发一次
@Scheduled(cron = "0 5 * * * ?")
public void hourlyReconciliation() {
// 1. 从 MySQL 按小时窗口拉取增量账单数据
List<BillingRecord> bills = billRepository.findByHourWindow(
LocalDateTime.now().minusHours(1), LocalDateTime.now());
// 2. 从上游交易系统拉取对应窗口的交易流水
List<TradeRecord> trades = tradeClient.queryByTimeWindow(
LocalDateTime.now().minusHours(1), LocalDateTime.now());
// 3. 基于商户ID做分组对账
Map<String, ReconciliationResult> results = reconciliationEngine
.compare(bills, trades);
// 4. 差异记录写入差异表,触发告警
results.values().stream()
.filter(r -> !r.isMatched())
.forEach(r -> alertService.sendAlert(r));
}
同时,对于大促等极端流量场景,我们还部署了一套 多级容灾降级策略:
- Level 1(正常):微批对账每小时执行,ES 查询 + MySQL 验证。
- Level 2(降级):ES 不可用时,查询降级到 MySQL DRDS,延迟略有上升但功能可用。
- Level 3(兜底):MySQL 某些分片不可用时,读取本地缓存 + Redis 缓存中的近期账单数据,标记为"部分数据"返回商户。
改造完成后,对账时效性从 T+1D(隔日)提升到 T+1H(隔小时),商户投诉率下降 73%,且系统在当年的双 11 大促期间平稳运行,峰值 TPS 达到 12 万。