【问题标题】:Batching futures in DartDart 中的批处理期货
【发布时间】:2020-01-30 08:27:46
【问题描述】:

我想将多个 future 批处理到一个请求中,当达到最大批处理大小或达到自收到最早的 future 以来的最长时间时触发。

动机

在 Flutter 中,我有许多 UI 元素需要显示未来的结果,这取决于 UI 元素中的数据。

例如,我有一个地方小部件和一个显示步行到某个地方需要多长时间的子小部件。为了计算步行需要多长时间,我向 Google Maps API 发出请求以获取到该地点的行程时间。

将所有这些 API 请求批处理成一个批处理 API 请求会更加高效和经济。因此,如果小部件即时发出 100 个请求,则可以通过单个提供者代理未来,该提供者将未来批处理成对 Google 的单个请求,并将来自 Google 的结果解包到所有单独的请求中。

提供者需要知道何时停止等待更多的未来以及何时实际发出请求,这应该由最大“批量”大小(即旅行时间请求的数量)或最大时间量来控制你愿意等待批处理发生。

所需的 API 类似于:


// Client gives this to tell provider how to compute batch result.
abstract class BatchComputer<K,V> {
  Future<List<V>> compute(List<K> batchedInputs);
}

// Batching library returns an object with this interface
// so that client can submit inputs to completed by the Batch provider.
abstract class BatchingFutureProvider<K,V> {
  Future<V> submit(K inputValue);
}

// How do you implement this in dart???
BatchingFutureProvider<K,V> create<K,V>(
   BatchComputer<K,V> computer, 
   int maxBatchSize, 
   Duration maxWaitDuration,
);

Dart(或 pub 包)是否已经提供了这种批处理功能,如果没有,您将如何实现上面的 create 函数?

【问题讨论】:

  • 您的意思是Future.wait()?等待未来列表,完成后执行回调。
  • 所以你必须使用Future api - 注意then() 和timeout() 方法
  • @Tokenyet 我不明白 Future.wait() 是如何应用的,你能解释一下它如何帮助实现上面的create 吗?
  • @pskink 我觉得您的回答可能太低级,无法对这个问题有建设性。我对Future 的基本单位有所了解,但是我正在寻找以允许在符合人体工程学的 API 中进行批处理的方式将它们组合在一起

标签: flutter dart future batching


【解决方案1】:

这听起来很合理,但也非常专业。 您需要一种表示查询的方法,将这些查询组合成一个超级查询,然后将超级结果拆分为单独的结果,这就是您的BatchComputer 所做的。然后你需要一个队列,你可以在某些条件下刷新它。

有一点很清楚,您将需要使用 Completers 来获取结果,因为当您想要返回未来时,您总是需要它,然后才能获得价值或未来来完成它。

我会选择的方法是:

import "dart:async";

/// A batch of requests to be handled together.
///
/// Collects [Request]s until the pending requests are flushed.
/// Requests can be flushed by calling [flush] or by configuring
/// the batch to automatically flush when reaching certain 
/// tresholds.
class BatchRequest<Request, Response> {
  final int _maxRequests;
  final Duration _maxDelay;
  final Future<List<Response>> Function(List<Request>) _compute;
  Timer _timeout;
  List<Request> _pendingRequests;
  List<Completer<Response>> _responseCompleters;

  /// Creates a batcher of [Request]s.
  ///
  /// Batches requests until calling [flush]. At that pont, the
  /// [batchCompute] function gets the list of pending requests,
  /// and it should respond with a list of [Response]s.
  /// The response to the a request in the argument list
  /// should be at the same index in the response list, 
  /// and as such, the response list must have the same number
  /// of responses as there were requests.
  ///
  /// If [maxRequestsPerBatch] is supplied, requests are automatically
  /// flushed whenever there are that many requests pending.
  ///
  /// If [maxDelay] is supplied, requests are automatically flushed 
  /// when the oldest request has been pending for that long. 
  /// As such, The [maxDelay] is not the maximal time before a request
  /// is answered, just how long sending the request may be delayed.
  BatchRequest(Future<List<Response>> Function(List<Request>) batchCompute,
               {int maxRequestsPerBatch, Duration maxDelay})
    : _compute = batchCompute,
      _maxRequests = maxRequestsPerBatch,
      _maxDelay = maxDelay;

  /// Add a request to the batch.
  ///
  /// The request is stored until the requests are flushed,
  /// then the returned future is completed with the result (or error)
  /// received from handling the requests.
  Future<Response> addRequest(Request request) {
    var completer = Completer<Response>();
    (_pendingRequests ??= []).add(request);
    (_responseCompleters ??= []).add(completer);
    if (_pendingRequests.length == _maxRequests) {
      _flush();
    } else if (_timeout == null && _maxDelay != null) {
      _timeout = Timer(_maxDelay, _flush);
    }
    return completer.future;
  }

  /// Flush any pending requests immediately.
  void flush() {
    _flush();
  }

  void _flush() {
    if (_pendingRequests == null) {
      assert(_timeout == null);
      assert(_responseCompleters == null);
      return;
    }
    if (_timeout != null) {
      _timeout.cancel();
      _timeout = null;
    }
    var requests = _pendingRequests;
    var completers = _responseCompleters;
    _pendingRequests = null;
    _responseCompleters = null;

    _compute(requests).then((List<Response> results) {
      if (results.length != completers.length) {
        throw StateError("Wrong number of results. "
           "Expected ${completers.length}, got ${results.length}");
      }
      for (int i = 0; i < results.length; i++) {
        completers[i].complete(results[i]);
      }
    }).catchError((error, stack) {
      for (var completer in completers) {
        completer.completeError(error, stack);
      }
    });
  }
}

