一、我们为什么要聊这个

做过分布式系统开发的朋友,应该都遇到过这样一种情况:多个客户端同时往同一个地方写数据,结果后写的人把前一个的正常修改给覆盖了,数据莫名其妙丢失。放在 Durable Entities 这种长期运行、有状态的计算单元里,问题就更刺手——因为一个实体可能同时被好几个请求调用,每个请求都以为自己拿到了“最新”的数据,然后各自修改再保存,最后只有一个生效,其他的全白干。

你可能会想,加个全局锁不就行了?可分布式环境下的锁要么太慢,要么容易死锁,要么干脆不支持。业界普遍采用的方案是“乐观锁”,配合实体本身的版本控制机制,让冲突在提交阶段暴露出来,由业务侧决定是重试、合并还是报错。今天咱们就用最接地气的方式,把这种玩法聊透。

二、典型的冲突场景:银行转账&购物车

假设你有一个网上书店,用户可以把书加入购物车。购物车存成一个 Durable Entity,里面有一个 items 列表。用户 A 在手机上加了一本书,与此同时用户 B 在电脑上把购物车里的一本书删掉了。两个操作几乎是同时发生的,实体收到请求后,先读取当前的 items,A 操作在列表末尾追加,B 操作从列表中间移除,然后各自保存。如果实体没有做任何保护,最后保存的那个会把另一个的改动冲掉——要么书没加进去,要么书没删成功。

更常见的例子是银行账户扣款。假设账户余额是 100 元,两笔同时发起的转账分别扣 80 元和 30 元。如果互不干扰地读取余额→计算新余额→写入,最终余额可能变成 70(第二笔基于原始 100 算出的 70),而实际上应该是 -10(先扣 80 得 20,再扣 30 得 -10)或者需要拒绝第二笔。这种财务数据丢失是绝对不能接受的。

三、乐观锁的核心思想:不抢,先检查再落笔

乐观锁的想法很简单:我不在读写之间加锁,而是允许所有人同时读,但只在最后写的时候检查一下——数据在我读完之后有没有被别人改过?如果被别人改了,我就放弃本次写入,重新读取最新的再去计算,或者直接抛异常告诉调用方“你该重试”。

在 Durable Entities 里,这个“检查别人有没有改过”靠的就是一个版本号。每次实体状态改变,版本号就加一。读取的时候把版本号也读出来,写入的时候带上这个旧版本号,告诉存储:只有当你当前版本号等于我读到的这个版本号时,才允许写入。如果不等于,说明有别人抢先改了,写入失败。

四、实体版本控制机制怎么玩

4.1 版本号到底是什么

简单说,版本号就是一个整数字段,放在实体的状态里面。初始值是 0(或者 1),每次 SignalEntityCallEntity 触发状态修改,实体在保存之前把版本号增加到下一个值。

4.2 在 Durable Entities 中实现的两种方式

方式一:手动在实体代码里维护版本号

这是最灵活的方式,你可以自己控制检查逻辑。实体状态类里加一个 Version 属性,在操作入口判断传入的版本号是否等于当前版本号,如果不一致就抛出 VersionConflictException 或返回一个特殊的错误结果。

方式二:利用底层存储的 ETag(实体标签)

对于 Azure 存储来说,Durable Entities 背后是 Table Storage 或 Blob,默认就支持 ETag。每次读取实体状态都会带一个 ETag 值,写入时如果指定的 ETag 与当前不符,存储会返回 412 Precondition Failed。Durable Functions 运行时其实已经内置了这种冲突检测,只不过默认情况下它帮你自动重试(最多 10 次),你几乎感觉不到。但如果你想自己控制冲突时的行为,就需要自己管理 ETag。

本文重点讲第一种手动方式,因为它更透明,方便你理解原理。

五、完整示例:购物车实体加版本锁

下面的例子使用 C# (.NET 6)Microsoft.Azure.WebJobs.Extensions.DurableTask 包。代码里包含详细的注释,让你一步步看到版本号是怎么流动的。

