Ch.5: 内存表示 ↔ 字节序列 — 兼容性让滚动升级成为可能, 幂等让重试变得安全, 消息队列把"同步调用"变"异步事实"
编码格式像寄快递的打包方式: JSON是原样打包 — 收件人拆开就能看懂, 但泡沫(引号、字段名)占了一半体积; Protobuf是压缩打包 — 每件物品贴编号 (field tag), 体积省 60%, 但编号永远不能复用。兼容性则是"新旧两套仓库同时营业"的能力: 滚动升级时新代码读旧数据 (向后)、旧代码读新数据 (向前) 都不能崩。而 RPC像把"喊同事帮忙"当"自己动手" — 但喊出去可能没人听见, 也可能办完了你没收到回音 — 所以要么幂等, 要么别重试。
import json b = json.dumps({"id": 42}) # 编码 → b'{"id": 42}' d = json.loads(b) # 解码 → {"id": 42} # pickle 跨服务? 等于让对端执行任意代码
# 向后: v2 代码读 v1 记录 → 缺字段给默认 ✓ # 向前: v1 代码读 v2 记录 → 未知字段忽略 ✓ # 任何一头崩了 → 滚动升级失败
# 集群 5 节点逐个升: 3 新 2 旧共存 1 小时 # 旧节点收到的每条消息都必须能解
9007199254740993 # 2^53+1 # Python: json 往返无损 # JS: JSON.parse('9007199254740993') → 9007199254740992! # 正解: {"order_id": "9007199254740993"} 字符串化
# 1337 的 varint: 1337 = 0b10100111001 # → B9 0A (低7位+续传位 | 高位) 共 2 字节 # JSON: "count":1337 → 13 字节; Protobuf → 4 字节
// proto 演进 message User { string name = 1; reserved 2, 15; // 删掉的 tag 占住不复用 optional string tier = 3; // 新增, 旧 reader 忽略 }
writer: {userName, age, fav}
reader: {user_name, age, tier=free}
# 按名匹配 → 新 reader 也能读旧记录
# 缺 tier → 默认 free; writer 的 fav 多余 → 忽略# 注册新版本被拒: # Schema being registered is incompatible: # field "userName" renamed to "user_name" w/o default # → CI 阻断发布, 兼容事故消失在上线前
# 本地函数: 要么返回要么抛, 状态确定 # 远程调用: 超时 → 三种可能 (没执行/执行了/执行一半) # 语义差距就是事故来源
# SET x=5 → 幂等 (重试安全) # x = x + 1 → 非幂等 (重试 ×3 = +3) # 对: INSERT ... ON CONFLICT(idempotency_key) DO NOTHING
producer.send(topic="orders", msg) # 发完即走 # group A (计费): 各得全部消息 # group B (风控): 各得全部消息 ← topic 扇出 # 组内 3 消费者分摊 ← queue 语义
# durable workflow 伪代码 step1 = charge(card) # 结果已记日志 if replayed: step1 = load_from_log() # 不重复扣 step2 = ship(order)
后端 Java long 的订单号经 JSON 到 JS 前端, 尾位悄悄变了 — 客服查无此单。
# 后端发出: {"order_id": 9007199254740993} // 前端收到: obj.order_id === 9007199254740992 ← 尾位变了! # 根因: 2^53+1 无法用 IEEE 754 double 精确表示 # 修复: 大 ID 一律字符串化 (序列化层全局约定) {"order_id": "9007199254740993"} # 同理: 雪花 ID / 微秒时间戳 / 链上数值 都要字符串化
建立编码规约: int64 一律 string — 一次约定, 终身免疫。
升级窗口 1 小时, 3 新 2 旧共存 — 用测试验证双向兼容再发车。
# CI 兼容性测试矩阵 for w in (v1, v2): for r in (v1, v2): data = encode_v(w)(sample) # w 写 assert decode_v(r)(data).ok # r 读 # 4 个组合全绿才允许滚动升级 # Protobuf 侧配合: 新字段只增不删, 删字段 reserved 占位
升级顺序: 先发"能读新数据"的 reader, 再发 writer — 永远快半步。
同事把 userName 改成 user_name, 若无 registry, 下游消费组集体爆炸。
$ curl -X POST registry/subjects/orders/versions \
-d @orders-v2.avsc
# 409 Conflict
# Schema being registered is incompatible:
# old field "userName" removed; new "user_name" has no default
# 正确演进: 1) 加新字段 user_name 带 default
# 2) 双写过渡 3) 下个版本删旧字段
registry 把"口头约定"变成"机器强制" — 兼容事故清零。
节点间 3 万 QPS 的订单同步, JSON 序列化既是带宽又是 CPU 大头。
# 同一条订单: # JSON: 81 字节 (字段名重复传输, 数字文本化) # Protobuf: 33 字节 (tag + varint) # 收益 (3万 qps × 1KB 载荷): json_bps = 30000 * 1024 # ≈ 293 Mbps proto_bps = 30000 * 420 # ≈ 120 Mbps, -59% # 附赠: 编解码 CPU 降一半 (无文本解析)
迁移路径: 网关对外保持 JSON, 内部 gRPC 双向兼容灰度切换。
湖里百万条记录一个文件, 头部带一次 schema — 新 reader 按名对齐旧数据。
import fastavro # 2023 年的文件, schema v1 (无 tier 字段) with open("orders-2023.avro", "rb") as f: reader = fastavro.reader(f, reader_schema=v2_schema) for rec in reader: # v2 读 v1: print(rec["tier"]) # → "free" (默认值) # writer schema 从文件头读 — 不需要"知道当年"
这正是 schema-on-read + 自描述格式的组合价值 (呼应模型页)。
支付请求超时后不敢重试 — 用幂等键把"不知道执行没执行"变成"重试安全"。
# 客户端: 生成幂等键, 重试时复用 key = uuid4() POST /pay {"amount": 100, "idempotency_key": key} # 超时 → 重试, 带同一个 key # 服务端: key 唯一约束, 重复请求返回首次结果 INSERT INTO payments(idem_key, amount, status) VALUES ($key, 100, 'done') ON CONFLICT (idem_key) DO NOTHING RETURNING *; # 已存在 → 读出并返回旧结果
规则: 所有非幂等 RPC 必须带幂等键 — 这是重试的入场券。
下单峰值 5 万 QPS, 库存库只能扛 1 万 — 队列把"瞬时洪峰"变"持续消费"。
# 生产者: API 收单后立刻返回 (用户无感) producer.send("order-created", order) # 5万/s 洪峰入队 # 消费者: 固定速率消费 while True: batch = consumer.poll(max=500) # 1万/s 稳定消费 db.insert_orders(batch) # 攒批写库 # 峰值期间 lag 涨到 40 万条, 10 分钟消化完 — 系统无感 # 红线: lag 增长率告警 (响应总纲页背压思想)
代价: 落库延迟从 5ms 变为秒级 — 业务必须接受"最终一致的下单"。
退款流程 5 步跨 3 天, 第 4 步崩了重跑不能重复退款 — 每步落日志。
# Temporal 式 workflow (伪代码) async def refund(order): # 每步结果写入持久日志, 重放时跳过已完成 check = await step(verify_order) # 已完成 → 直接取结果 pay = await step(reverse_charge) # 幂等 + 结果落盘 await step(notify_user) # 崩溃重启: 引擎重放日志, 未完成的步骤才真正执行
本质: 把"执行到哪了"从内存搬进日志 — 长流程的断点续传。
"这个服务在哪"的答案晚 1 秒没关系, 但不能全体查不到 — 共识成本别乱花。
# 需求分析: # 错误路由到刚下线的实例 → 重试即恢复 (自愈) # 全体查不到注册表 → 全站故障 (灾难) # → 要 HA + 低延迟, 不需要强一致 # 落地: 客户端负载均衡 + 最终一致注册表 (Eureka/Consul) instances = discovery.lookup("price-svc") # 允许略旧 target = balancer.pick(instances) # 失败重试兜底
对照: 选主/元数据才需要 ZooKeeper/etcd 级别的共识 (见共识页)。
字段要破坏性变更 — 用版本头分流, 新旧 consumer 各读各的, 稳定后再删旧。
# 网关按版本头路由 if header("X-Api-Version") == "2": route("orders-v2") # 新响应结构 else: route("orders-v1") # 旧结构, 保留 90 天 # 灰度节奏: 5% → 25% → 100%, 监控两版 P99/error 对比 # v1 流量归零后下线 — 演进有退出机制
配合 schema registry: 消息类用兼容演进, API 类用版本并行 — 二选一别硬扛。
# 错: {"order_id": 9007199254740993} # 对: {"order_id": "9007199254740993"}
# 错: "2026/9/26" 与 1695724800 与 "09-26" 并存 # 对: 统一 epoch_ms: 1695724800000
# Proto: name → user_name (tag=1 不变) ✓ # Avro: 同样改名 = 删字段+加字段, 必须带 default
# 错: tag=2 删了又给新字段用 # 对: reserved 2; 新字段用 tag=7
# 错: required string tier = 3; (旧 reader 报错) # 对: optional + 服务端校验
# 错: "schema 文件在群里, 大家记得更新" # 对: registry BACKWARD 校验不过不许发
# 错: 只 assert 不抛异常 # 对: assert rec["tier"] == expected_default
# 错: 只测 v2读v1 # 对: v1读v1/v1读v2/v2读v1/v2读v2
# 错: result = stub.Charge(x) 无超时无重试策略 # 对: timeout + retry budget + 幂等键
# 错: 超时 → retry(POST /pay) # 对: /pay 带 idempotency_key 再谈重试
# 错: for msg: registry.fetch(schema_id) # 对: lru_cache(schema_id) 本地持有
# 错: 兼容性只对 v2 验证 # 对: registry 记录在用版本, 全部纳入 CI
# 错: 发布即全量新格式 # 对: 5% 灰度 → 对比解析成功率 → 全量
# 错: "订单数据在 Kafka 里, 不用存库" # 对: 消费落库, log compaction 只存最新状态
# 错: send(topic, msg) 轮询分区 # 对: send(topic, key=order_id) 同键同分区
# 错: 处理完才 commit offset → 崩溃后重投 # 对: INSERT ... ON CONFLICT DO NOTHING
# 错: mailbox 无限收 # 对: bounded mailbox + overflow → DLQ
# 错: 全网格化后 P99 +3ms 没人发现 # 对: mesh 收益 vs 延迟税分开记账
# 错: 直接改 /orders 响应结构 # 对: /v2/orders 并行, 旧版给下线时间表
# 错: await send_sms(user) 直接调 # 对: step(send_sms) 结果落工作流日志