数据层云原生化:MongoDB、Aerospike 与分布式数据底座

2022-02-11T10:00:00+08:00 | 11分钟阅读 | 更新于 2022-02-11T10:00:00+08:00

@

数据层云原生化(非关系型篇):从「搬家公司」到「跨城双写」

提到数据层的云原生化,大部分人第一反应是 MySQL、Redis、Kafka。但真实的生产链路里,还有两类很容易被面试深挖、却少有人讲透的组件:文档型的 MongoDB内存+SSD 混合的 Aerospike。它们不是「换个部署方式」那么简单,而是要把「均衡器、事务边界、跨机房复制」这些分布式内核机制,和运行在 Kubernetes 上的生命周期(Operator、CronJob、ConfigMap)对齐。

本篇用两条主线串起来:

  • MongoDB 线:分片集群里的 Balancer 像一家「搬家公司」,平时帮你把数据搬匀,大促时却要请它「放假」;多文档事务是「一次打包的银行转账」;Chunk 过大变成 Jumbo 就像箱子塞爆了搬不动,需要手动劈箱。
  • Aerospike 线:namespace 是「带内存缓存的仓库」,二级索引是「货架上贴的标签」,XDR 是「跨城市同步货车」,而回环写入冲突就是「A 城搬来的货又被 B 城搬回 A 城」的死循环。

下面每一节都先给直觉类比,再给可运行的命令 / 代码 / YAML,最后落到面试怎么答。


一、MongoDB on Kubernetes(5.3):分片、均衡与多文档事务

直觉:Balancer 是搬家公司,Chunk 是箱子,Jumbo 是塞爆的箱子

MongoDB 分片集群里有三类角色:

  • Mongos:前台「调度员」,业务只跟它说话,它负责把请求路由到正确的分片。
  • Config Server:后台「户口簿」,记着每个 Chunk(数据箱)归哪个分片。
  • Shard:真正存数据的「仓库」。

数据按 shard key 切成一个个 Chunk(默认 64MB 一个)。Balancer 是一个后台进程,负责发现「某个分片箱子太多、另一个太少」时,把箱子从胖分片搬到瘦分片——这就是「搬家公司」。问题来了:大促期间你最怕它搬家,因为搬迁要占用网络带宽和磁盘 IO,会让在线读写抖动。所以我们要给搬家公司排班:只在凌晨低峰期上班

graph LR
  App[业务 Pod] --> Mongos[Mongos 路由/调度员]
  Mongos --> Cfg[(Config Server
户口簿:Chunk 归属)] Mongos --> S1[(Shard1 仓库)] Mongos --> S2[(Shard2 仓库)] Mongos --> S3[(Shard3 仓库)] Balancer[Balancer 搬家公司] -. 仅在窗口内搬运 Chunk .-> Cfg Cfg -. 变更归属 .-> S1 Cfg -. 变更归属 .-> S2

5.3.1 用 Balancer 窗口避免大促期间搬迁

最简单稳妥的做法是直接在 config 库里设置 activeWindow,让 Balancer 只在 23:00–06:00 工作。注意:设置窗口后,窗口之外 Balancer 自动停摆,所以大促如果横跨白天,搬迁天然不会发生。

// 1) 连上任意一个 Mongos(K8s 里就是 mongodb 的 mongos Service)
// 2) 设置 Balancer 的活跃窗口:只在每天 23:00 到次日 06:00 搬迁
use config;

db.settings.updateOne(
  { _id: "balancer" },
  { $set: { activeWindow: { start: "23:00", stop: "06:00" } } },
  { upsert: true }   // 没有就插入,有就更新
);

// 3) 确认一下当前 Balancer 状态
sh.isBalancerRunning();          // false 表示当前不在窗口内、已停摆
db.settings.findOne({ _id: "balancer" });

但在 Kubernetes + MongoDB Operator 的场景里,更工程化的做法是:用 CronJob 在大促前「强制关门」、大促后「开门」。这样即使有人改了窗口配置,也能用显式的 stop/start 兜底。下面这段 YAML 定义一个大促前夜停 Balancer 的 Job(开门只是把 stop 换成 start)。

