DDIA · 编码与服务数据流

Ch.5: 内存表示 ↔ 字节序列 — 兼容性让滚动升级成为可能, 幂等让重试变得安全, 消息队列把"同步调用"变"异步事实"

同一条记录的四种体积 (field tag 的功劳) JSON 81B MsgPack 66B Protobuf 33B Avro 32B 省在哪: 字段名 → tag 编号; 整数 → varint 变长 例: 1337 → varint 2 字节; JSON 里 4 字符还带引号 向后兼容: 新代码读旧数据 (易) 向前兼容: 旧代码读新数据 (难, 忽略未知) Avro: writer schema + reader schema 共同解码 writer schema v1: userName: string age: int fav: string 编码时用它 — 随数据携带/注册表取 reader schema v2: user_name: string age: int tier: string (默认 free) 按字段名匹配: 缺→默认 多→忽略 对齐 schema registry 集中存版本 + 校验 不兼容发布 → 阻断 规则示例: BACKWARD / FULL (Confluent) 无 tag 格式 (Avro) 按字段名解析 — 改名 = 删字段+加字段, 必须给默认值 有 tag 格式 (Protobuf) 按编号解析 — tag 永不复用, 未知字段必须保留透传 文件头带一次 schema 可存百万条记录 — 对象容器文件 (object container file) 服务间数据流: 同步调用 vs 异步事实 REST / RPC 请求-响应, 同步等待 gRPC 服务 Protobuf 编解码 RPC 的最大谎言: "像调本地函数一样" 超时后结果未知: 可能没执行, 也可能执行了但响应丢了 → 重试前先问: 这个操作幂等吗? message broker: 生产者 → [queue/topic] → 消费者 queue: 每条消息一个消费者 (工作分摊) topic: 每个订阅者各得一份 (fan-out 广播) 解耦: 发送方不必等消费方在线 — 削峰、重投、异步 actor 模型: 私有状态 + 逐条异步消息 (Akka/Orleans) — 把队列内化到进程里 兼容性 = 滚动升级的门票 发布中: 新旧实例并存 v1 ←请求→ v2 的数据互通必须无损 数据比代码活得久 5 年前的文件用今天的 reader 读 安全的 schema 演进 (Protobuf 为例) ✓ 新增字段: 给新 tag, reader 旧版忽略未知 ✓ 删字段: reserved 保留 tag, 永不复用 ✗ 改字段编号 ✗ 改类型 ✗ required 改可选 发布顺序: 先让 reader 兼容新格式 → 再升级 writer 即: 双向兼容是滚动升级的前提, 不是可选优化 durable execution (Temporal 式): 把每步结果记日志 失败重放时跳过已完成的 RPC — 长流程的"进度条"存在日志里

编码三大家

  • • JSON: 人读友好, 但大整数/二进制/体积全有坑
  • • Protobuf: tag+varint, 33B, 未知字段透传
  • • Avro: 无 tag 按名解析, writer/reader 双 schema
  • • 体积: 81 → 66 → 33 → 32 字节

兼容性规则

  • • 向后 (新读旧) 易; 向前 (旧读新) 难
  • • tag 永不复用; 按名解析禁裸改名
  • • registry 挡住不兼容发布
  • • 双向兼容 = 滚动升级的门票

数据流三形态

  • • REST/RPC: 同步, 超时后结果未知
  • • broker: queue 分摊 / topic 广播
  • • actor: 状态私有 + 消息逐条
  • • 重试安全的前提只有一个: 幂等

💡 一句话理解

编码格式像寄快递的打包方式: JSON是原样打包 — 收件人拆开就能看懂, 但泡沫(引号、字段名)占了一半体积; Protobuf是压缩打包 — 每件物品贴编号 (field tag), 体积省 60%, 但编号永远不能复用。兼容性则是"新旧两套仓库同时营业"的能力: 滚动升级时新代码读旧数据 (向后)、旧代码读新数据 (向前) 都不能崩。而 RPC像把"喊同事帮忙"当"自己动手" — 但喊出去可能没人听见, 也可能办完了你没收到回音 — 所以要么幂等, 要么别重试。

🧠 必知必会 必考 & 必会