// 1. 定义实体状态
public class ShoppingCartState
{
    public List<string> Items { get; set; } = new List<string>();
    public int Version { get; set; } = 0;  // 版本号,每次修改+1
}

// 2. 定义实体操作输入(用于携带调用方期望的版本号)
public class CartOperation
{
    public string Action { get; set; }           // "Add" 或 "Remove"
    public string Item { get; set; }
    public int ExpectedVersion { get; set; }     // 调用者上次读到的版本号
}

// 3. 定义实体函数本身
[JsonObject(MemberSerialization.OptIn)]
public class ShoppingCartEntity : IDurableEntity
{
    [JsonProperty("state")]
    private ShoppingCartState state = new ShoppingCartState();

    // 实体入口,处理所有操作
    public void Dispatch(IDurableEntityContext ctx)
    {
        // 根据操作类型分发
        switch (ctx.OperationName)
        {
            case "AddItem":
                var addOp = ctx.GetInput<CartOperation>();
                AddItem(addOp);
                break;
            case "RemoveItem":
                var removeOp = ctx.GetInput<CartOperation>();
                RemoveItem(removeOp);
                break;
            case "GetState":
                // 查询操作不修改版本号,直接返回
                ctx.Return(state);
                break;
        }
    }

    private void AddItem(CartOperation op)
    {
        // 核心:版本校验
        if (op.ExpectedVersion != state.Version)
        {
            // 版本不匹配,抛出异常,调用方可以捕获后重试
            throw new InvalidOperationException(
                $"版本冲突:当前版本 {state.Version},期望版本 {op.ExpectedVersion}");
        }

        // 校验通过,执行逻辑
        state.Items.Add(op.Item);
        state.Version++;  // 版本号递增
    }

    private void RemoveItem(CartOperation op)
    {
        if (op.ExpectedVersion != state.Version)
        {
            throw new InvalidOperationException(
                $"版本冲突:当前版本 {state.Version},期望版本 {op.ExpectedVersion}");
        }

        state.Items.Remove(op.Item);
        state.Version++;
    }
}

// 4. 客户端调用示例(比如在 Orchestrator 或 HTTP 触发器中)
public static async Task CallShoppingCart([DurableClient] IDurableEntityClient client)
{
    string entityId = new EntityId(nameof(ShoppingCartEntity), "cart_user123");

    // 先读取当前状态,拿到版本号
    var currentState = await client.ReadEntityStateAsync<ShoppingCartState>(entityId);
    int currentVersion = currentState.EntityState?.Version ?? 0;

    // 构造操作,带上期望版本号
    var operation = new CartOperation
    {
        Action = "Add",
        Item = "《深入理解C#》",
        ExpectedVersion = currentVersion
    };

    // 发送信号(异步)或直接调用
    await client.SignalEntityAsync(entityId, "AddItem", operation);
    // 注意:Signal 是异步的,不会马上拿到异常。
    // 如果要捕获版本冲突,应该改用 CallEntityAsync(同步等待结果)
    // 或通过错误处理机制(比如重试策略)来兜底
}

上面的示例展示了手动版本控制的雏形。实际生产环境里,你可能会把重试逻辑封装成一个辅助方法,让调用方自动处理版本冲突。下面是一个更完整的客户端重试示例(放在 Orchestrator 函数里):

// 5. 在 Orchestrator 中带重试的调用
[FunctionName("PlaceOrderOrchestrator")]
public static async Task RunOrchestrator(
    [OrchestrationTrigger] IDurableOrchestrationContext context)
{
    var entityId = new EntityId(nameof(ShoppingCartEntity), "cart_user123");
    const int maxRetries = 5;
    int attempt = 0;

    while (attempt < maxRetries)
    {
        // 每次重试都重新读取最新状态
        var currentState = await context.CallEntityAsync<ShoppingCartState>(
            entityId, "GetState");
        int version = currentState.Version;

        var op = new CartOperation
        {
            Action = "Add",
            Item = "《设计模式》",
            ExpectedVersion = version
        };

        try
        {
            // 调用实体,如果版本冲突实体内部会抛异常
            await context.CallEntityAsync(entityId, "AddItem", op);
            break; // 成功,退出循环
        }
        catch (InvalidOperationException ex) when (ex.Message.Contains("版本冲突"))
        {
            attempt++;
            if (attempt >= maxRetries)
            {
                // 重试耗尽,记录日志或抛给上游
                throw;
            }
            // 延迟一小段时间再重试(避免忙等)
            await context.CreateTimer(context.CurrentUtcDateTime.AddSeconds(1), CancellationToken.None);
        }
    }
}