# mongo-balancer-stop.yaml —— 大促前夜 20:00 强制停 Balancer
apiVersion: batch/v1
kind: CronJob
metadata:
  name: mongo-balancer-stop
  namespace: mongodb
spec:
  schedule: "0 20 10 11 *"   # 11 月 10 日 20:00(双十一前一天),按需改
  jobTemplate:
    spec:
      template:
        spec:
          restartPolicy: OnFailure
          containers:
            - name: mongo-cli
              image: mongo:6.0
              # 通过 mongosh 连 mongos Service,执行 stopBalancer
              command:
                - /bin/sh
                - -c
                - |
                  mongosh "mongodb://mongos.mongodb.svc.cluster.local:27017/admin" \
                    --eval 'sh.stopBalancer(); print("balancer stopped")'                  

面试要点:Balancer 窗口是「软约束」(靠时间),CronJob 的 stopBalancer() 是「硬开关」。大促保障要用硬开关兜底;窗口之外记得 startBalancer(),否则平时数据倾斜了也不会自动均衡。

5.3.2 分片集群下实现多文档事务并限制超时 3 秒(Java)

分片集群上的多文档事务,本质是把「跨多个分片/集合的若干写操作」包成一个原子单元:要么全成功,要么全回滚——就像「从 A 账户扣钱、给 B 账户加钱」必须一起完成。关键是两件事:

  1. ClientSession 开启事务;
  2. TransactionOptions.maxCommitTime提交阶段设 3 秒上限,避免长事务把锁和资源拖死。
import com.mongodb.ClientSessionOptions;
import com.mongodb.TransactionOptions;
import com.mongodb.client.ClientSession;
import com.mongodb.client.MongoClient;
import com.mongodb.client.MongoClients;
import com.mongodb.client.MongoCollection;
import org.bson.Document;
import java.util.concurrent.TimeUnit;

public class ShardedTxnExample {

    public void transfer(MongoClient client, String from, String to, double amount) {
        // 1) 事务选项:快照读 + 多数派写 + 提交最多 3 秒
        TransactionOptions txnOpts = TransactionOptions.builder()
                .readConcern(com.mongodb.ReadConcern.SNAPSHOT)   // 快照隔离,避免脏读
                .writeConcern(com.mongodb.WriteConcern.MAJORITY) // 多数派确认,防丢
                .maxCommitTime(3, TimeUnit.SECONDS)              // 关键:提交超 3 秒直接失败
                .build();

        // 2) 开会话、起事务
        try (ClientSession session = client.startSession()) {
            session.startTransaction(txnOpts);

            MongoCollection<Document> accounts =
                    client.getDatabase("shop").getCollection("accounts");

            // 3) 两个分片上的写,都在同一个 session 里
            accounts.updateOne(session,
                    new Document("_id", from),
                    new Document("$inc", new Document("balance", -amount)));
            accounts.updateOne(session,
                    new Document("_id", to),
                    new Document("$inc", new Document("balance", amount)));

            // 4) 提交;若 3 秒内没完成,驱动抛 MongoTransactionException
            session.commitTransaction();
        } catch (Exception e) {
            // 5) 任何异常都中止,保证原子性
            // 注意:session 用 try-with-resources 关闭,abort 由驱动在关闭时兜底
            System.err.println("事务失败,已回滚: " + e.getMessage());
            throw e;
        }
    }
}

踩坑提醒:maxCommitTime 限制的是**提交(commit)**耗时,不是整个事务生命周期。事务本身还有服务端 transactionLifetimeLimitSeconds(默认 60 秒),写代码时业务要在更短时间内完成,否则会被服务端强制 kill。

5.3.3 Chunk 超过 64MB 变成 Jumbo 时手动拆分并保证均衡

当一个 Chunk 体积超过默认 64MB,或因为 shard key 是单调递增(如时间戳)导致所有新数据都落进最后一个箱子,Balancer 会把它标记为 jumbo(塞爆了),然后跳过它——再也不搬。结果是这个分片越来越胖,变成热点。