您可以使用它,例如:

void main() async {
  var b = BatchRequest<int, int>(_compute, 
      maxRequestsPerBatch: 5, maxDelay: Duration(seconds: 1));
  var sw = Stopwatch()..start();
  for (int i = 0; i < 8; i++) {
    b.addRequest(i).then((r) {
      print("${sw.elapsedMilliseconds.toString().padLeft(4)}: $i -> $r");
    });
  }
}
Future<List<int>> _compute(List<int> args) => 
    Future.value([for (var x in args) x + 1]);

【讨论】:

    【解决方案2】:

    见https://pub.dev/packages/batching_future/versions/0.0.2

    我的答案与@lrn 几乎完全相同,但我努力使主线同步,并添加了一些文档。

    /// Exposes [createBatcher] which batches computation requests until either
    /// a max batch size or max wait duration is reached.
    ///
    import 'dart:async';
    
    import 'dart:collection';
    
    import 'package:quiver/iterables.dart';
    import 'package:synchronized/synchronized.dart';
    
    /// Converts input type [K] to output type [V] for every item in
    /// [batchedInputs]. There must be exactly one item in output list for every
    /// item in input list, and assumes that input[i] => output[i].
    abstract class BatchComputer<K, V> {
      const BatchComputer();
      Future<List<V>> compute(List<K> batchedInputs);
    }
    
    /// Interface to submit (possible) batched computation requests.
    abstract class BatchingFutureProvider<K, V> {
      Future<V> submit(K inputValue);
    }
    
    /// Returns a batcher which computes transformations in batch using [computer].
    /// The batcher will wait to compute until [maxWaitDuration] is reached since
    /// the first item in the current batch is received, or [maxBatchSize] items
    /// are in the current batch, whatever happens first.
    /// If [maxBatchSize] or [maxWaitDuration] is null, then the triggering
    /// condition is ignored, but at least one condition must be supplied.
    ///
    /// Warning: If [maxWaitDuration] is not supplied, then it is possible that
    /// a partial batch will never finish computing.
    BatchingFutureProvider<K, V> createBatcher<K, V>(BatchComputer<K, V> computer,
        {int maxBatchSize, Duration maxWaitDuration}) {
      if (!((maxBatchSize != null || maxWaitDuration != null) &&
          (maxWaitDuration == null || maxWaitDuration.inMilliseconds > 0) &&
          (maxBatchSize == null || maxBatchSize > 0))) {
        throw ArgumentError(
            "At least one of {maxBatchSize, maxWaitDuration} must be specified and be positive values");
      }
      return _Impl(computer, maxBatchSize, maxWaitDuration);
    }
    
    // Holds the input value and the future to complete it.
    class _Payload<K, V> {
      final K k;
      final Completer<V> completer;
    
      _Payload(this.k, this.completer);
    }
    
    enum _ExecuteCommand { EXECUTE }
    
    /// Implements [createBatcher].
    class _Impl<K, V> implements BatchingFutureProvider<K, V> {
      /// Queues computation requests.
      final controller = StreamController<dynamic>();
    
      /// Queues the input values with their futures to complete.
      final queue = Queue<_Payload>();
    
      /// Locks access to [listen] to make queue-processing single-threaded.
      final lock = Lock();
    
      /// [maxWaitDuration] timer, as a stored reference to cancel early if needed.
      Timer timer;
    
      /// Performs the input->output batch transformation.
      final BatchComputer computer;
    
      /// See [createBatcher].
      final int maxBatchSize;
    
      /// See [createBatcher].
      final Duration maxWaitDuration;
      _Impl(this.computer, this.maxBatchSize, this.maxWaitDuration) {
        controller.stream.listen(listen);
      }
    
      void dispose() {
        controller.close();
      }
    
      @override
      Future<V> submit(K inputValue) {
        final completer = Completer<V>();
        controller.add(_Payload(inputValue, completer));
        return completer.future;
      }
    
      // Synchronous event-processing logic.
      void listen(dynamic event) async {
        await lock.synchronized(() {
          if (event.runtimeType == _ExecuteCommand) {
            if (timer?.isActive ?? true) {
              // The timer got reset, so ignore this old request.
              // The current timer needs to inactive and non-null
              // for the execution to be legitimate.
              return;
            }
            execute();
          } else {
            addPayload(event as _Payload);
          }
          return;
        });
      }
    
      void addPayload(_Payload _payload) {
        if (queue.isEmpty && maxWaitDuration != null) {
          // This is the first item of the batch.
          // Trigger the timer so we are guaranteed to start computing
          // this batch before [maxWaitDuration].
          timer = Timer(maxWaitDuration, triggerTimer);
        }
        queue.add(_payload);
        if (maxBatchSize != null && queue.length >= maxBatchSize) {
          execute();
          return;
        }
      }
    
      void execute() async {
        timer?.cancel();
        if (queue.isEmpty) {
          return;
        }
        final results = await computer.compute(List<K>.of(queue.map((p) => p.k)));
        for (var pair in zip<Object>([queue, results])) {
          (pair[0] as _Payload).completer.complete(pair[1] as V);
        }
        queue.clear();
      }
    
      void triggerTimer() {
        listen(_ExecuteCommand.EXECUTE);
      }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-10-17
      • 1970-01-01
      • 2020-05-27
      • 2020-08-21
      • 1970-01-01
      • 1970-01-01
      • 2020-05-09
      • 2017-09-30
      相关资源
      最近更新 更多