引言
大多数人第一次意识到「数据传递」是个问题时,工作流已经在生产上跑了几周。表现是实例列表页加载越来越慢、引擎数据库的磁盘占用线性增长、某个跑了一个月的实例重放要花 8 秒。排查后发现根源是某一步任务返回了 20 万行明细,而引擎把每一步的输入输出都持久化了。
工作流引擎与普通服务在数据传递上有一个根本差异:引擎会持久化每一步的输入与输出,因为它们既是重放的依据,也是断点恢复的凭据。这个设计带来了可观测性与可靠性,也带来了两个硬约束:单次传递的数据有体积上限,累计传递的数据会无限增长。
第二个差异是数据的生命周期远长于代码。一个跨月实例里保存的 JSON 是三个月前写的代码产生的,期间的 Schema 变更必须保证它仍能正确反序列化。这与微服务之间的 API 契约不同——API 调用是即时的,工作流的状态是持久的。
本文按「传什么 → 怎么传 → 怎么演进 → 怎么保护」的顺序展开:先讲值传递与引用传递的边界(Claim Check 模式),再讲 Schema 定义与兼容规则,然后是序列化格式的选型与实测数据,最后是敏感数据的加密、脱敏与生命周期治理。引擎的序列化实现差异参见 Temporal 与持久化执行 ,数据管道的编排视角参见 Airflow DAG 调度体系 。
目录
- 数据传递为什么是隐藏难点
- 三种传递方式:值、引用与共享存储
- 数据契约的设计原则
- 值传递的边界与体积上限
- Claim Check 模式
- 中间结果的存放位置选择
- Schema 定义与注册中心
- Schema 演进的兼容规则
- 序列化格式选型
- 序列化的性能与体积实测
- 数据在实例状态中的累积
- 敏感数据的识别与分级
- 加密与密钥管理
- 脱敏、掩码与访问控制
- 数据生命周期与清理
- 跨引擎的传递机制差异
- 类型系统与运行时校验
- 契约测试与数据演练
- 观测与容量规划
- 落地路线图
- 权衡取舍
- 常见坑清单
- 小结
1. 数据传递为什么是隐藏难点
工作流引擎的持久化模型决定了数据传递的特殊性。以持久化执行引擎为例,每次任务调度的输入参数与返回结果都会被写进事件历史,因为重放时需要它们来「跳过已完成的部分」并重建状态。这意味着:
一次任务调用 = 历史里的一条 Scheduled 事件 + 一条 Completed 事件(含完整结果)
一个 50 步的流程 = 至少 100 条事件,每条都可能携带数据
一个每天跑 1000 次的流程 = 每天 10 万条事件,数据量取决于单次返回大小
三个后果随之而来。体积上限:多数引擎对单个载荷有硬限制(Temporal 默认 2 MB,Step Functions 是 256 KB),超过会直接失败。累积增长:单次返回 10 KB 看起来很小,但 50 步乘 10 万次就是 50 GB。兼容性约束:历史里的数据是旧结构,新代码必须能读它。
还有第四个更隐蔽的后果:数据传递的错误往往在很久以后才暴露。一个体积过大的返回值不会立刻报错,它只是让历史慢慢变大,直到某天重放超时或撞上历史上限。因此数据传递必须有前置的设计规范,而不是等到出问题再补。
2. 三种传递方式:值、引用与共享存储
任务之间传递数据只有三种方式,选择标准是「数据的大小与生命周期」:
| 方式 | 机制 | 适用数据量 | 优点 | 缺点 |
|---|---|---|---|---|
| 值传递 | 数据直接放在实例状态/事件里 | < 100 KB | 简单、可重放、可观测 | 有体积上限、拖慢重放 |
| 引用传递 | 传 URI 或 key,数据在外部存储 | 任意 | 无上限、状态轻量 | 需要外部存储与清理 |
| 共享存储 | 上下游读写同一位置(表、目录) | 大 | 零拷贝、适合大数据 | 隐式依赖、难追溯 |
引用传递(Claim Check 模式)是生产系统的默认选择。原则很简单:工作流状态里只放「控制信息」(ID、状态、路径、计数、错误码),业务数据放外部存储。判断标准是「这个数据是否需要参与重放决策」——需要(比如分支条件依赖的金额)就传值,不需要(比如展示用的明细列表)就传引用。
共享存储的典型形态是「上游写分区、下游读分区」,在数据编排场景里很常见(Airflow 的 DAG 里任务之间通过数仓表传递)。它的优点是零拷贝,缺点是依赖关系隐含在表名里,出问题时难以定位「是哪个任务写了这张表」。
3. 数据契约的设计原则
数据契约是任务之间对「我传给你什么」的明确约定。没有契约的工作流会在多人协作时迅速腐化:A 改了返回结构,B 的解析代码静默出错。
四条设计原则:
1. 每个任务的输入输出都有显式的类型定义(类、struct、schema),不用裸 Map
2. 契约的字段命名带业务语义,不用 data1 / payload / result 这类泛名
3. 契约一经发布即视为对外接口,变更走版本化流程
4. 契约里不放「只有本步骤知道含义」的内部字段,避免下游误用
第二条最容易被忽视,但它的收益最大。payload 这种字段名在下游代码里会变成 input.payload.items[0].value,半年后没人能说清 value 是什么。而 orderAmountCents、refundReasonCode 这类名字自带语义,且天然适合做类型检查。
契约的载体应该与语言类型系统绑定,而不是靠文档。在 Java 里是一个不可变的 record,在 Python 里是 dataclass 或 Pydantic 模型,在 Go 里是 struct。类型定义本身就是契约,文档只是补充。
4. 值传递的边界与体积上限
值传递有明确的硬上限,且不同引擎的默认值差异很大,迁移时必须重新评估:
| 引擎 | 单载荷上限 | 历史/状态上限 | 超限行为 |
|---|---|---|---|
| Temporal | 2 MB(可配置) | 50 MB 或 50K 事件 | 调度失败 / 强制截断 |
| Step Functions | 256 KB | 无(状态在外部) | 执行失败 |
| Camunda 7 | 无硬限制(存数据库) | 受数据库行大小限制 | 数据库写入失败 |
| Airflow XCom | 元数据库字段限制 | 元数据库容量 | 性能急剧下降 |
| Zeebe | 4 MB(可配置) | 无 | 命令被拒绝 |
经验阈值是 100 KB:低于这个量级传值没有明显代价,超过之后重放时间与存储成本开始非线性增长。一个粗略的估算方法是「单次载荷 × 步骤数 × 日实例数 × 保留天数」,超过 10 GB 就应该考虑改为引用传递。
// 反例:把明细列表直接放进实例状态
public class ExtractResult {
List<OrderRow> rows; // 20 万行,几十 MB
}
// 正例:只传引用与统计信息
public class ExtractResult {
String storageUri; // s3://bucket/staging/orders/2026-10-07.parquet
long rowCount; // 100000
String checksum; // sha256:...
Instant producedAt;
}
正例里的 rowCount 与 checksum 不是冗余——它们是下游做校验的依据(行数不符或校验和不匹配时可以直接失败,而不是处理错误数据)。这些「控制信息」正是值传递该承载的内容。
5. Claim Check 模式
Claim Check 模式的名字来自行李寄存:把大件行李存起来,只带一张票据走。在工作流里,票据就是对象存储的 URI。
@task
def extract(data_interval_start, data_interval_end) -> str:
df = query_orders(data_interval_start, data_interval_end)
uri = f"s3://staging/orders/{data_interval_start:%Y-%m-%d}.parquet"
write_parquet(df, uri)
return uri # 只返回路径
@task
def transform(uri: str) -> str:
df = read_parquet(uri)
out = f"{uri}.transformed.parquet"
write_parquet(normalize(df), out)
return out
这个模式的四个工程要点:
一是路径的确定性。路径必须由「数据区间 + 任务名 + 版本」唯一确定,而不是随机 UUID。确定性的路径让任务天然幂等(重跑覆盖同一路径),也让排查时能直接从路径推断出是哪个区间、哪一步产生的。
二是写入的原子性。写对象存储时先写临时路径再原子重命名,避免下游读到写了一半的文件。S3 没有原生 rename,做法是「先写 _SUCCESS 标记文件,下游先检查标记再读」。
三是清理策略。中间文件不会自动消失,必须有生命周期规则(S3 Lifecycle 或定时清理任务),按「最后一次被引用的时间 + 保留期」删除。没有清理策略的对象存储账单会失控。
四是引用失效的处理。如果中间文件被清理了但实例还在跑,下游读取会失败,防御方式是让引用里带上「最小保留截止时间」,清理任务据此跳过。
6. 中间结果的存放位置选择
外部存储的选择影响成本、延迟与运维复杂度:
| 存储 | 延迟 | 成本 | 适用 |
|---|---|---|---|
| 对象存储(S3/OSS) | 10~100 ms | 极低 | 大文件、Parquet、归档 |
| 数据库表 | 1~10 ms | 中 | 结构化中间结果、需要事务 |
| 分布式缓存(Redis) | < 1 ms | 高 | 临时小数据、TTL 短 |
| 分布式文件系统(HDFS) | 5~50 ms | 中 | 大数据集、批处理 |
| 消息队列 | 1~10 ms | 低 | 流式传递、需要扇出 |
选择标准是**「下游怎么消费」**:下游要按条件筛选就用数据库(能下推过滤),下游要整批扫描就用对象存储(吞吐高成本低),下游要低延迟随机读就用缓存。
一个常见的混合模式是「对象存储放明细 + 数据库放索引」:明细写 Parquet 到对象存储,同时在数据库里写一条元数据记录(路径、行数、区间、校验和),下游先查元数据再决定是否读文件。这样既避免了把大对象塞进数据库,又保留了「按条件查询有哪些数据集」的能力。
CREATE TABLE flow_artifact (
artifact_id VARCHAR(64) PRIMARY KEY,
run_id VARCHAR(64) NOT NULL,
step_name VARCHAR(64) NOT NULL,
storage_uri VARCHAR(512) NOT NULL,
row_count BIGINT,
checksum VARCHAR(80),
expires_at TIMESTAMP NOT NULL,
INDEX idx_run_step (run_id, step_name),
INDEX idx_expires (expires_at)
);
7. Schema 定义与注册中心
当契约数量增长到几十个、由多个团队维护时,需要一个集中的 Schema 注册中心来管理版本与兼容性校验。它的核心能力是在写入时校验兼容性:生产者注册新版本 Schema 时,注册中心检查它与历史版本的兼容关系,不兼容则拒绝注册。
# 注册新版本 schema,指定兼容性级别
curl -X POST "$REGISTRY/subjects/order-extracted-value/versions" \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
-d '{"schema": "...", "schemaType": "AVRO"}'
# 检查兼容性(不实际注册)
curl -X POST "$REGISTRY/compatibility/subjects/order-extracted-value/versions/latest" \
-d '{"schema": "..."}'
四种兼容性级别,工作流场景的选择很明确:
BACKWARD 新代码能读老数据(工作流最需要:老实例的数据被新代码读)
FORWARD 老代码能读新数据(回滚时需要)
FULL 双向兼容(最严格,推荐)
NONE 不校验(危险,只在明确的破坏性变更时用)
工作流的特殊之处在于两个方向都必须兼容:正常运行需要 BACKWARD(新代码读老状态),回滚需要 FORWARD(老代码读新状态)。所以默认应该选 FULL,除非有明确的理由放宽。
8. Schema 演进的兼容规则
兼容性级别落到具体操作上,就是一张「能不能做」的表:
| 变更 | BACKWARD | FORWARD | FULL |
|---|---|---|---|
| 新增可选字段(带默认值) | 可以 | 可以 | 可以 |
| 新增必填字段 | 不可以 | 可以 | 不可以 |
| 删除字段 | 可以 | 不可以 | 不可以 |
| 字段改类型 | 不可以 | 不可以 | 不可以 |
| 字段改名 | 不可以 | 不可以 | 不可以 |
| 扩大枚举范围 | 可以 | 不可以 | 不可以 |
| 缩小枚举范围 | 不可以 | 可以 | 不可以 |
| 提高字段精度(int→long) | 可以 | 不可以 | 不可以 |
实践中的三条经验规则:
规则一:所有新字段必须有默认值。这是唯一能同时满足两个方向的加法变更。没有默认值的必填字段会破坏 BACKWARD 兼容。
规则二:删除字段要先标记废弃,再物理删除。先用 deprecated 标注并保留至少两个发布周期,确认所有实例都不再读写它之后再删。Protobuf 用 reserved 关键字保留编号与名字,Avro 用 default 让老数据有值可填。
规则三:枚举的演进方向要提前想清楚。新增枚举值是 BACKWARD 兼容(老数据不会有新值),但删除或改语义是破坏性的。如果枚举可能频繁变化,考虑改用字符串 + 校验表。
// 向后兼容的字段新增:带默认值,且用包装类型区分「未设置」与「默认值」
public record OrderInput(
String orderId,
long amountCents,
@JsonProperty(defaultValue = "unknown") String channel, // 新增,带默认值
Optional<String> couponCode // 新增,可为空
) {}
9. 序列化格式选型
序列化格式决定体积、速度与可演进性,三者往往互相冲突:
| 格式 | 体积(相对) | 速度 | 可读性 | 可演进性 |
|---|---|---|---|---|
| JSON | 100% | 中 | 好 | 弱(无编号) |
| Protobuf | 30~50% | 快 | 差(需解码) | 强(字段编号) |
| Avro | 25~40% | 快 | 差 | 强(带 schema) |
| MessagePack | 60~80% | 快 | 差 | 弱 |
| Thrift | 30~50% | 快 | 差 | 强(字段编号) |
| FlatBuffers | 40~60% | 极快(零拷贝) | 差 | 中 |
工作流场景的选型逻辑与微服务略有不同:读写的频率远低于普通 API(一次任务调度读一次),因此「可演进性」与「体积」的权重高于「绝对速度」。这也是为什么 Protobuf 或 Avro 是更合适的选择,而 FlatBuffers 这种追求零拷贝的格式在大多数工作流里属于过度设计。
JSON 并非完全不可用。它的优势是调试友好:实例历史里的数据可以直接在 UI 里读,出问题时不需要工具就能看懂。对于「实例生命周期短、契约稳定、团队规模小」的场景,JSON 的工程收益可能超过它的体积劣势。决策点在于「实例会不会跨版本运行」——会,就必须用带编号的格式。
10. 序列化的性能与体积实测
给一组可用于估算的参考数字(同一份 1 KB 量级的订单对象,1000 次序列化 + 反序列化的总耗时):
格式 序列化耗时 反序列化耗时 体积
JSON 18 ms 26 ms 1024 B
Protobuf 9 ms 12 ms 412 B
Avro 11 ms 14 ms 368 B
MessagePack 13 ms 17 ms 720 B
结论有两条。一是体积差异远大于速度差异:Protobuf 的体积只有 JSON 的 40%,速度只有 2 倍差距。对工作流来说体积更重要,因为它直接决定存储成本与重放时间。二是反序列化普遍比序列化慢,这解释了为什么「读多写少」的工作流里格式选择要更看重解码效率。绝对数值都很小,真正的问题在体积导致的 I/O 与重放上,所以实测重点应该放在真实数据集的压缩后体积上,而不是微基准。
11. 数据在实例状态中的累积
即使单次传递的数据很小,累积起来也会失控。三种累积模式:
纵向累积:一个长驻实例不断接收信号,每次信号都追加数据
例:订单实例累积了 5000 条物流轨迹事件
横向累积:每步都把结果追加到一个列表里,传给下一步
例:每步 append 一条处理记录,50 步后列表有 50 项
重复累积:同一步骤被重试 N 次,每次的结果都被记录
例:一个失败重试 20 次的步骤产生 20 条失败记录
治理手段有三种。一是设置累积上限:任何「追加型」的字段都要有上限,超过时保留最近的 N 条并记录「已截断」标记。二是把累积数据外置:轨迹类数据写入独立的表,实例状态里只保留「最新状态 + 总数」。三是限制重试记录:重试的中间结果不需要全部持久化,只保留最后一次失败的摘要。
// 累积数据外置:实例状态只保留指针与计数
public class OrderState {
private String latestLogUri; // 轨迹明细在外部
private int logCount; // 总条数
private Instant lastEventAt;
public void appendLog(LogEntry entry) {
appendToExternal(latestLogUri, entry); // 追加到外部存储
this.logCount++;
this.lastEventAt = entry.occurredAt();
}
}
长驻实例还必须配合历史截断(Temporal 的 ContinueAsNew),否则事件数会撞上上限。截断时要把「外部存储的指针」作为状态传下去,而不是把数据本身带过去。
12. 敏感数据的识别与分级
工作流状态里经常混入不该持久化的数据:身份证号、银行卡号、健康信息、地址。问题在于开发者通常不会意识到「引擎会持久化一切」,于是把本应只在内存里流转的数据写进了实例状态。
治理的第一步是分级。四级分类够用:
| 级别 | 定义 | 示例 | 处理要求 |
|---|---|---|---|
| L1 公开 | 泄露无影响 | 订单号、商品 ID | 无需特殊处理 |
| L2 内部 | 泄露影响有限 | 内部编码、金额 | 访问控制 |
| L3 敏感 | 涉及个人隐私 | 手机号、地址、身份证 | 加密 + 脱敏 |
| L4 极敏感 | 强监管 | 银行卡、生物特征、密码 | 禁止持久化 |
L4 的原则是不进实例状态。做法是让任务在内存里完成处理,只把「是否通过校验」「掩码后的后四位」这类派生结果写进状态。任何需要跨步骤使用的 L4 数据都应该走专门的安全存储(密钥管理服务、专门的加密表),用一次性令牌引用。
识别手段有三种:代码评审时的字段名扫描(idCard、bankCard、password 这类命名)、自动化的静态扫描(把敏感字段名做成规则库接进 CI)、运行时的数据采样检测(对写入实例状态的载荷做正则匹配)。第三种最可靠但也最重,适合作为兜底。
13. 加密与密钥管理
对于确实需要持久化的 L3 数据,加密是唯一选择。加密的层次有三个,效果与成本递增:
传输加密(TLS):防中间人,不防存储泄露
存储加密(透明加密 / 磁盘加密):防物理介质泄露,不防数据库被拖库
字段级加密(应用层加密):防数据库泄露与运维越权,需要管理密钥
工作流场景至少要做到「传输 + 存储」两层,涉及 L3 数据的要上字段级加密。
// 字段级加密:加密后再放进实例状态,密文长度会膨胀
public class EncryptedField {
private String ciphertext; // Base64(AES-GCM(plaintext))
private String keyId; // 密钥版本,支持轮换
private String iv; // 每次加密随机生成
public static EncryptedField of(String plaintext, KmsClient kms, String keyId) {
byte[] iv = randomBytes(12);
byte[] ct = aesGcmEncrypt(plaintext, kms.dataKey(keyId), iv);
return new EncryptedField(base64(ct), keyId, base64(iv));
}
}
三个关键点。一是 keyId 必须持久化:密钥轮换后老数据仍要用老密钥解密,没有 keyId 就无法解密。二是 IV 每次随机且不可复用,复用 IV 在 GCM 模式下会直接导致密钥泄露。三是密钥要用 KMS 管理而不是配置文件,配置文件里的密钥会随代码泄露,且无法审计谁读取过。
加密的代价是体积膨胀与不可查询:Base64 编码使密文膨胀 33%,且加密后的字段无法被数据库索引或用于查询条件。因此需要查询的字段(比如手机号用于查重)不能简单加密,要用「确定性加密 + 单独索引表」或者「哈希值用于等值查询 + 密文用于展示」的组合。
14. 脱敏、掩码与访问控制
很多场景不需要「可解密的加密」,只需要「不可还原的脱敏」。区分三者的用途:
加密:需要还原(比如后续步骤要发短信给用户)
掩码:只需展示部分(界面上显示 138****8888)
哈希:只需等值比对(用手机号查重、做去重键)
选择标准是「后续步骤需不需要原始值」。需要就加密,不需要就掩码或哈希。把「其实不需要还原」的数据加密,是过度设计,会带来不必要的密钥管理负担。
掩码的实现建议在写入实例状态时做,而不是在展示时做。因为一旦明文进了实例状态,它就存在于数据库备份、日志、快照里,展示层的掩码只是视觉遮挡。正确做法是任务返回时就把值掩码掉:
def mask_phone(phone: str) -> str:
return phone[:3] + "****" + phone[-4:] if len(phone) == 11 else "***"
访问控制是最后一层,原则是最小可见性:实例状态的完整内容只有该流程的运维与开发可见,业务方通过专门的只读接口查看,且接口层做字段过滤。可观测性与数据隐私天然冲突,必须在设计时明确边界,相关讨论见 工作流可观测与调试 。
15. 数据生命周期与清理
工作流数据有三份副本,各有各的清理策略:实例状态、外部中间结果、日志。任何一份没有清理策略,都会变成成本黑洞。
| 数据 | 保留期 | 清理方式 |
|---|---|---|
| 运行中实例的状态 | 与实例生命周期一致 | 实例结束时随实例归档 |
| 已结束实例的历史 | 7~90 天(按合规要求) | 引擎的保留策略自动清理 |
| 外部中间结果 | 1~30 天 | 对象存储生命周期规则 |
| 任务日志 | 7~30 天 | 日志平台的保留策略 |
| 审计记录 | 1~7 年 | 独立归档,不可删除 |
清理的关键约束是**「不能清理还在被引用的数据」**。中间结果的清理必须检查引用:如果实例还在运行且引用了某个文件,就不能删。实现方式是在 flow_artifact 表里记录 expires_at(等于「实例最晚结束时间 + 缓冲」),清理任务只删 expires_at < now() 的记录。
保留期还要与合规要求对齐。金融与医疗场景通常要求「交易数据保留 5 年以上」,但「实例状态」与「业务数据」是两回事:业务数据在业务系统里长期保留,实例状态只需保留到「不可能再需要重放」为止。
16. 跨引擎的传递机制差异
不同引擎的传递机制差异很大,迁移时这部分往往需要重写:
| 引擎 | 传递机制 | 体积限制 | 持久化位置 |
|---|---|---|---|
| Temporal | Activity 参数/返回值 | 2 MB | 事件历史 |
| Airflow | XCom | 元数据库限制 | 元数据库(可外置) |
| Camunda 7 | 流程变量 | 数据库行限制 | 数据库变量表 |
| Zeebe | 变量 | 4 MB | 引擎状态 |
| Dagster | IO Manager | 无(可配) | 可插拔(S3/DB/内存) |
| Step Functions | JSON 状态 | 256 KB | 引擎内部 |
Dagster 的 IO Manager 是最灵活的抽象:它把「任务输出存哪、下游怎么读」抽成了一个可插拔的组件,因此同一个资产可以用内存传递(开发环境)或对象存储传递(生产环境),业务代码不变。数据编排侧的资产与 IO 设计参见 Dagster 与 Prefect 数据编排 。
17. 类型系统与运行时校验
静态类型只能在编译期保护「本服务内的代码」,无法保护「历史数据」。因此运行时的反序列化校验是必须的:
from pydantic import BaseModel, Field, ValidationError
class OrderInput(BaseModel):
order_id: str = Field(min_length=1, max_length=64)
amount_cents: int = Field(ge=0)
channel: str = "unknown" # 默认值保证老数据能解析
model_config = {"extra": "ignore"} # 忽略未知字段,容忍新版本写入的字段
try:
inp = OrderInput.model_validate(raw_state)
except ValidationError as e:
# 关键:区分「数据格式错误」与「代码 bug」,前者走人工处理
raise NonRetryableError(f"invalid state: {e}") from e
三个设计要点。一是 extra = "ignore":新版本写入的字段被老版本代码读到时不应报错,这是回滚兼容的基础。二是校验失败要抛不可重试错误:格式不对的数据重试一百次也不会变对,应该直接转人工或死信。三是区分「缺字段」与「字段值非法」:前者可能是老版本数据(用默认值兜住),后者是真实的数据问题(必须暴露)。
18. 契约测试与数据演练
契约的破坏往往在「生产者改了返回结构、消费者没同步改」时发生。防御手段是契约测试:把契约定义作为测试的输入,同时验证生产者产出的结构与消费者期望的结构一致。
def test_extract_output_matches_contract():
result = extract(data_interval_start=DATETIME, data_interval_end=DATETIME)
# 生产者产出必须符合契约
OrderExtractResult.model_validate(result)
def test_consumer_accepts_old_schema():
# 用三个月前的真实历史样本验证消费者仍能解析
old_sample = json.load(open("fixtures/order_input_v1.json"))
assert OrderInput.model_validate(old_sample).order_id is not None
第二个测试(老数据兼容)比第一个更重要,因为它直接模拟了「老实例 + 新代码」的真实场景。样本库应该从生产环境定期采样并纳入版本控制,覆盖每个历史版本至少一条。
除了契约测试,还建议做一次数据演练:把生产环境的实例状态导出到测试环境,用新代码跑一遍完整的反序列化与重放,确认没有兼容问题。这类演练最好在每次涉及数据结构的发布前做一次。
19. 观测与容量规划
数据传递的成本要能被观测到,否则问题总是被动发现。四个指标:
| 指标 | 含义 | 关注点 |
|---|---|---|
| 单次载荷体积 P99 | 任务输入输出的字节数 | 超过 100 KB 告警 |
| 实例状态体积 P99 | 单个实例的累计状态大小 | 超过 1 MB 告警 |
| 历史事件数 P99 | 单个实例的事件条数 | 超过 5000 告警 |
| 外部中间结果总量 | 对象存储的占用 | 按增长率预测账单 |
| 重放耗时 P99 | 实例重放的时间 | 超过 1 秒告警 |
「单次载荷体积」应该按任务类型分组看,而不是看全局。一个全局 P99 正常但某个任务返回 5 MB 的情况很常见,分组才能定位。
容量规划的估算公式:
日存储增量 = 日实例数 × 平均步骤数 × 平均单步载荷体积 × 2(输入+输出)× 副本数
保留量 = 日存储增量 × 保留天数
用这个公式算一遍,通常会发现「保留 90 天历史」的成本远超预期,从而推动把保留期缩短或把大载荷改为引用传递。容量规划与成本模型的更多讨论参见 工作流引擎全景与选型 。
20. 落地路线图
- 第 1 周:给所有任务的输入输出加上显式类型定义,消灭裸
Map与dict。 - 第 2 周:接入体积监控,找出当前超过 100 KB 的载荷,评估改为引用传递的收益。
- 第 3 周:为大载荷任务实现 Claim Check,建立外部中间结果的命名规范与清理规则。
- 第 4 周:把契约定义接进 Schema 注册中心,开启 FULL 兼容性校验。
- 第 5 周:做一次敏感数据扫描,识别 L3/L4 字段,L4 改为不持久化、L3 上字段级加密。
- 第 6 周:补齐契约测试与老数据兼容测试,纳入 CI。
顺序上「加类型定义」必须最先做,因为没有类型就无法讨论契约、无法做兼容校验、无法做字段级加密。这一步的收益也最直接:类型定义本身就是最好的文档。
21. 权衡取舍
| 选择 | 收益 | 代价 |
|---|---|---|
| 值传递 | 简单、可重放、可观测 | 有体积上限、拖慢重放 |
| 引用传递 | 无上限、状态轻量 | 需要外部存储、清理与失效处理 |
| 共享存储 | 零拷贝、适合大数据 | 隐式依赖、难追溯 |
| JSON | 调试友好、无需工具 | 体积大、无字段编号 |
| Protobuf / Avro | 体积小、可演进 | 需 schema registry 与代码生成 |
| FULL 兼容 | 支持灰度与回滚 | 变更受限,有时需要多版本并行 |
| 字段级加密 | 防数据库泄露 | 体积膨胀、无法索引查询 |
| 掩码替代加密 | 实现简单、体积不增 | 不可还原,需要原始值时无法补救 |
| 累积数据外置 | 实例状态轻量 | 多一次外部读写、需管理引用 |
| 缩短保留期 | 成本低、合规风险小 | 失去长期重放与审计能力 |
22. 常见坑清单
- 任务返回几十万行明细,实例历史膨胀到几十 MB,重放要几秒甚至超时。
- 用
Map<String, Object>做任务参数,字段名与类型全靠约定,改一处崩一片。 - 新增必填字段没有默认值,老实例反序列化失败后卡住。
- 用 JSON 存跨月实例的状态,删除字段后编号被复用,老数据被错误解析。
- 中间结果路径用随机 UUID,重跑产生重复文件,且无法从路径推断来源。
- 写对象存储不做原子重命名,下游读到写了一半的文件。
- 外部中间结果没有生命周期规则,对象存储账单逐月翻倍。
- 长驻实例不断追加轨迹数据,撞上事件数上限或状态体积上限。
- 身份证号、银行卡号直接写进实例状态,存在于备份与日志里无法清除。
- 字段级加密只存密文不存
keyId,密钥轮换后老数据永久无法解密。 - AES-GCM 复用 IV,直接导致密钥泄露。
- 对加密字段建索引或做查询条件,查询结果为空且报错难以理解。
- 校验失败的数据抛可重试错误,重试一百次后进死信,浪费大量调度资源。
- 契约只写文档不做类型定义,生产端改了返回结构,消费端静默出错。
23. 小结
工作流的数据传递有一条贯穿始终的原则:实例状态里只放「参与决策的控制信息」,业务数据放外部存储并用引用关联。这条原则同时解决了体积、成本、重放性能与数据治理四个问题,是投入产出比最高的设计决策。
Schema 演进的核心规则是「只加可选字段、不删不改、用带编号的格式」。这条规则看起来保守,但它是唯一能同时支持灰度发布与代码回滚的方案。配合 Schema 注册中心的 FULL 兼容性校验,可以把绝大多数契约破坏挡在发布之前。
数据安全上要先分级再处理:L4 不持久化、L3 加密或掩码、L2 做访问控制。分级的意义在于避免「一刀切加密」带来的密钥管理与查询能力损失。可观测性、成本与数据治理是同一件事的三个侧面,设计数据传递方案时应该一起考虑,而不是等账单或审计来提醒。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。