多实例Azure Web API迁移SQL大字段至Azure Blob的EF Core并发方案咨询
多实例环境下EF Core批量迁移数据至Blob的并发解决方案
核心思路
采用乐观锁+原子化抢占的方案,避免长时间锁表导致的超时问题,同时确保各实例处理互不重叠的批次数据,适配上传Blob的耗时差异。
具体实现步骤
1. 优化表结构
在目标业务表中新增两个字段:
ProcessingInstanceId:NVARCHAR(50),标记当前处理该记录的实例唯一标识ProcessingStartTime:DATETIME,记录实例开始处理的时间
2. 原子化抢占未处理记录
通过EF Core结合原生SQL完成短事务的抢占操作,确保同一批记录只会被一个实例获取:
// 获取当前实例的唯一标识(Azure Web App可直接用环境变量,自托管实例生成GUID并缓存) var instanceId = Environment.GetEnvironmentVariable("WEBSITE_INSTANCE_ID") ?? Guid.NewGuid().ToString(); var batchSize = 1000; var timeoutMinutes = 5; // 处理超时阈值,超过则允许其他实例重试 using var transaction = await _dbContext.Database.BeginTransactionAsync(IsolationLevel.ReadCommitted); try { // 原子更新:抢占未处理、未被其他实例锁定或锁定超时的记录 var affectedRows = await _dbContext.Database.ExecuteSqlRawAsync(@" UPDATE TOP({0}) YourTableName SET ProcessingInstanceId = {1}, ProcessingStartTime = GETDATE() WHERE MarkedAsSent = 0 AND (ProcessingInstanceId IS NULL OR ProcessingStartTime < DATEADD(MINUTE, -{2}, GETDATE())) ", batchSize, instanceId, timeoutMinutes); // 查询当前实例抢占到的记录 var recordsToProcess = await _dbContext.YourTable .Where(r => r.ProcessingInstanceId == instanceId && r.MarkedAsSent == false) .ToListAsync(); await transaction.CommitAsync(); // 独立处理抢占到的记录,不持有长事务 await ProcessAndUpdateRecords(recordsToProcess); } catch { await transaction.RollbackAsync(); throw; }
3. 记录处理与状态更新
处理阶段不持有数据库事务,避免因Blob上传耗时过长导致超时:
private async Task ProcessAndUpdateRecords(List<YourTableEntity> records) { var blobServiceClient = new BlobServiceClient("your-blob-connection-string"); var containerClient = blobServiceClient.GetBlobContainerClient("your-container-name"); foreach (var record in records) { try { // 上传大字符串至Blob var blobClient = containerClient.GetBlobClient($"record-{record.Id}.txt"); await blobClient.UploadAsync(new MemoryStream(Encoding.UTF8.GetBytes(record.LargeStringColumn)), overwrite: true); // 更新状态:标记为已迁移,清空抢占标记 record.MarkedAsSent = true; record.ProcessingInstanceId = null; record.ProcessingStartTime = null; record.BlobUrl = blobClient.Uri.ToString(); // 可选:保存Blob访问地址 } catch (Exception ex) { // 处理失败:释放抢占标记,允许其他实例重试 record.ProcessingInstanceId = null; record.ProcessingStartTime = null; // 可选:记录错误日志用于排查 } } // 批量提交状态更新 await _dbContext.SaveChangesAsync(); }
4. 关键注意事项
- 实例标识唯一性:确保每个实例的
ProcessingInstanceId唯一,避免不同实例抢占时冲突。 - 超时阈值配置:根据Blob上传的平均耗时调整
timeoutMinutes,防止因个别大文件处理超时导致数据无法被重新处理。 - 批量操作优化:通过
SaveChangesAsync批量提交更新,减少数据库交互次数;Blob上传可并行处理(需注意Blob存储的并发限制)。 - 隔离级别选择:使用
ReadCommitted即可,无需更高隔离级别,避免过度锁表引发超时。
内容的提问来源于stack exchange,提问作者Aether
相关产品推荐
相关产品推荐