encoding 双向转换
序列化: 内存对象 → 字节序列; 反序列化反向。语言内置格式 (Java Serializable/pickle) 只在"同语言可信环境"安全, 跨服务禁用。
import json
b = json.dumps({"id": 42})   # 编码 → b'{"id": 42}'
d = json.loads(b)              # 解码 → {"id": 42}
# pickle 跨服务? 等于让对端执行任意代码
向后 / 向前兼容
向后: 新代码读旧数据 (较易); 向前: 旧代码读新数据, 靠"忽略不认识的内容" (较难)。部署时两者都要。
# 向后: v2 代码读 v1 记录 → 缺字段给默认 ✓
# 向前: v1 代码读 v2 记录 → 未知字段忽略 ✓
# 任何一头崩了 → 滚动升级失败
滚动升级需要它
逐节点重启期间新旧实例并存: v2 实例可能把数据写给 v1, v1 又要读 — 兼容性是硬要求不是优化。
# 集群 5 节点逐个升: 3 新 2 旧共存 1 小时
# 旧节点收到的每条消息都必须能解
JSON 精度坑
超过 2^53 的整数经 IEEE 754 double 失真 — JS 里订单号变错号。大 ID 一律传字符串。
9007199254740993  # 2^53+1
# Python: json 往返无损
# JS:     JSON.parse('9007199254740993') → 9007199254740992!
# 正解: {"order_id": "9007199254740993"} 字符串化
field tag 与 varint
二进制格式用 tag 编号替代字段名 (省体积), 整数用 varint: 每字节 7 位有效 + 续传标志, 1337 只要 2 字节。
# 1337 的 varint: 1337 = 0b10100111001
# → B9 0A (低7位+续传位 | 高位)  共 2 字节
# JSON: "count":1337 → 13 字节; Protobuf → 4 字节
Protobuf 演进铁律
tag 永不复用、未知字段保留透传; 删字段用 reserved 占位。违反任何一条 = 兼容性静默崩坏。
// proto 演进
message User {
  string name = 1;
  reserved 2, 15;          // 删掉的 tag 占住不复用
  optional string tier = 3; // 新增, 旧 reader 忽略
}
Avro 双 schema
无 tag、值按序拼接: 解码需要 writer schema (随数据携带) + reader schema (代码持有) 对字段名匹配; 缺填默认、多则忽略。
writer: {userName, age, fav}
reader: {user_name, age, tier=free}
# 按名匹配 → 新 reader 也能读旧记录
# 缺 tier → 默认 free;  writer 的 fav 多余 → 忽略
schema registry
集中存 schema 版本并校验兼容性 (BACKWARD/FULL): 不兼容的 schema 提交直接被拒 — 把"约定"变"强制"。
# 注册新版本被拒:
# Schema being registered is incompatible:
#   field "userName" renamed to "user_name" w/o default
# → CI 阻断发布, 兼容事故消失在上线前
REST vs RPC
REST: HTTP 资源 + 统一语义, 适合对外; RPC/gRPC: 像调函数, 内部高频调用省事省带宽。但"location transparency"是错的 — 远程会超时、会部分失败。
# 本地函数: 要么返回要么抛, 状态确定
# 远程调用: 超时 → 三种可能 (没执行/执行了/执行一半)
# 语义差距就是事故来源
idempotence 幂等
执行多次效果同一次 — 网络重试安全的前提。天然幂等: GET/SET/DELETE by id; 非幂等: 扣款/加计数 — 要幂等键。
# SET x=5   → 幂等 (重试安全)
# x = x + 1 → 非幂等 (重试 ×3 = +3)
# 对: INSERT ... ON CONFLICT(idempotency_key) DO NOTHING
message broker
缓冲、重投、扇出消息的中间件: queue 每条一个消费者 (分摊), topic 每订阅者一份 (广播); 发送方与消费方解耦, 削峰填谷。
producer.send(topic="orders", msg)   # 发完即走
# group A (计费): 各得全部消息
# group B (风控): 各得全部消息   ← topic 扇出
# 组内 3 消费者分摊               ← queue 语义
actor 与 durable execution
actor: 私有状态 + 逐条异步处理, 无锁并发; durable execution: 每步结果记日志, 失败重放跳过已完成步骤 (Temporal) — 长流程的容错。
# durable workflow 伪代码
step1 = charge(card)          # 结果已记日志
if replayed: step1 = load_from_log()  # 不重复扣
step2 = ship(order)

🏭 生产实战 real world

场景 1 · 大整数精度事故: 订单号"变异"

后端 Java long 的订单号经 JSON 到 JS 前端, 尾位悄悄变了 — 客服查无此单。

# 后端发出: {"order_id": 9007199254740993}
// 前端收到: obj.order_id === 9007199254740992  ← 尾位变了!

# 根因: 2^53+1 无法用 IEEE 754 double 精确表示
# 修复: 大 ID 一律字符串化 (序列化层全局约定)
{"order_id": "9007199254740993"}
# 同理: 雪花 ID / 微秒时间戳 / 链上数值 都要字符串化