这个重试逻辑看起来有点笨,但它是保障数据一致性最可靠的办法。Durable Functions 的编排器天生支持重试和定时器,所以写起来很自然。

六、技术优缺点分析

优点:

  1. 无锁高并发:读操作完全不阻塞,写操作只在最后一刻检查,系统吞吐量远高于悲观锁。
  2. 实现简单:在实体状态里加一个整数,比较一下就行,不需要引入外部锁服务(如 Redis/ ZooKeeper)。
  3. 天然适合分布式:每个实体自己管理自己的版本号,不依赖全局协调,非常适合 Durable Entities 这种分片存储模型。
  4. 错误可预见:冲突发生时调用方明确知道失败了,可以自由选择重试、合并或报错,而不是静默丢失数据。

缺点:

  1. 写冲突频繁时重试开销大:如果几个请求总是同时到达同一个实体,乐观锁会导致大量回滚和重试,系统延迟增加。这时候可以考虑降低重试间隔或改用悲观锁。
  2. 需要调用方配合:调用方必须负责读取版本号、捕获异常、重试,业务代码会变得更啰嗦。好在 Durable Functions 的编排器可以把重试逻辑复用起来。
  3. 长事务不友好:如果调用方读取版本号之后,经过很长时间才发起写入,中间被别的请求修改的概率大大增加,重试成本也高。

七、注意事项

  1. 不要把版本号暴露给外部客户端:版本号应该由服务端内部维护,客户端只管携带从服务端获取到的版本号。如果客户端自己随便填个数字,就失去了校验意义。
  2. 实体操作要幂等:即使加入了版本控制,你也要确保重试不会产生副作用。比如“添加商品”操作如果重复执行,应该在状态里判断商品是否已经存在,避免重复添加。
  3. 版本号溢出:整数版本号理论上有上限(2^31-1),但日常应用几乎不会达到。如果你实在不放心,可以用 long 或 GUID 替代数字版本号。
  4. Durable Entities 默认有内置重试:Azure 的 Durable Functions 底层在遇到存储 ETag 冲突时会自动重试(最多 10 次)。如果你用内置的分布式实体,可能不需要手动实现乐观锁。但手动控制可以在冲突后执行更复杂的逻辑(比如合并或告警)。
  5. 考虑死锁防范:乐观锁虽然避免了死锁,但如果两个操作互相依赖(A 等待 B,B 等待 A),依然可能产生逻辑死锁。这种场景要避免在实体内部调用其他实体。
  6. 测试环境要模拟高并发:本地开发时用单线程很难触发冲突,建议在单元测试中使用多线程或集成测试里用 Task.WhenAll 模拟并发,确保你的版本控制逻辑正确。

八、文章总结

乐观锁加上实体版本控制是处理 Durable Entities 并发更新冲突最实用的方案,它轻量、无锁、易扩展,尤其适合读多写少、冲突概率低的场景。你只需要在实体状态里藏一个版本号,每次修改前做一次比较,然后让调用方接收失败信号并重试,就能大幅降低数据覆盖的风险。

当然它也不是万能的,如果写冲突特别频繁,或者业务上不允许重试(比如金融交易需要实时失败),你可能需要搭配更严格的悲观锁策略。但大多数互联网业务,乐观锁足够用。记住:让冲突早暴露、早重试,远比悄无声息地丢失数据要好。