处理四步法:

  1. 找到 jumbo chunk;
  2. sh.splitAt / sh.splitFind 手动劈箱(前提是还有可分的 shard key 区间,单调递增键要先调整);
  3. 清掉 jumbo 标记;
  4. 让 Balancer 重新搬匀。
// 1) 找出所有被标记为 jumbo 的箱子
use config;
db.chunks.find({ jumbo: true }).pretty();

// 2) 手动拆分:假设 shard key 是 { orderId: 1 },在指定值处劈开
//    splitAt 在「精确等于该值」的边界切一刀
sh.splitAt("shop.orders", { orderId: 500000 });

// 如果是范围型热点,用 splitFind 自动在块内找中点切
sh.splitFind("shop.orders", { orderId: 750000 }).split();

// 3) 清掉 jumbo 标记(拆完后箱子变小,标记就没必要了)
db.chunks.updateOne(
  { ns: "shop.orders", min: { orderId: 500000 }, jumbo: true },
  { $unset: { jumbo: "" } }
);

// 4) 确认 Balancer 在窗口内会重新均衡
sh.startBalancer();
sh.status();   // 观察 chunks 在各 shard 的分布是否趋近均匀

根因提醒:如果 jumbo 是因为单一超大数据文档shard key 区分度太低(如所有文档同值),劈箱治标不治本。面试要补一句:「长期方案是选高基数的 shard key,或对单调键加哈希前缀打散。」


二、Aerospike on Kubernetes(5.5):内存+SSD 底座、二级索引与跨城双写

直觉:namespace 是「带内存索引的仓库」,XDR 是「跨城同步货车」

Aerospike 的设计哲学是:主键索引永远在内存里(快),数据本体可以放在 SSD 上(省)。一个 namespace 就像一座仓库:你先决定「货架(内存)多大、仓库(SSD)多大」,比例错了要么内存爆、要么盘浪费。

XDR(Cross Datacenter Replication) 是它跨机房同步的能力。设想 A 城和 B 城各有一座仓库,两边都要写。如果 A 把货搬给 B,B 收到后又原样搬回 A,就形成了回环写入——货车永远在路上。Aerospike 的解法很聪明:只搬运「本地原始写」,从 XDR 同步过来的写不会被再次搬运,回环自然断掉。

graph LR
  subgraph A[A 机房 Cluster]
    NA[namespace: ads
内存索引+SSD 数据] end subgraph B[B 机房 Cluster] NB[namespace: ads
内存索引+SSD 数据] end NA -- "XDR: 仅 ship 本地写" --> NB NB -- "XDR: 仅 ship 本地写" --> NA X[(XDR 标记:
来自同步的写不回 ship)]

5.5.1 namespace 级内存与 SSD 比例设置

核心是两个参数:

  • memory-size:该 namespace 的主索引 +(可选)常驻内存数据占用。经验公式:主键索引约 64 字节/条记录,再加你打算放内存的数据量。
  • storage-engine device:SSD 上的数据文件,用 filesize 控制单文件大小,write-block-size 控制写块(通常 128KB–1MB,越大写放大越小但删除/更新开销越高)。

下面是一段放进 ConfigMapaerospike.conf 片段(Aerospike K8s Operator 通过 ConfigMap 注入配置)。假设 namespace ads 有 5 亿条记录,主键索引 ≈ 500M × 64B ≈ 32GB,数据本体约 200GB 放 SSD,并保留 20% 内存余量。

# aerospike-namespace.yaml —— 通过 ConfigMap 注入 namespace 配置
apiVersion: v1
kind: ConfigMap
metadata:
  name: aerospike-conf
  namespace: aerospike
