工作流数据传递与 Schema

本文系统讲解工作流的任务间数据传递与 Schema 演进,回答大数据集该传值还是传引用、Schema 变更怎样不破坏运行中实例、敏感字段如何加密与脱敏。覆盖值传递与 Claim Check 模式、中间结果存储选型、Schema 注册中心与兼容规则、序列化格式的体积与性能对比、事件历史累积治理、字段分级与密钥管理、生命周期清理与契约测试,并给出可落地的代码与配置片段。

引言

大多数人第一次意识到「数据传递」是个问题时,工作流已经在生产上跑了几周。表现是实例列表页加载越来越慢、引擎数据库的磁盘占用线性增长、某个跑了一个月的实例重放要花 8 秒。排查后发现根源是某一步任务返回了 20 万行明细,而引擎把每一步的输入输出都持久化了。

工作流引擎与普通服务在数据传递上有一个根本差异:引擎会持久化每一步的输入与输出,因为它们既是重放的依据,也是断点恢复的凭据。这个设计带来了可观测性与可靠性,也带来了两个硬约束:单次传递的数据有体积上限,累计传递的数据会无限增长。

第二个差异是数据的生命周期远长于代码。一个跨月实例里保存的 JSON 是三个月前写的代码产生的,期间的 Schema 变更必须保证它仍能正确反序列化。这与微服务之间的 API 契约不同——API 调用是即时的,工作流的状态是持久的。

本文按「传什么 → 怎么传 → 怎么演进 → 怎么保护」的顺序展开:先讲值传递与引用传递的边界(Claim Check 模式),再讲 Schema 定义与兼容规则,然后是序列化格式的选型与实测数据,最后是敏感数据的加密、脱敏与生命周期治理。引擎的序列化实现差异参见 Temporal 与持久化执行 ,数据管道的编排视角参见 Airflow DAG 调度体系 。

目录

  1. 数据传递为什么是隐藏难点
  2. 三种传递方式:值、引用与共享存储
  3. 数据契约的设计原则
  4. 值传递的边界与体积上限
  5. Claim Check 模式
  6. 中间结果的存放位置选择
  7. Schema 定义与注册中心
  8. Schema 演进的兼容规则
  9. 序列化格式选型
  10. 序列化的性能与体积实测
  11. 数据在实例状态中的累积
  12. 敏感数据的识别与分级
  13. 加密与密钥管理
  14. 脱敏、掩码与访问控制
  15. 数据生命周期与清理
  16. 跨引擎的传递机制差异
  17. 类型系统与运行时校验
  18. 契约测试与数据演练
  19. 观测与容量规划
  20. 落地路线图
  21. 权衡取舍
  22. 常见坑清单
  23. 小结

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. 值传递的边界与体积上限

值传递有明确的硬上限,且不同引擎的默认值差异很大,迁移时必须重新评估:

引擎单载荷上限历史/状态上限超限行为
Temporal2 MB(可配置)50 MB 或 50K 事件调度失败 / 强制截断
Step Functions256 KB无(状态在外部)执行失败
Camunda 7无硬限制(存数据库)受数据库行大小限制数据库写入失败
Airflow XCom元数据库字段限制元数据库容量性能急剧下降
Zeebe4 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 演进的兼容规则

兼容性级别落到具体操作上,就是一张「能不能做」的表:

变更BACKWARDFORWARDFULL
新增可选字段(带默认值)可以可以可以
新增必填字段不可以可以不可以
删除字段可以不可以不可以
字段改类型不可以不可以不可以
字段改名不可以不可以不可以
扩大枚举范围可以不可以不可以
缩小枚举范围不可以可以不可以
提高字段精度(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. 序列化格式选型

序列化格式决定体积、速度与可演进性,三者往往互相冲突:

格式体积(相对)速度可读性可演进性
JSON100%中好弱(无编号)
Protobuf30~50%快差(需解码)强(字段编号)
Avro25~40%快差强(带 schema)
MessagePack60~80%快差弱
Thrift30~50%快差强(字段编号)
FlatBuffers40~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. 跨引擎的传递机制差异

不同引擎的传递机制差异很大,迁移时这部分往往需要重写:

引擎传递机制体积限制持久化位置
TemporalActivity 参数/返回值2 MB事件历史
AirflowXCom元数据库限制元数据库(可外置)
Camunda 7流程变量数据库行限制数据库变量表
Zeebe变量4 MB引擎状态
DagsterIO Manager无(可配)可插拔(S3/DB/内存)
Step FunctionsJSON 状态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. 常见坑清单

  1. 任务返回几十万行明细,实例历史膨胀到几十 MB,重放要几秒甚至超时。
  2. 用 Map<String, Object> 做任务参数,字段名与类型全靠约定,改一处崩一片。
  3. 新增必填字段没有默认值,老实例反序列化失败后卡住。
  4. 用 JSON 存跨月实例的状态,删除字段后编号被复用,老数据被错误解析。
  5. 中间结果路径用随机 UUID,重跑产生重复文件,且无法从路径推断来源。
  6. 写对象存储不做原子重命名,下游读到写了一半的文件。
  7. 外部中间结果没有生命周期规则,对象存储账单逐月翻倍。
  8. 长驻实例不断追加轨迹数据,撞上事件数上限或状态体积上限。
  9. 身份证号、银行卡号直接写进实例状态,存在于备份与日志里无法清除。
  10. 字段级加密只存密文不存 keyId,密钥轮换后老数据永久无法解密。
  11. AES-GCM 复用 IV,直接导致密钥泄露。
  12. 对加密字段建索引或做查询条件,查询结果为空且报错难以理解。
  13. 校验失败的数据抛可重试错误,重试一百次后进死信,浪费大量调度资源。
  14. 契约只写文档不做类型定义,生产端改了返回结构,消费端静默出错。

23. 小结

工作流的数据传递有一条贯穿始终的原则:实例状态里只放「参与决策的控制信息」,业务数据放外部存储并用引用关联。这条原则同时解决了体积、成本、重放性能与数据治理四个问题,是投入产出比最高的设计决策。

Schema 演进的核心规则是「只加可选字段、不删不改、用带编号的格式」。这条规则看起来保守,但它是唯一能同时支持灰度发布与代码回滚的方案。配合 Schema 注册中心的 FULL 兼容性校验,可以把绝大多数契约破坏挡在发布之前。

数据安全上要先分级再处理:L4 不持久化、L3 加密或掩码、L2 做访问控制。分级的意义在于避免「一刀切加密」带来的密钥管理与查询能力损失。可观测性、成本与数据治理是同一件事的三个侧面,设计数据传递方案时应该一起考虑,而不是等账单或审计来提醒。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「工作流引擎」更多文章

  1. 工作流成本优化
  2. 执行器与资源隔离
  3. 调度、回填与补数