【问题标题】:How to perform parallel Insert/Update/Delete operations on same table and context in hosted services?如何在托管服务中对同一表和上下文执行并行插入/更新/删除操作?
【发布时间】:2019-08-23 23:12:43
【问题描述】:

所以基本上我有 4 个环境 SB、NP、PR 和 DR,并且每个环境都有一个 API 端点。我要做的是调用每个 API 端点并获取一些数据并将其保存到数据库中。因此,对于所有环境,只有一张表将分别存储基础 ID 为 1、2、3 和 4 的所有 apps,类似地,我有 orgs 和 spaces 的表。

1) 我需要根据从每个 API 端点获得的数据更新这些表,并且我需要调用每个 API every second。因此,在新数据到达之前,以前的数据很有可能仍未存储,这是我不知道如何解决的问题。

2) 为了每秒调用这些端点,我为每个环境创建了托管服务,该服务还具有将相应数据保存在其表中的逻辑,如下所示。

这种方法的问题是我对所有四个托管服务都使用相同的 DBContext 并出现以下错误,我知道 EF 不支持对相同 DBcontext 的并行操作,但我不知道如何解决这个问题,从文档中可以使用await 解决,但我不确定如何将它用于所有 4 个托管服务。

我正在尝试实施火灾并忘记了一些事情。

“DBContext.Organizations”与 'DBContext.Organizations'

public class TokenService : DelegatingHandler, IHostedService
{
    public IConfiguration Configuration { get; }
    protected IMemoryCache _cache;
    private Timer _timer;
    public IHttpClientFactory _clientFactory;
    protected HttpClient _client_NP;
    private readonly IServiceScopeFactory _scopeFactory;

    public TokenService(IConfiguration configuration, IMemoryCache memoryCache, IHttpClientFactory clientFactory, IServiceScopeFactory scopeFactory)
    {
        Configuration = configuration;
        _cache = memoryCache;
        _clientFactory = clientFactory;
        _scopeFactory = scopeFactory;

        // NamedClients foreach Env.
        _client_NP = _clientFactory.CreateClient("NonProductionEnv");
    }

    public Task StartAsync(CancellationToken cancellationToken)
    {
        _timer = new Timer(GetAccessToken, null, 0, 3300000);
        _timer = new Timer(Heartbeat, null, 1000, 1000);
        return Task.CompletedTask;
    }

    public Task StopAsync(CancellationToken cancellationToken)
    {
        //Timer does not have a stop. 
        _timer?.Change(Timeout.Infinite, 0);
        return Task.CompletedTask;
    }

    public async Task<Token> GetToken(Uri authenticationUrl, Dictionary<string, string> authenticationCredentials)
    {
        HttpClient client = new HttpClient();
        FormUrlEncodedContent content = new FormUrlEncodedContent(authenticationCredentials);
        HttpResponseMessage response = await client.PostAsync(authenticationUrl, content);

        if (response.StatusCode != System.Net.HttpStatusCode.OK)
        {
            string message = String.Format("POST failed. Received HTTP {0}", response.StatusCode);
            throw new ApplicationException(message);
        }

        string responseString = await response.Content.ReadAsStringAsync();
        Token token = JsonConvert.DeserializeObject<Token>(responseString);

        return token;
    }

    private void GetAccessToken(object state)
    {
        Dictionary<string, string> authenticationCredentials_np = Configuration.GetSection("NonProductionEnvironment:Credentials").GetChildren().Select(x => new KeyValuePair<string, string>(x.Key, x.Value)).ToDictionary(x => x.Key, x => x.Value);
        Token token_np = GetToken(new Uri(Configuration["NonProductionEnvironment:URL"]), authenticationCredentials_np).Result;

        _client_NP.DefaultRequestHeaders.Add("Authorization", $"Bearer {token_np.AccessToken}");
    }

    public void Heartbeat(object state)
    {
        // Discard the result
        _ = GetOrg();
    }

    public async Task GetOrg()
    {
        var request = new HttpRequestMessage(HttpMethod.Get, "organizations");
        var response = await _client_NP.SendAsync(request);
        var json = await response.Content.ReadAsStringAsync();
        OrganizationsClass.OrgsRootObject model = JsonConvert.DeserializeObject<OrganizationsClass.OrgsRootObject>(json);

        using (var scope = _scopeFactory.CreateScope())
        {
            var _DBcontext = scope.ServiceProvider.GetRequiredService<DBContext>();

            foreach (var item in model.resources)
            {
                var g = Guid.Parse(item.guid);
                var x = _DBcontext.Organizations.FirstOrDefault(o => o.OrgGuid == g);
                if (x == null)
                {
                    _DBcontext.Organizations.Add(new Organizations
                    {
                        OrgGuid = g,
                        Name = item.name,
                        CreatedAt = item.created_at,
                        UpdatedAt = item.updated_at,
                        Timestamp = DateTime.Now,
                        Foundation = 2
                    });
                }
                else if (x.UpdatedAt != item.updated_at)
                {
                    x.CreatedAt = item.created_at;
                    x.UpdatedAt = item.updated_at;
                    x.Timestamp = DateTime.Now;
                }
            }

            await GetSpace();
            await _DBcontext.SaveChangesAsync();
        }
    }

