【问题标题】:How to implement an async task queue with multiple concurrent workers (async) in dart如何在 dart 中实现具有多个并发工作者(异步)的异步任务队列
【发布时间】:2020-11-02 19:36:38
【问题描述】:

我的目标是在 dart 中创建一种网络爬虫。为此,我想维护一个任务队列,其中存储了需要抓取的元素(例如 URL)。在 crawl 函数中对元素进行爬取,该函数返回需要处理的更多元素的列表。因此,这些元素被添加到队列中。示例代码:

import "dart:collection";
final queue = Queue<String>();
main() async{
  queue
    ..add("...")
    ..add("...")
    ..add("...");
  while (queue.isNotEmpty) {
    results = await crawl(queue.removeFirst());
    queue.addAll(results);
  }
}

Future<List<String>> crawl(String x) async {
  ...
  res = await http.get(x)
  ...
  return results;
}

这个粗略的代码一次只处理一个元素。但是,我希望有一个工作人员池(例如 5 个)将元素从队列中取出并同时处理它们,然后将结果添加回队列中。由于瓶颈是 HTTP 请求,我认为 Future.wait() 调用多个工作人员可以加快执行速度。但是我不想让服务器超载,因此我也想限制工人的数量。

这可以通过基本的异步原语和信号量来实现吗?如果可能,我想避免隔离,以使解决方案尽可能简单。

【问题讨论】:

    标签: flutter asynchronous dart queue


    【解决方案1】:

    我不知道那里是否已经有一个包提供了这个功能,但是因为编写你自己的逻辑并不复杂,所以我做了以下示例:

    import 'dart:async';
    import 'dart:collection';
    import 'dart:math';
    
    class TaskRunner<A, B> {
      final Queue<A> _input = Queue();
      final StreamController<B> _streamController = StreamController();
      final Future<B> Function(A) task;
    
      final int maxConcurrentTasks;
      int runningTasks = 0;
    
      TaskRunner(this.task, {this.maxConcurrentTasks = 5});
    
      Stream<B> get stream => _streamController.stream;
    
      void add(A value) {
        _input.add(value);
        _startExecution();
      }
    
      void addAll(Iterable<A> iterable) {
        _input.addAll(iterable);
        _startExecution();
      }
    
      void _startExecution() {
        if (runningTasks == maxConcurrentTasks || _input.isEmpty) {
          return;
        }
    
        while (_input.isNotEmpty && runningTasks < maxConcurrentTasks) {
          runningTasks++;
          print('Concurrent workers: $runningTasks');
    
          task(_input.removeFirst()).then((value) async {
            _streamController.add(value);
    
            while (_input.isNotEmpty) {
              _streamController.add(await task(_input.removeFirst()));
            }
    
            runningTasks--;
            print('Concurrent workers: $runningTasks');
          });
        }
      }
    }
    
    Random _rnd = Random();
    Future<List<String>> crawl(String x) =>
        Future.delayed(Duration(seconds: _rnd.nextInt(5)), () => x.split('-'));
    
    void main() {
      final runner = TaskRunner(crawl, maxConcurrentTasks: 3);
    
      runner.stream.forEach((listOfString) {
        if (listOfString.length == 1) {
          print('DONE: ${listOfString.first}');
        } else {
          print('PUTTING STRINGS ON QUEUE: $listOfString');
          runner.addAll(listOfString);
        }
      });
    
      runner.addAll(['1-2-3-4-5-6-7-8-9', '10-20-30-40-50-60-70-80-90']);
    }
    

    哪些输出:

    Concurrent workers: 1
    Concurrent workers: 2
    Concurrent workers: 1
    PUTTING STRINGS ON QUEUE: [1, 2, 3, 4, 5, 6, 7, 8, 9]
    Concurrent workers: 2
    Concurrent workers: 3
    Concurrent workers: 4
    PUTTING STRINGS ON QUEUE: [10, 20, 30, 40, 50, 60, 70, 80, 90]
    DONE: 3
    DONE: 5
    DONE: 1
    DONE: 2
    DONE: 7
    DONE: 4
    DONE: 6
    DONE: 10
    DONE: 8
    DONE: 9
    DONE: 30
    DONE: 20
    DONE: 40
    DONE: 50
    Concurrent workers: 3
    DONE: 90
    Concurrent workers: 2
    DONE: 60
    Concurrent workers: 1
    DONE: 80
    Concurrent workers: 0
    DONE: 70
    

    我确信该类的可用性可以提高,但我认为核心概念很容易理解。概念是我们定义了一个Queue,每次我们向这个Queue 添加东西时,我们都会检查是否可以开始执行新的异步任务。否则我们只是跳过它,因为我们确保每个当前正在运行的异步任务都会在“关闭”之前检查 Queue 上的更多内容。

    结果由您可以订阅的Stream 返回,例如根据我在示例中显示的结果向TaskRunner 添加更多内容。返回数据的顺序是基于它们完成的顺序。

    重要的是,这不是在多个线程中运行任务的方法。所有代码都在单个 Dart 隔离线程中运行,但由于 HTTP 请求是 IO 延迟的,因此尝试生成多个 Future 并等待结果是有意义的。

    【讨论】:

    • 太棒了,非常感谢!即使在您的示例中 maxConcurrentTasks=3,为什么会显示“并发工作人员:4”?
    • 我猜while (_input.isNotEmpty &amp;&amp; runningTasks &lt;= maxConcurrentTasks) { 应该是while (_input.isNotEmpty &amp;&amp; runningTasks &lt; maxConcurrentTasks) {
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-01-05
    • 2021-11-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多