解决PostgreSQL查询时的NpgsqlOperationInProgressException异常
问题描述
执行查询时反复遇到以下异常:
Npgsql.NpgsqlOperationInProgressException (0x80004005): A command is already in progress.
我使用的是Dapper和PostgreSQL,以下是CarsService中的数据访问代码:
public class CarsService : ICarsService { private readonly IDbConnection _connection; public CarsService(IDbConnection connection) => _connection = connection; public Task<IEnumerable<int>> GetAvailableCarsIdsAsync(int carModelId, DateTime startDate, DateTime endDate) { return _connection.QueryAsync<int>(@" select c.id from cars c where c.car_model_id = @CarModelId and c.id not in ( select r.car_id from rentals r where r.car_id = c.id and r.start_date <= @EndDate and r.end_date >= @StartDate )", new { CarModelId = carModelId, StartDate = startDate, EndDate = endDate }); } public Task<IEnumerable<int>> GetCarIdsAsync(int carModelId) { return _connection.QueryAsync<int>(@" select c.id from cars c where c.car_model_id = @CarModelId", new { CarModelId = carModelId }); } }
调用该服务的方法如下:
public async Task<IEnumerable<GetCarModelResponse>> Handle(GetCarModelsQuery query, CancellationToken cancellationToken) { var carModels = await _carModelsService.GetAsync(query); return await Task.WhenAll(carModels.Select ( async carModel => { var availableCarIds = await _carsService.GetAvailableCarsIdsAsync(carModel.Id, query.StartDate, query.EndDate); return new GetCarModelResponse ( carModel.Id, carModel.Model, carModel.Color, availableCarIds.Count(), availableCarIds, await _carsService.GetCarIdsAsync(carModel.Id) ); })); }
IDbConnection是以Transient方式注入的:
builder.Services.AddTransient<IDbConnection>(_ => new NpgsqlConnection("connection string"));
尝试改为Scoped或Singleton仍无法解决问题,请问该如何修复?
问题原因
核心问题是同一个NpgsqlConnection实例被多个异步操作并行使用:
NpgsqlConnection本身不是线程安全的,且PostgreSQL不支持多活动结果集(MARS),同一连接无法同时执行多个未完成的命令。- 尽管注册了Transient的
IDbConnection,但如果CarsService是Scoped生命周期(默认注入方式),则每个请求中CarsService只会被实例化一次,其依赖的_connection也会是同一个实例。 - 在
Task.WhenAll的并行逻辑中,多个carModel会同时调用CarsService的两个异步方法,导致同一连接被并发调用,触发异常。
解决方案
方案1:使用连接工厂模式,确保每个操作使用独立连接
定义连接工厂接口和实现:
public interface IDbConnectionFactory { IDbConnection CreateConnection(); } public class NpgsqlConnectionFactory : IDbConnectionFactory { private readonly string _connectionString; public NpgsqlConnectionFactory(string connectionString) { _connectionString = connectionString; } public IDbConnection CreateConnection() { var connection = new NpgsqlConnection(_connectionString); connection.Open(); // Dapper会自动打开连接,但显式打开可避免重复开销 return connection; } }
注册连接工厂为Singleton:
builder.Services.AddSingleton<IDbConnectionFactory>(_ => new NpgsqlConnectionFactory("your connection string"));
修改CarsService,每次数据库操作创建新连接并自动释放:
public class CarsService : ICarsService { private readonly IDbConnectionFactory _connectionFactory; public CarsService(IDbConnectionFactory connectionFactory) => _connectionFactory = connectionFactory; public async Task<IEnumerable<int>> GetAvailableCarsIdsAsync(int carModelId, DateTime startDate, DateTime endDate) { using var connection = _connectionFactory.CreateConnection(); return await connection.QueryAsync<int>(@" select c.id from cars c where c.car_model_id = @CarModelId and c.id not in ( select r.car_id from rentals r where r.car_id = c.id and r.start_date <= @EndDate and r.end_date >= @StartDate )", new { CarModelId = carModelId, StartDate = startDate, EndDate = endDate }); } public async Task<IEnumerable<int>> GetCarIdsAsync(int carModelId) { using var connection = _connectionFactory.CreateConnection(); return await connection.QueryAsync<int>(@" select c.id from cars c where c.car_model_id = @CarModelId", new { CarModelId = carModelId }); } }
方案2:合并查询减少数据库调用(更优)
并行调用多个小查询不仅容易引发连接冲突,还会增加数据库往返开销。可以通过一个SQL查询一次性获取所有车型的总车辆数、可用车辆数及对应ID列表:
// 在CarsService中新增方法 public async Task<IEnumerable<CarModelStats>> GetCarModelsStatsAsync(IEnumerable<int> carModelIds, DateTime startDate, DateTime endDate) { using var connection = _connectionFactory.CreateConnection(); return await connection.QueryAsync<CarModelStats>(@" select c.car_model_id as ModelId, count(c.id) as TotalCars, count(case when r.car_id is null then 1 end) as AvailableCars, array_agg(c.id) as AllCarIds, array_remove(array_agg(case when r.car_id is null then c.id end), null) as AvailableCarIds from cars c left join rentals r on c.id = r.car_id and r.start_date <= @EndDate and r.end_date >= @StartDate where c.car_model_id = any(@CarModelIds) group by c.car_model_id", new { CarModelIds = carModelIds.ToArray(), StartDate = startDate, EndDate = endDate }); } // 对应的DTO类 public class CarModelStats { public int ModelId { get; set; } public int TotalCars { get; set; } public int AvailableCars { get; set; } public int[] AllCarIds { get; set; } public int[] AvailableCarIds { get; set; } }
修改Handle方法,避免并行调用:
public async Task<IEnumerable<GetCarModelResponse>> Handle(GetCarModelsQuery query, CancellationToken cancellationToken) { var carModels = await _carModelsService.GetAsync(query); var modelIds = carModels.Select(m => m.Id).ToArray(); var stats = await _carsService.GetCarModelsStatsAsync(modelIds, query.StartDate, query.EndDate); var statsDict = stats.ToDictionary(s => s.ModelId); return carModels.Select(carModel => { statsDict.TryGetValue(carModel.Id, out var modelStats); return new GetCarModelResponse ( carModel.Id, carModel.Model, carModel.Color, modelStats?.AvailableCars ?? 0, modelStats?.AvailableCarIds ?? Array.Empty<int>(), modelStats?.AllCarIds ?? Array.Empty<int>() ); }); }
内容的提问来源于stack exchange,提问作者Kacper Wyczawski
相关产品推荐
相关产品推荐