data:
  aerospike.conf: |
    namespace ads {
      # 内存:主索引 32GB + 20% 余量 ≈ 38GB;data-in-memory 关闭,数据走 SSD
      memory-size 38G
      # 复制因子 2:每个对象在两节点各有副本
      replication-factor 2
      # SSD 存储引擎
      storage-engine device {
        device /dev/nvme0n1        # K8s 里用 PVC/裸盘挂载
        filesize 220G              # 单文件略大于数据量,留碎片空间
        write-block-size 1M        # 1MB 写块,降低写放大
        data-in-memory false       # 数据本体在 SSD,仅索引在内存
        # 内存与 SSD 比例 ≈ 38G : 220G,约 1:6,符合「索引在内存、数据在盘」
      }
    }    

面试要点:比例不是拍脑袋。内存 = 主索引(64B×记录数) + 你要常驻的热数据,SSD = 全量数据 × 副本数 ÷ 节点数。内存算少了索引放不下直接 OOM,算多了浪费。

5.5.2 Secondary Index + Spark SQL 实现广告人群包查询

「人群包」就是:从海量用户里,按一堆属性(年龄、城市、兴趣标签)圈出一批人去做广告投放。Aerospike 的 Secondary Index(二级索引) 相当于在「年龄」「城市」这些非主键字段上贴了货架标签,让 Spark 能快速过滤而不用全表扫。

Aerospike 官方提供 Spark Connector,可以把 namespace 读成 DataFrame,直接跑 Spark SQL。下面用 PySpark 圈一个「北京、25–35 岁、对游戏感兴趣」的人群包。

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("audience-pack") \
    .config("spark.aerospike.namespace", "ads") \
    .config("spark.aerospike.seedhosts", "aerospike.aerospike.svc.cluster.local:3000") \
    .config("spark.aerospike.port", "3000") \
    .getOrCreate()

# 1) 把 users 这个 set 读成 DataFrame(依赖 aerospike-spark-connector 包)
users = spark.read.format("aerospike") \
    .option("aerospike.set", "users") \
    .load()

# 2) 注册成临时视图,用标准 Spark SQL 圈人
users.createOrReplaceTempView("users")

# 3) 二级索引生效点:city / age / interest 上已建好 SI,过滤走索引不扫全表
audience = spark.sql("""
    SELECT user_id
    FROM   users
    WHERE  city    = 'Beijing'
      AND  age BETWEEN 25 AND 35
      AND  array_contains(interest, 'game')
""")

# 4) 把人群包写回 Aerospike 的另一个 set,供投放系统读取
audience.write.format("aerospike") \
    .option("aerospike.set", "pack_game_beijing") \
    .mode("overwrite") \
    .save()

spark.stop()

建二级索引的命令(在 Aerospike 侧,一次性):

# 在 ads 命名空间的 users set 上,给 city 建字符串二级索引
asadm -e "manage indexes create string city on ads.users"

面试要点:二级索引适合等值/范围过滤;人群包查询的本质是「多条件交集」,用 Spark 把结果物化回 Aerospike 比让业务端逐条查快几个数量级。

5.5.3 集群跨机房双写时用 XDR 解决回环写入冲突

双向 XDR(A⇄B 互相同步)最怕回环:A 写 → 同步到 B → B 又当作本地写同步回 A → 无限循环。Aerospike 的默认机制是:XDR 只 ship「本地产生的写」,从远端同步来的写会被打上标记、不再被本端 XDR 搬运,回环因此断掉。

下面是两个机房的 XDR 配置片段(同样进 ConfigMap)。关键点:每个 DC 把对端声明为 datacenter,并依赖「ship-only-local-writes」的默认语义。

# aerospike-xdr.yaml —— A 机房的 XDR 配置(B 机房镜像对称配置)
apiVersion: v1
kind: ConfigMap
metadata:
  name: aerospike-xdr-conf
  namespace: aerospike
data:
  aerospike.conf: |
    namespace ads {
      memory-size 38G
      replication-factor 2
      storage-engine device {
        device /dev/nvme0n1
        filesize 220G
        write-block-size 1M
      }
      # XDR 开启:把本 namespace 的写同步到远端 DC
      xdr {
        enable-xdr true
        # 声明远端 B 机房为一个 datacenter
        datacenter dc-b {
          seed-nodes 10.20.0.10 10.20.0.11  # B 机房节点 IP
          # 关键:默认只 ship 本地写;来自 dc-b 的写不会回 ship,避免回环
          ship-only-local-writes true
        }
      }
    }    