建立编码规约: int64 一律 string — 一次约定, 终身免疫。

场景 2 · 滚动升级演练: 新旧实例互读写数据

升级窗口 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 — 永远快半步。

场景 3 · schema registry 拦截不兼容发布

同事把 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 把"口头约定"变成"机器强制" — 兼容事故清零。

场景 4 · JSON 换 Protobuf: 内网带宽省 60%

节点间 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 双向兼容灰度切换。

场景 5 · Avro 容器文件: 数仓数据 3 年后还能读

湖里百万条记录一个文件, 头部带一次 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 + 自描述格式的组合价值 (呼应模型页)。

场景 6 · RPC 超时的"结果未知": 幂等键设计

支付请求超时后不敢重试 — 用幂等键把"不知道执行没执行"变成"重试安全"。

# 客户端: 生成幂等键, 重试时复用
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 必须带幂等键 — 这是重试的入场券。

场景 7 · 消息队列削峰: 大促订单的异步落库

下单峰值 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 变为秒级 — 业务必须接受"最终一致的下单"。

场景 8 · durable execution: 长流程失败从断点继续

退款流程 5 步跨 3 天, 第 4 步崩了重跑不能重复退款 — 每步落日志。

# Temporal 式 workflow (伪代码)
async def refund(order):
    # 每步结果写入持久日志, 重放时跳过已完成
    check = await step(verify_order)       # 已完成 → 直接取结果
    pay   = await step(reverse_charge)     # 幂等 + 结果落盘
    await step(notify_user)
# 崩溃重启: 引擎重放日志, 未完成的步骤才真正执行

本质: 把"执行到哪了"从内存搬进日志 — 长流程的断点续传。

场景 9 · 服务发现选型: HA 和低延迟比一致性重要

"这个服务在哪"的答案晚 1 秒没关系, 但不能全体查不到 — 共识成本别乱花。

# 需求分析:
#   错误路由到刚下线的实例 → 重试即恢复 (自愈)
#   全体查不到注册表      → 全站故障 (灾难)
# → 要 HA + 低延迟, 不需要强一致

# 落地: 客户端负载均衡 + 最终一致注册表 (Eureka/Consul)
instances = discovery.lookup("price-svc")   # 允许略旧
target = balancer.pick(instances)            # 失败重试兜底

对照: 选主/元数据才需要 ZooKeeper/etcd 级别的共识 (见共识页)。

场景 10 · API 版本灰度: header 路由而不是复制接口

字段要破坏性变更 — 用版本头分流, 新旧 consumer 各读各的, 稳定后再删旧。

# 网关按版本头路由
if header("X-Api-Version") == "2":
    route("orders-v2")       # 新响应结构
else:
    route("orders-v1")       # 旧结构, 保留 90 天

# 灰度节奏: 5% → 25% → 100%, 监控两版 P99/error 对比
# v1 流量归零后下线 — 演进有退出机制

配合 schema registry: 消息类用兼容演进, API 类用版本并行 — 二选一别硬扛。

⚠️ 编码注意与常见坑 pitfalls

