【问题标题】:Multithreading Entity Framework: The connection was not closed. The connection's current state is connecting多线程实体框架:连接未关闭。连接的当前状态是正在连接
【发布时间】:2012-09-13 08:30:52
【问题描述】:

所以我有一个执行工作流过程的 Windows 服务进程。后端在 Entity Framework 之上使用 Repository 和 UnitofWork Pattern 和 Unity,以及从 edmx 生成的实体类。我不会详细介绍它,因为它不是必需的,但基本上有 5 个工作流程要经过的步骤。一个特定的过程可能在任何时间点的任何阶段(当然是按顺序)。第一步只是为第二步生成数据,第二步通过一个长时间运行的进程验证数据到另一台服务器。然后在那里生成包含该数据的pdf。对于每个阶段,我们都会生成一个计时器,但是可以配置为允许为每个阶段生成多个计时器。问题就在于此。当我将处理器添加到特定阶段时,它随机出现以下错误:

连接没有关闭。连接的当前状态为正在连接。

阅读这似乎很明显,这是因为上下文试图从两个线程访问同一个实体。但这就是让我陷入困境的地方。我能找到的所有信息都表明我们应该为每个线程使用一个实例上下文。据我所知,我在做什么(见下面的代码)。我没有使用单例模式或静态或任何东西,所以我不确定为什么会发生这种情况或如何避免它。我已经在下面发布了我的代码的相关部分供您查看。

基础存储库:

 public class BaseRepository
{
    /// <summary>
    /// Initializes a repository and registers with a <see cref="IUnitOfWork"/>
    /// </summary>
    /// <param name="unitOfWork"></param>
    public BaseRepository(IUnitOfWork unitOfWork)
    {
        if (unitOfWork == null) throw new ArgumentException("unitofWork");
        UnitOfWork = unitOfWork;
    }


    /// <summary>
    /// Returns a <see cref="DbSet"/> of entities.
    /// </summary>
    /// <typeparam name="TEntity">Entity type the dbset needs to return.</typeparam>
    /// <returns></returns>
    protected virtual DbSet<TEntity> GetDbSet<TEntity>() where TEntity : class
    {

        return Context.Set<TEntity>();
    }

    /// <summary>
    /// Sets the state of an entity.
    /// </summary>
    /// <param name="entity">object to set state.</param>
    /// <param name="entityState"><see cref="EntityState"/></param>
    protected virtual void SetEntityState(object entity, EntityState entityState)
    {
        Context.Entry(entity).State = entityState;
    }

    /// <summary>
    /// Unit of work controlling this repository.       
    /// </summary>
    protected IUnitOfWork UnitOfWork { get; set; }

    /// <summary>
    /// 
    /// </summary>
    /// <param name="entity"></param>
    protected virtual void Attach(object entity)
    {
        if (Context.Entry(entity).State == EntityState.Detached)
            Context.Entry(entity).State = EntityState.Modified;
    }

    protected virtual void Detach(object entity)
    {
        Context.Entry(entity).State = EntityState.Detached;
    }

    /// <summary>
    /// Provides access to the ef context we are working with
    /// </summary>
    internal StatementAutoEntities Context
    {
        get
        {                
            return (StatementAutoEntities)UnitOfWork;
        }
    }
}

StatementAutoEntities 是自动生成的 EF 类。

存储库实现:

public class ProcessingQueueRepository : BaseRepository, IProcessingQueueRepository
{

    /// <summary>
    /// Creates a new repository and associated with a <see cref="IUnitOfWork"/>
    /// </summary>
    /// <param name="unitOfWork"></param>
    public ProcessingQueueRepository(IUnitOfWork unitOfWork) : base(unitOfWork)
    {
    }

    /// <summary>
    /// Create a new <see cref="ProcessingQueue"/> entry in database
    /// </summary>
    /// <param name="Queue">
    ///     <see cref="ProcessingQueue"/>
    /// </param>
    public void Create(ProcessingQueue Queue)
    {
        GetDbSet<ProcessingQueue>().Add(Queue);
        UnitOfWork.SaveChanges();
    }

