一、分布式锁概述
分布式锁是分布式系统中控制多个节点访问共享资源的机制,在数据密集型应用中,分布式锁是保障数据一致性和避免并发冲突的关键。
二、分布式锁特性
2.1 分布式锁特性
| 特性 | 描述 | 重要性 | 实现难点 |
|---|---|---|---|
| 互斥性 | 同一时间只有一个客户端持有锁 | 高 | 多节点竞争 |
| 锁超时 | 锁自动释放防止死锁 | 高 | 超时时间设置 |
| 可重入性 | 同一客户端可重复获取锁 | 中 | 锁计数管理 |
| 公平性 | 按顺序获取锁 | 中 | 队列管理 |
| 高性能 | 获取和释放锁快速 | 高 | 减少网络开销 |
| 高可用 | 锁服务不可用时降级 | 高 | 故障转移 |
三、Redis分布式锁实现
3.1 Redis分布式锁基础实现
public class RedisDistributedLock : IDistributedLock
{
private readonly IDistributedCache _cache;
private readonly string _lockKey;
private readonly string _lockValue;
private readonly TimeSpan _expiry;
public RedisDistributedLock(IDistributedCache cache, string lockKey, TimeSpan expiry)
{
_cache = cache;
_lockKey = lockKey;
_lockValue = Guid.NewGuid().ToString();
_expiry = expiry;
}
public async Task AcquireAsync()
{
var result = await _cache.StringSetAsync(
_lockKey,
_lockValue,
_expiry,
When.NotExists);
return result;
}
public async Task ReleaseAsync()
{
var luaScript = @"
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
";
var result = await _cache.ExecuteAsync(luaScript,
new RedisKey[] { _lockKey },
new RedisValue[] { _lockValue });
return (long)result == 1;
}
public async Task RenewAsync()
{
var luaScript = @"
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('expire', KEYS[1], ARGV[2])
else
return 0
end
";
var result = await _cache.ExecuteAsync(luaScript,
new RedisKey[] { _lockKey },
new RedisValue[] { _lockValue, (int)_expiry.TotalSeconds });
return (long)result == 1;
}
public async Task ExecuteWithLockAsync(Func> func)
{
if (!await AcquireAsync())
{
throw new LockAcquisitionException("Failed to acquire lock");
}
try
{
return await func();
}
finally
{
await ReleaseAsync();
}
}
}
3.2 Redis分布式锁优化
public class RedisDistributedLockOptimized : IDistributedLock
{
private readonly IDistributedCache _cache;
private readonly string _lockKey;
private readonly string _lockValue;
private readonly TimeSpan _expiry;
private CancellationTokenSource _renewalTokenSource;
public async Task AcquireAsync(int retryCount = 3, TimeSpan retryDelay = default)
{
retryDelay = retryDelay == default ? TimeSpan.FromMilliseconds(100) : retryDelay;
for (int i = 0; i < retryCount; i++)
{
if (await AcquireOnceAsync())
{
StartRenewal();
return true;
}
if (i < retryCount - 1)
{
await Task.Delay(retryDelay);
}
}
return false;
}
private async Task AcquireOnceAsync()
{
var result = await _cache.StringSetAsync(
_lockKey,
_lockValue,
_expiry,
When.NotExists);
return result;
}
private void StartRenewal()
{
_renewalTokenSource = new CancellationTokenSource();
_ = Task.Run(async () =>
{
while (!_renewalTokenSource.Token.IsCancellationRequested)
{
await Task.Delay(_expiry / 3);
await RenewAsync();
}
});
}
public async Task ReleaseAsync()
{
_renewalTokenSource?.Cancel();
var luaScript = @"
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
";
var result = await _cache.ExecuteAsync(luaScript,
new RedisKey[] { _lockKey },
new RedisValue[] { _lockValue });
return (long)result == 1;
}
}
四、ZooKeeper分布式锁实现
4.1 ZooKeeper分布式锁原理
graph TD
A[ZooKeeper锁节点] --> B[/locks/]
B --> C[lock-000000001]
B --> D[lock-000000002]
B --> E[lock-000000003]
F[客户端1] --> G[创建临时顺序节点]
G --> H[获取最小节点]
H --> I{是否最小?}
I -->|是| J[获取锁]
I -->|否| K[监听前一个节点]
K --> L[等待通知]
L --> H
M[客户端2] --> G
4.2 ZooKeeper分布式锁实现
public class ZooKeeperDistributedLock : IDistributedLock
{
private readonly IZooKeeperClient _zooKeeper;
private readonly string _lockPath;
private readonly string _nodePath;
private EventWaitHandle _waitHandle;
public ZooKeeperDistributedLock(IZooKeeperClient zooKeeper, string lockPath)
{
_zooKeeper = zooKeeper;
_lockPath = lockPath;
}
public async Task AcquireAsync()
{
_nodePath = await _zooKeeper.CreateAsync(
$"{_lockPath}/lock-",
data: null,
ZooDefs.Ids.OPEN_ACL_UNSAFE,
CreateMode.EPHEMERAL_SEQUENTIAL);
var lockNodeName = _nodePath.Substring(_lockPath.Length + 1);
var children = await _zooKeeper.GetChildrenAsync(_lockPath);
var sortedChildren = children.OrderBy(c => c).ToList();
var index = sortedChildren.IndexOf(lockNodeName);
if (index == 0)
{
return true;
}
var previousNode = $"{_lockPath}/{sortedChildren[index - 1]}";
_waitHandle = new AutoResetEvent(false);
await _zooKeeper.SubscribeAsync(previousNode, async (watcher) =>
{
if (watcher.Type == Watcher.Event.EventType.NodeDeleted)
{
_waitHandle.Set();
}
});
_waitHandle.WaitOne();
return true;
}
public async Task ReleaseAsync()
{
_waitHandle?.Dispose();
await _zooKeeper.DeleteAsync(_nodePath);
return true;
}
}
五、etcd分布式锁实现
5.1 etcd分布式锁原理
public class EtcdDistributedLock : IDistributedLock
{
private readonly IEtcdClient _etcdClient;
private readonly string _lockKey;
private readonly string _leaseId;
private readonly TimeSpan _expiry;
public EtcdDistributedLock(IEtcdClient etcdClient, string lockKey, TimeSpan expiry)
{
_etcdClient = etcdClient;
_lockKey = lockKey;
_expiry = expiry;
}
public async Task AcquireAsync()
{
var leaseResponse = await _etcdClient.LeaseGrantAsync((int)_expiry.TotalSeconds);
_leaseId = leaseResponse.ID.ToString();
var txnResponse = await _etcdClient.TxnAsync(txn => txn
.If(Compare.Version(_lockKey).Equal(0))
.Then(Op.Put(_lockKey, _leaseId, leaseResponse.ID))
.Else(Op.Get(_lockKey)));
if (!txnResponse.Succeeded)
{
await WaitForLockAsync();
return true;
}
StartLeaseKeepAlive();
return true;
}
private async Task WaitForLockAsync()
{
var watcher = _etcdClient.Watch(_lockKey);
await foreach (var response in watcher)
{
foreach (var @event in response.Events)
{
if (@event.Type == EventType.Delete)
{
await AcquireAsync();
return;
}
}
}
}
private void StartLeaseKeepAlive()
{
_ = Task.Run(async () =>
{
var keepAlive = await _etcdClient.LeaseKeepAliveAsync(_leaseId);
await foreach (var response in keepAlive)
{
if (response == null)
{
break;
}
}
});
}
public async Task ReleaseAsync()
{
await _etcdClient.DeleteAsync(_lockKey);
await _etcdClient.LeaseRevokeAsync(_leaseId);
return true;
}
}
六、分布式锁方案对比
6.1 分布式锁方案对比表
| 方案 | 性能 | 可靠性 | 公平性 | 实现复杂度 | 适用场景 |
|---|---|---|---|---|---|
| Redis | 极高 | 高 | 低 | 低 | 高并发 |
| ZooKeeper | 中等 | 极高 | 高 | 中等 | 强一致性 |
| etcd | 高 | 极高 | 高 | 中等 | K8s环境 |
七、分布式锁最佳实践
7.1 锁超时策略
public class DistributedLockBestPractices
{
public TimeSpan CalculateLockTimeout(TimeSpan estimatedOperationTime)
{
return estimatedOperationTime + TimeSpan.FromSeconds(10);
}
public async Task ExecuteWithTimeoutAsync(Func> func, TimeSpan timeout)
{
using var cts = new CancellationTokenSource(timeout);
var task = func();
var completedTask = await Task.WhenAny(task, Task.Delay(timeout));
if (completedTask == task)
{
return await task;
}
throw new TimeoutException("Operation timed out");
}
}
7.2 锁重入实现
public class ReentrantDistributedLock : IDistributedLock
{
private readonly IDistributedLock _innerLock;
private int _lockCount;
public async Task AcquireAsync()
{
if (_lockCount > 0)
{
_lockCount++;
return true;
}
var acquired = await _innerLock.AcquireAsync();
if (acquired)
{
_lockCount = 1;
}
return acquired;
}
public async Task ReleaseAsync()
{
if (_lockCount > 1)
{
_lockCount--;
return true;
}
_lockCount = 0;
return await _innerLock.ReleaseAsync();
}
}
7.3 锁降级策略
public class LockDegradationService
{
public async Task ExecuteWithDegradationAsync(
Func> lockedOperation,
Func> degradedOperation)
{
try
{
return await lockedOperation();
}
catch (LockAcquisitionException)
{
return await degradedOperation();
}
catch (Exception)
{
return await degradedOperation();
}
}
public async Task ExecuteWithFallbackAsync(
Func> primaryOperation,
Func> fallbackOperation,
int maxRetries = 3)
{
for (int i = 0; i < maxRetries; i++)
{
try
{
return await primaryOperation();
}
catch (Exception)
{
if (i == maxRetries - 1)
{
return await fallbackOperation();
}
await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, i)));
}
}
return await fallbackOperation();
}
}
7.4 锁监控与告警
public class DistributedLockMonitor
{
public async Task GetMetricsAsync()
{
var metrics = new LockMetrics();
var locks = await _lockRepository.GetAllLocksAsync();
foreach (var @lock in locks)
{
metrics.TotalLocks++;
if (@lock.IsAcquired)
{
metrics.AcquiredLocks++;
if (@lock.AcquiredAt.Add(@lock.Expiry) < DateTime.UtcNow)
{
metrics.ExpiredLocks++;
}
}
metrics.TotalWaitTime += @lock.WaitTime;
}
metrics.AverageWaitTime = locks.Count > 0
? metrics.TotalWaitTime / locks.Count
: TimeSpan.Zero;
return metrics;
}
public async Task MonitorAsync()
{
var metrics = await GetMetricsAsync();
if (metrics.ExpiredLocks > 0)
{
await _alertService.SendAlert("分布式锁过期",
$"过期锁数: {metrics.ExpiredLocks}");
}
if (metrics.AverageWaitTime > TimeSpan.FromSeconds(5))
{
await _alertService.SendAlert("分布式锁等待过长",
$"平均等待时间: {metrics.AverageWaitTime}");
}
}
}
public class LockMetrics
{
public int TotalLocks { get; set; }
public int AcquiredLocks { get; set; }
public int ExpiredLocks { get; set; }
public TimeSpan TotalWaitTime { get; set; }
public TimeSpan AverageWaitTime { get; set; }
}
八、总结
分布式锁是分布式系统中控制多个节点访问共享资源的关键机制。Redis适合高并发场景,ZooKeeper适合强一致性场景,etcd适合Kubernetes环境。通过合理选择锁方案、实现锁超时和重入机制、做好监控和告警,能够构建可靠的分布式锁系统。