You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

解决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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 22:32:48