    /// <summary>
    /// Updates a <see cref="ProcessingQueue"/> entry in database
    /// </summary>
    /// <param name="queue">
    ///     <see cref="ProcessingQueue"/>
    /// </param>
    public void Update(ProcessingQueue queue)
    {
        //Attach(queue);
        UnitOfWork.SaveChanges();
    }

    /// <summary>
    /// Delete a <see cref="ProcessingQueue"/> entry in database
    /// </summary>
    /// <param name="Queue">
    ///     <see cref="ProcessingQueue"/>
    /// </param>
    public void Delete(ProcessingQueue Queue)
    {
        GetDbSet<ProcessingQueue>().Remove(Queue);  
        UnitOfWork.SaveChanges();
    }

    /// <summary>
    /// Gets a <see cref="ProcessingQueue"/> by its unique Id
    /// </summary>
    /// <param name="id"></param>
    /// <returns></returns>
    public ProcessingQueue GetById(int id)
    {
        return (from e in Context.ProcessingQueue_SelectById(id) select e).FirstOrDefault();
    }

    /// <summary>
    /// Gets a list of <see cref="ProcessingQueue"/> entries by status
    /// </summary>
    /// <param name="status"></param>
    /// <returns></returns>
    public IList<ProcessingQueue> GetByStatus(int status)
    {
        return (from e in Context.ProcessingQueue_SelectByStatus(status) select e).ToList();
    }

    /// <summary>
    /// Gets a list of all <see cref="ProcessingQueue"/> entries
    /// </summary>
    /// <returns></returns>
    public IList<ProcessingQueue> GetAll()
    {
        return (from e in Context.ProcessingQueue_Select() select e).ToList();
    }

    /// <summary>
    /// Gets the next pending item id in the queue for a specific work        
    /// </summary>
    /// <param name="serverId">Unique id of the server that will process the item in the queue</param>
    /// <param name="workTypeId">type of <see cref="WorkType"/> we are looking for</param>
    /// <param name="operationId">if defined only operations of the type indicated are considered.</param>
    /// <returns>Next pending item in the queue for the work type or null if no pending work is found</returns>
    public int GetNextPendingItemId(int serverId, int workTypeId, int? operationId)
    {
        var id = Context.ProcessingQueue_GetNextPending(serverId, workTypeId,  operationId).SingleOrDefault();
        return id.HasValue ? id.Value : -1;
    }

    /// <summary>
    /// Returns a list of <see cref="ProcessingQueueStatus_dto"/>s objects with all
    /// active entries in the queue
    /// </summary>
    /// <returns></returns>
    public IList<ProcessingQueueStatus_dto> GetActiveStatusEntries()
    {
        return (from e in Context.ProcessingQueueStatus_Select() select e).ToList();
    }
    /// <summary>
    /// Bumps an entry to the front of the queue 
    /// </summary>
    /// <param name="processingQueueId"></param>
    public void Bump(int processingQueueId)
    {
        Context.ProcessingQueue_Bump(processingQueueId);
    }
}

我们使用Unity进行依赖注入,例如一些调用代码:

#region Members
    private readonly IProcessingQueueRepository _queueRepository;       
    #endregion

    #region Constructors
    /// <summary>Initializes ProcessingQueue services with repositories</summary>
    /// <param name="queueRepository"><see cref="IProcessingQueueRepository"/></param>        
    public ProcessingQueueService(IProcessingQueueRepository queueRepository)
    {
        Check.Require(queueRepository != null, "processingQueueRepository is required");
        _queueRepository = queueRepository;

    }
    #endregion