冲突处理补充:即使断掉回环,A、B 两端同时改同一条记录仍可能产生冲突。XDR 默认后写覆盖(last-write-wins),靠记录的 generation/version 裁决。若业务需要更精细的冲突解决,要在应用层做合并逻辑,面试时可点出这一点体现深度。


自测题与动手练习

  1. 概念题:MongoDB 的 Balancer activeWindowstopBalancer() 有什么区别?大促保障为什么建议用后者兜底?
  2. 代码题:在分片集群上用 Java 写一个多文档事务,要求「扣 A 加 B」且提交不超过 3 秒。说明 maxCommitTime 限制的是哪一段耗时。
  3. 排错题sh.status() 显示某分片 Chunk 被标 jumbo 且不再迁移,给出手动拆分 + 清除标记的命令顺序,并说明什么情况下劈箱无效。
  4. 计算题:Aerospike namespace 预计 8 亿条记录、全量数据 320GB、副本因子 2、3 个节点,求单节点 memory-sizefilesize 的参考值(主索引按 64B/条)。
  5. 设计题:A、B 双机房用 XDR 双向同步,如何避免回环?若两端同时改同一条记录,默认冲突策略是什么?

动手练习:在本地用 Docker 起一个 3 分片的 MongoDB(或 MongoDB Kubernetes Operator 的 minikube 部署),插入 10 万条带时间戳 shard key 的数据,观察 Chunk 分布,手动触发一次 splitAt 并清 jumbo 标记。

本章小结

  • MongoDB 分片集群的 Balancer 要「排班 + 硬开关」双保险,大促前 stopBalancer() 兜底,事后 startBalancer() 恢复均衡。
  • 分片多文档事务用 ClientSession + TransactionOptions.maxCommitTime 控提交上限,但要配合合理的事务生命周期,避免被服务端 60 秒上限 kill。
  • Jumbo Chunk 是「箱子塞爆」,手动 splitAt 劈箱 + 清标记 + 重启 Balancer;根因多是 shard key 区分度低或单调键,需从键设计治本。
  • Aerospike 的 namespace 内存/SSD 比例 = 主索引(64B×记录数) + 热数据 : 全量数据,先算后配。
  • 二级索引 + Spark Connector 把人群包圈选变成一条 SQL;XDR 靠「只 ship 本地写」天然断回环,冲突默认后写覆盖。
复习提示:
  • MongoDB 分片集群面试要点:Balancer activeWindow(软排班)vs stopBalancer()(硬开关),大促前用后者兜底;Jumbo Chunk 手动 splitAt 劈箱是常见考点。
  • Aerospike namespace 配置:内存 = 主索引(64B×记录数) + 热数据,SSD = 全量数据;比例算错会导致 OOM 或浪费 SSD。
  • XDR 双机房ship-only-local-writes true 是断回环的关键;跨写冲突默认 last-write-wins,需要应用层合并策略。
  • 下一篇讲 KEDA 弹性伸缩——它和可观测性天然衔接,Prometheus 指标驱动 Pod 扩缩。
面试官
MongoDB 多文档事务的 maxCommitTime 和 MySQL 的事务超时有什么区别?面试时怎么答出深度?
候选人
好问题!两者的本质区别在于:

MySQL 的事务超时是全局的,由 innodb_transaction_timeout 控制,所有事务共享。
MongoDBmaxCommitTime单事务级的,每个 ClientSession 独立设置,粒度更细。

面试加分点:MongoDB 服务端有 60 秒硬上限,maxCommitTime 设 3 秒是为了在业务层尽早失败、避免长时间持有锁。如果 commit 阶段超过 60 秒会被服务端强制 kill,这时候客户端拿到的是"TransactionTooOld"错误,需要在应用层做重试或补偿。知道这个细节说明你踩过坑。
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!