    public async Task GetSpace()
        {
            var request = new HttpRequestMessage(HttpMethod.Get, "spaces");
            var response = await _client_NP.SendAsync(request);
            var json = await response.Content.ReadAsStringAsync();
            SpacesClass.SpaceRootObject model = JsonConvert.DeserializeObject<SpacesClass.SpaceRootObject>(json);

            using (var scope = _scopeFactory.CreateScope())
            {
                var _DBcontext = scope.ServiceProvider.GetRequiredService<DBContext>();

                foreach (var item in model.resources)
                {
                    var g = Guid.Parse(item.guid);
                    var x = _DBcontext.Spaces.FirstOrDefault(o => o.SpaceGuid == g);
                    if (x == null)
                    {
                        _DBcontext.Spaces.Add(new Spaces
                        {
                            SpaceGuid = Guid.Parse(item.guid),
                            Name = item.name,
                            CreatedAt = item.created_at,
                            UpdatedAt = item.updated_at,
                            OrgGuid = Guid.Parse(item.relationships.organization.data.guid),
                            Foundation = 2,
                            Timestamp = DateTime.Now
                        });
                    }

                    else if (x.UpdatedAt != item.updated_at)
                    {
                        x.CreatedAt = item.created_at;
                        x.UpdatedAt = item.updated_at;
                        x.Timestamp = DateTime.Now;
                    }
                }

                await GetApps();
            }
        }

        public async Task GetApps()
        {
            var request = new HttpRequestMessage(HttpMethod.Get, "apps?per_page=200");
            var response = await _client_NP.SendAsync(request);
            var json = await response.Content.ReadAsStringAsync();
            AppsClass.AppsRootobject model = JsonConvert.DeserializeObject<AppsClass.AppsRootobject>(json);
            using (var scope = _scopeFactory.CreateScope())
            {
                var _DBcontext = scope.ServiceProvider.GetRequiredService<DBContext>();

                foreach (var item in model.resources)
                {
                    var g = Guid.Parse(item.guid);
                    var x = _DBcontext.Apps.FirstOrDefault(o => o.AppGuid == g);

                    if (x == null)
                    {
                        _DBcontext.Apps.Add(new Apps
                        {
                            AppGuid = Guid.Parse(item.guid),
                            Name = item.name,
                            State = item.state,
                            CreatedAt = item.created_at,
                            UpdatedAt = item.updated_at,
                            SpaceGuid = Guid.Parse(item.relationships.space.data.guid),
                            Foundation = 2,
                            Timestamp = DateTime.Now
                        });
                    }

                    else if (x.UpdatedAt != item.updated_at)
                    {
                        x.State = item.state;
                        x.CreatedAt = item.created_at;
                        x.UpdatedAt = item.updated_at;
                        x.DeletedAt = null;
                        x.Timestamp = DateTime.Now;
                    }
                }
            }
        }
}

启动:

services.AddDbContext<DBContext>(options => options.UseSqlServer(Configuration.GetConnectionString("DefaultConnection")));

除了foundation 和httpclient 之外,基本上所有服务的代码都是相同的。

有人可以指导正确的方向吗?

【问题讨论】:

  • 您是否遇到构建错误或运行时错误?哪一行代码给出了错误?
  • 有趣的是,我在 _DBcontext.Spaces/Apps/Orgs 上看到了一个错误,但应用程序构建并运行,但它没有将数据保存在数据库中。
  • DBContext是如何在startup中注册的? (例如单例、作用域、瞬态)
  • 每个服务只需要不同的 DbContext。
  • 请澄清实际问题。目前还不清楚,听起来像是 XY 问题。首先,每个方法都使用 separate 数据库上下文实例,因此不存在线程问题。其次,只有GetOrg 保存所做的更改,GetSpace 和GetApps 缺少SaveChanges[Async],因此什么也不做。还有哪一行抛出有问题的异常,确切的异常消息和异常堆栈跟踪是什么?

标签: c# entity-framework asp.net-core .net-core asp.net-core-webapi


【解决方案1】:
  1. 我认为您要求不要每秒都访问数据库有点奇怪。如果您想为每一天的每一秒保存数据,那么无论哪种情况,您每天都会为每个环境创建 86400 条记录。这些数据将需要输入,如果不是每秒,那么每 x 秒你将做相同数量的工作。在任何数据库上每秒 4 次点击也没有那么多,甚至每秒 12 次,就像查​​看您的代码一样,这就是它正在做的事情。有一个很好的变化是数据将在一秒钟内保存。在我正在研究 API 调用限制的项目中,在我们研究其他存储/处理数据的方式之前,我们的大多数调用都通过 webapi,在 SQL DB 中存储/改变数据并更新 Elasticsearch 索引在一秒钟内完成,每天在高峰时段拥有 1000 多个用户。

  2. 与其按照您现在的方式拆分这些方法,不如对 API(GetOrg、GetSpace、GetApps)进行三次调用并返回到一个位置,然后创建一个 DBContext 来保存所有三个结果。保存并处理 DBContext 并在一秒钟后重新执行此操作!每秒一个事务中的一个数据库调用。

所以在 sudo 代码中:

Heartbeat()
{
    var orgResult = GetOrgs()
    var spaceResult = GetSpace()
    var appResult = GetApps()

    var dbContext = blah()

    dbContext.Add(orgResult)
    dbContext.Add(spaceResult)
    dbContext.Add(appResult)

    dbContext.SaveAllThatStuff()
}

例如

GetOrgs()
{
    var response = httpClient.Send(orgRequest)
    var orgResult = jsondeserialize(response.Content)

    return orgResult
}

【讨论】:

  • 有道理,但你的伪代码没有帮助,我需要像我在代码中显示的那样遍历每个项目。
猜你喜欢
  • 1970-01-01
  • 2011-07-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-10-03
相关资源
最近更新 更多