windows服务中启动定时器的代码如下:

            _staWorkTypeConfigLock.EnterReadLock();
        foreach (var timer in from operation in (from o in _staWorkTypeConfig.WorkOperations where o.UseQueueForExecution && o.AssignedProcessors > 0 select o) 
                              let interval = operation.SpawnInternval < 30 ? 30 : operation.SpawnInternval 
                              select new StaTimer
                            {
                                Interval = _runImmediate ? 5000 : interval*1000,
                                Operation = (ProcessingQueue.RequestedOperation) operation.OperationId
                            })
        {
            timer.Elapsed += ApxQueueProcessingOnElapsedInterval;
            timer.Enabled = true;
            Logger.DebugFormat("Queue processing for operations of type {0} will execute every {1} seconds", timer.Operation, timer.Interval/1000);                
        }
        _staWorkTypeConfigLock.ExitReadLock();

StaTimer 只是对定时器添加操作类型的封装。 ApxQueueProcessingOnElapsedInterval 然后基本上只是根据操作将工作分配给进程。

我还将在我们生成任务的地方添加一些 ApxQueueProcessingOnElapsedInterval 代码。

            _staTasksLock.EnterWriteLock();
        for (var x = 0; x < tasksNeeded; x++)
        {
            var t = new Task(obj => ProcessStaQueue((QueueProcessConfig) obj),
                             CreateQueueProcessConfig(true, operation), _cancellationToken);


            _staTasks.Add(new Tuple<ProcessingQueue.RequestedOperation, DateTime, Task>(operation, DateTime.Now,t));

            t.Start();
            Thread.Sleep(300); //so there are less conflicts fighting for jobs in the queue table
        }
        _staTasksLock.ExitWriteLock();

【问题讨论】:

  • 你是 IoC 容器 Dispose 上下文实例吗?

标签: multithreading entity-framework dbcontext


【解决方案1】:

看起来您的服务、存储库和上下文应该在您的应用程序的整个生命周期中都存在,但这是不正确的。您可以同时触发多个计时器。这意味着多个线程将并行使用您的服务,并且它们将在其线程中执行您的服务代码 = 上下文在多个线程之间共享 => 异常,因为上下文不是线程安全的。

唯一的选择是为每个要执行的操作使用一个新的上下文实例。例如,您可以更改您的类以接受上下文工厂而不是上下文,并为每个操作获取一个新的上下文。

【讨论】:

  • 我添加了一些代码来显示服务是如何产生的。我正在使用一个生成 ProcessStaQueue 的任务,然后它会计算出它应该运行什么操作。我知道 Context 不是线程安全的,但鉴于我没有使用静态或单例模式或任何东西,不应该为每个执行的任务独立实例化一个新的上下文吗?
  • 服务最终被实例化如下:ServiceLocator.Current.GetInstance();这可能与它有关吗?我的猜测是,也许 ServiceLocator 与此有​​关。
  • 但是,如果它在您每次请求时创建一个新的服务实例,或者它只为第一个请求创建实例而不是重用该实例,则取决于服务定位器的实现。跨度>
  • 配置文件中我的所有对象的生命周期都设置为“perthread”。
  • 在 .NET Core 2.1 发布之前,如果您有一个针对多个存储库方法调用执行事务的服务方法,这将是一个问题。如果你尝试使用 TransactionScope(从 2.0 开始),它会给你一个例外。
【解决方案2】:

如果这对任何人都有帮助:

就我而言,我确实确保了非线程安全的 DbContext 有一个 TransientLifetime(使用 Ninject),但它仍然会导致并发问题!事实证明,在我的一些自定义ActionFilters 中,我使用依赖注入来访问构造函数中的DbContext,但ActionFilters 有一个生命周期,可以让它们在多个请求中实例化,因此没有重新创建上下文.

我通过在 OnActionExecuting 方法中而不是在构造函数中手动解决依赖关系来修复它,以便它每次都是一个新实例。

【讨论】:

    【解决方案3】:

    就我而言,我遇到了这个问题,因为我在我的一个 DAL 函数调用之前忘记了 await 关键字。将await 放在那里解决了它。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-05-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-01-07
      相关资源
      最近更新 更多