坑 1 · 大整数当 JSON number — 2^53 以上尾位静默变错. 原因: IEEE 754 double 精度. 正解: int64 一律字符串化。
# 错: {"order_id": 9007199254740993}
# 对: {"order_id": "9007199254740993"}
坑 2 · 时间字段格式五花八门 — 秒/毫秒/ISO 混用, 时区各写各的. 原因: 没有编码规约. 正解: 全司 UTC ISO8601 或 epoch 毫秒二选一。
# 错: "2026/9/26" 与 1695724800 与 "09-26" 并存
# 对: 统一 epoch_ms: 1695724800000
坑 3 · Protobuf 字段"改名"当兼容 — tag 没变名变了其实兼容, 但 Avro 按名匹配就是破坏. 原因: 两种机制混淆. 正解: Proto 靠 tag, Avro 改名=删+加带默认。
# Proto: name → user_name (tag=1 不变) ✓
# Avro:  同样改名 = 删字段+加字段, 必须带 default
坑 4 · 复用已删字段的 tag — 新字段用了旧 tag, 旧数据被解释成新类型. 原因: 不知道 reserved. 正解: 删字段必须 reserved 占位。
# 错: tag=2 删了又给新字段用
# 对: reserved 2; 新字段用 tag=7
坑 5 · required 字段锁死演进 — 加字段想 required, 旧 reader 直接拒读. 原因: 语义理解反了. 正解: 一律 optional + 应用层校验必填。
# 错: required string tier = 3;  (旧 reader 报错)
# 对: optional + 服务端校验
坑 6 · Avro 无 registry 手工同步 — writer/reader schema 各改各的, 上线才炸. 原因: 约定靠人. 正解: registry 强制校验 + CI 阻断。
# 错: "schema 文件在群里, 大家记得更新"
# 对: registry BACKWARD 校验不过不许发
坑 7 · 忽略缺省值语义 — reader 给了默认值, writer 旧数据没这字段, 语义悄悄漂移. 原因: 只测"能读"不测"读对". 正解: 兼容测试断言字段值。
# 错: 只 assert 不抛异常
# 对: assert rec["tier"] == expected_default
坑 8 · 向前兼容当免费 — "旧代码读新数据没事吧" — 遇到必填缺失/类型变化就崩. 原因: 只测了向后. 正解: 兼容矩阵 4 组合全测。
# 错: 只测 v2读v1
# 对: v1读v1/v1读v2/v2读v1/v2读v2
坑 9 · 把 RPC 当本地函数 — 忽略超时/部分失败/重试副作用. 原因: location transparency 迷信. 正解: 显式超时 + 幂等键 + 降级路径。
# 错: result = stub.Charge(x)  无超时无重试策略
# 对: timeout + retry budget + 幂等键
坑 10 · 重试非幂等 RPC — 网络抖动重试, 用户被扣三次款. 原因: 重试模板无脑套. 正解: 非幂等操作先做幂等化改造。
# 错: 超时 → retry(POST /pay)
# 对: /pay 带 idempotency_key 再谈重试
坑 11 · 每请求重复解析 schema — 每条消息都读注册表/编译 schema, CPU 空转. 原因: 图省事. 正解: schema 按 id 本地缓存 (version → parsed)。
# 错: for msg: registry.fetch(schema_id)
# 对: lru_cache(schema_id) 本地持有
坑 12 · 只测新 reader 不测旧 reader — 升级 writer 后还活着的旧消费者没人管. 原因: 视野只有最新版. 正解: 在线旧版本清单纳入兼容测试。
# 错: 兼容性只对 v2 验证
# 对: registry 记录在用版本, 全部纳入 CI
坑 13 · schema 变更全量直切 — 新格式一把梭, 出问题回不去. 原因: 无灰度. 正解: 双写/双读灰度, 验证后切流。
# 错: 发布即全量新格式
# 对: 5% 灰度 → 对比解析成功率 → 全量
坑 14 · 消息队列当数据库 — 业务状态靠"消息还在队列里"维持, 堆积数周. 原因: 职责错位. 正解: 队列只做传输, 状态进库, lag 有上限。
# 错: "订单数据在 Kafka 里, 不用存库"
# 对: 消费落库, log compaction 只存最新状态
坑 15 · topic 无分区键 — 同一订单的消息散到不同分区, 顺序全乱. 原因: 默认轮询分区. 正解: 按业务 key 分区 (order_id)。
# 错: send(topic, msg) 轮询分区
# 对: send(topic, key=order_id) 同键同分区
坑 16 · 消费者无幂等 — rebalance/重投后同一条消息处理两次. 原因: at-least-once 语义没消化. 正解: 消费端幂等 (唯一键/版本号)。
# 错: 处理完才 commit offset → 崩溃后重投
# 对: INSERT ... ON CONFLICT DO NOTHING
坑 17 · actor 信箱无上限 — 慢 actor 的信箱把内存吃穿. 原因: 异常消息堆积. 正解: 信箱有界 + 满则背压/丢弃策略。
# 错: mailbox 无限收
# 对: bounded mailbox + overflow → DLQ
坑 18 · service mesh 当万能药 — 每跳加 sidecar 两跳, P99 敏感服务雪上加霜. 原因: 跟风上 Istio. 正解: 度量延迟税, 核心链路可直连。
# 错: 全网格化后 P99 +3ms 没人发现
# 对: mesh 收益 vs 延迟税分开记账
坑 19 · REST API 无版本 — 破坏性变更直接上线, 客户端集体 400. 原因: 没有演进预案. 正解: 版本头/路径版本 + 双版本过渡期。
# 错: 直接改 /orders 响应结构
# 对: /v2/orders 并行, 旧版给下线时间表
坑 20 · durable execution 外部副作用不记录 — 重放时把"发短信"又执行一遍. 原因: 只记了数据库步骤. 正解: 所有外部调用结果都进日志。
# 错: await send_sms(user) 直接调
# 对: step(send_sms) 结果落工作流日志