背景

我们在产生某种数据。这些数据有ID作为标识。每30秒会产生新一轮的数据。我们将数据存储在AWS OpenSearch里。所以我们的文档的ID是[文档ID]_[时间戳]。每30秒会产生新一轮的文档。我们需要按照时间的范围查询并处理10个小时之内的文档。现在10个小时创建的文档有8百万之多。我们要求每30秒完成一次对所有文档的处理。

问题分析

我们的第一个问题是处理能力不足。OpenSearch的query每次可以返回的document数量是由限制的。我们所使用的版本每次query最多可以返回1万个document。这样800万需要按照分页查询800次。如果我们认为每次查询的时延是100ms,总的查询时间是:

100 * 800 = 80,000ms = 80s

而事实上因为大规模的持续查询,对于ECS/Fargate和OpenSearch带来极大的资源消耗。所以实际需要的查询时间远远大于80s。在实际环境中,我们需要40分钟才能执行完成。

我们的第二个问题是我们的service不能很好的scale。为了加快我们的查询速度,我们可以选择增强我们的ECS/Fargate的处理能力,也就是scale up我们的service。但是我们知道我们的service的处理能力终究是有限的。所以scale up的策略并不能很好的处理我们的数据。

解决方案

下面是针对上面两个问题的解决方案。

问题一,我们的解决方法是缩短我们查询的时长和查询的内容。因为我们的数据每30秒会产生一次,我们只需要查询30秒的结果就可以得到所有的待处理的文档ID。这样,我们第一次查询的内容从获得10小时的所有文档内容变成了获得30秒内文档的ID。我们大大简化了第一次查询所需要的时间的资源。

问题二,我们将service从scale up的方案变成scale out的方案。我们基于第一步查询的结果,引入一个AWS SQS。我们将所有的ID分成小的batch发送到SQS里。同时,我们创建另一个ECS/Fargate来监听来在SQS的消息。针对每个batch,新的ECS/Fargate并行处理多个ID的文档。利用OpenSearch来查询该ID在过去十小时的所有文档然后加以处理。这样做有两个好处。第一,我们的service现在可以通过增加ECS/Fargate task来解决scale的问题,service的处理能力不再是我们的处理能力瓶颈。第二,我们的query规模小了很多。这样ECS/Fargate和OpenSearch都可以在query结束后释放掉已经使用过资源,从而可以在健康的状态下运行。

结论和思考

我们在这里首先回答两个和设计有关的问题:

1. 为什么不选择Lambda function而选择ECS/Fargate?因为我们有30秒的时间来处理所有的文档,使用ECS/Fargate通过调整我们的batch的大小和每个Task的并行处理逻辑,我们可以更好的对整个系统的并发性进行控制。而Lambda function会让整个系统的并发性随着数据的变大线性增加。同时Lambda function有cold start的问题。为了解决这个问题,我们需要预先评估系统的并发性预留concurrent instance。这样系统控制的复杂度和费用都会增加。

2. 假如一个ID的10个小时的数据量也无法在30秒内处理完成应该怎样?我们注意到这是一个滑动窗口的问题。所以我们可以通过更加复杂的设计来减少计算量。比如下面是一种思路。我们将每30秒产生的数据存储在一个key-value-pair的数据库中,比如使用DynamoDB。同时我们将10小时计算的结果也存储在key-value-pair的数据库中。这样我们每次更新10小时的结果,理论上只需要3个读操作和一个写操作。如下图所示:

更多推荐