Datastore事务未锁定:如何实现分布式锁并触发并发事务失败?
问题分析与解决方案
核心问题:事务内使用非事务性写入
你代码里的关键错误是在事务函数内部调用了client.Put而非事务对象tx的Put方法:
client.Put是非事务性写入,会直接绕过当前事务的乐观并发控制机制,独立修改实体。- Datastore的乐观并发冲突检测仅针对事务间的版本不一致,非事务写入不会被纳入事务的冲突检测逻辑,这就是为什么你看到的失败率极低。
实现并发读写强制冲突的正确方式
要让并发操作触发datastore: concurrent transaction错误,必须严格遵循以下规则:
- 所有对目标实体的读写操作必须在事务内部完成,使用事务对象的
Get和Put方法。 - 启动多个并发执行的事务,每个事务都执行「读取实体 → 修改实体 → 提交事务」的逻辑。
- Datastore的乐观并发机制会在事务提交时校验:实体当前版本是否与事务读取时的版本一致,若不一致则抛出并发错误。
修正后的示例代码
以下是模拟高并发场景的代码,会触发明显的事务冲突:
package main import ( "context" "fmt" "sync" "cloud.google.com/go/datastore" ) type Counter struct { Count int } func main() { ctx := context.Background() client, err := datastore.NewClientWithDatabase(ctx, "project-id", "database-id") if err != nil { panic(err) } defer client.Close() key := datastore.NameKey("Counter", "singleton", nil) // 初始化计数器 initCounter(ctx, client, key) var wg sync.WaitGroup concurrency := 10 wg.Add(concurrency) for i := 0; i < concurrency; i++ { go func(id int) { defer wg.Done() var count int // 禁用客户端自动重试,确保冲突错误直接抛出 txOpts := datastore.TransactionOptions{ Attempts: 1, // 只尝试一次,不重试 } if _, err = client.RunInTransaction(ctx, func(tx *datastore.Transaction) error { var x Counter if err := tx.Get(key, &x); err != nil && err != datastore.ErrNoSuchEntity { return err } // 在事务内修改并写入 x.Count++ if _, err = tx.Put(key, &x); err != nil { return err } count = x.Count return nil }, &txOpts); err != nil { fmt.Printf("Goroutine %d 事务失败: %v\n", id, err) return } fmt.Printf("Goroutine %d 成功,Count: %d\n", id, count) }(i) } wg.Wait() // 清理实体 if err = client.Delete(ctx, key); err != nil { panic(err) } } func initCounter(ctx context.Context, client *datastore.Client, key *datastore.Key) { var x Counter if err := client.Get(ctx, key, &x); err == datastore.ErrNoSuchEntity { x.Count = 0 if _, err := client.Put(ctx, key, &x); err != nil { panic(err) } } else if err != nil { panic(err) } }
关于分布式锁的实现
如果你的目标是用Datastore实现分布式锁,有两种常见方案:
- 乐观锁(推荐):就是上面的事务方式,依赖Datastore的版本冲突检测,冲突时重试事务即可。适合读多写少、冲突概率低的场景。
- 模拟悲观锁:通过在实体中添加
Locked布尔字段,事务中先读取实体,检查Locked是否为false,若为false则设置为true并提交事务;操作完成后再将Locked设为false。这种方式需要注意锁超时,避免死锁。
内容的提问来源于stack exchange,提问作者M. Becerra
相关产品推荐
相关产品推荐

