Flutter如何实现并行多线程任务 复刻C# Parallelizer并行能力
Flutter 并行任务实现方案
Flutter 完全支持实现你需要的限定并发数的并行处理能力,不需要引入第三方依赖,基于Dart自带的异步、Stream、Isolate能力即可1:1对齐参考的功能逻辑。
概念对应关系
C#实现里的核心能力,在Dart中都有原生对应:
- C#
Task对应DartFuture - C#
CancellationTokenSource/CancellationToken可通过Dart内置的Stream取消能力实现,支持超时、手动取消 - C# 事件钩子(NewResult/Completed/Error等)可通过Dart广播Stream实现,监听更灵活
- 并发度控制可通过异步任务队列原生实现,CPU密集型任务丢入Isolate执行即可避免阻塞UI
完整实现代码
以下代码完全对齐示例逻辑,支持设置最大并发数、跳过前置任务、单任务结果回调、单任务错误回调、全局异常回调、完成回调、超时取消能力:
import 'dart:async'; import 'dart:collection'; // 结果返回结构 class ResultDetails<TInput, TOutput> { final TInput item; final TOutput result; ResultDetails(this.item, this.result); } // 单任务错误结构 class ErrorDetails<TInput> { final TInput item; final Object exception; ErrorDetails(this.item, this.exception); } // 并行任务执行器 class Parallelizer<TInput, TOutput> { final Iterable<TInput> workItems; final Future<TOutput> Function(TInput, CancellationToken) workFunction; final int degreeOfParallelism; final int skip; // 事件流 final StreamController<ResultDetails<TInput, TOutput>> _onNewResult = StreamController.broadcast(); final StreamController _onCompleted = StreamController.broadcast(); final StreamController<Object> _onError = StreamController.broadcast(); final StreamController<ErrorDetails<TInput>> _onTaskError = StreamController.broadcast(); Stream<ResultDetails<TInput, TOutput>> get onNewResult => _onNewResult.stream; Stream get onCompleted => _onCompleted.stream; Stream<Object> get onError => _onError.stream; Stream<ErrorDetails<TInput>> get onTaskError => _onTaskError.stream; final Completer _completionCompleter = Completer(); int _runningTaskCount = 0; bool _started = false; late Queue<TInput> _taskQueue; Parallelizer._({ required this.workItems, required this.workFunction, required this.degreeOfParallelism, required this.skip, }); // 工厂构造 factory Parallelizer.create({ required Iterable<TInput> workItems, required Future<TOutput> Function(TInput, CancellationToken) workFunction, required int degreeOfParallelism, int skip = 0, }) { return Parallelizer._( workItems: workItems, workFunction: workFunction, degreeOfParallelism: degreeOfParallelism, skip: skip, ); } // 启动任务 Future<void> start() async { if (_started) return; _started = true; _taskQueue = Queue.of(workItems.skip(skip)); try { _scheduleNextTasks(); } catch (e) { _onError.add(e); if (!_completionCompleter.isCompleted) _completionCompleter.completeError(e); } } // 调度下一批任务,控制并发数 void _scheduleNextTasks() { while (_runningTaskCount < degreeOfParallelism && _taskQueue.isNotEmpty) { final item = _taskQueue.removeFirst(); _runningTaskCount++; final token = CancellationToken(); _runSingleTask(item, token).whenComplete(() { _runningTaskCount--; if (_taskQueue.isEmpty && _runningTaskCount == 0) { if (!_completionCompleter.isCompleted) { _completionCompleter.complete(); _onCompleted.add(null); } _closeStreams(); } else { _scheduleNextTasks(); } }); } } // 执行单个任务 Future<void> _runSingleTask(TInput item, CancellationToken token) async { try { final result = await workFunction(item, token); _onNewResult.add(ResultDetails(item, result)); } catch (e) { _onTaskError.add(ErrorDetails(item, e)); } } // 等待所有任务完成,支持取消 Future<void> waitCompletion(CancellationToken token) async { final cancelSub = token.onCancel.listen((_) { _taskQueue.clear(); if (!_completionCompleter.isCompleted) _completionCompleter.complete(); _onCompleted.add(null); _closeStreams(); }); await _completionCompleter.future; cancelSub.cancel(); } void _closeStreams() { _onNewResult.close(); _onCompleted.close(); _onError.close(); _onTaskError.close(); } } // 取消令牌实现 class CancellationToken { final StreamController _cancelController = StreamController.broadcast(); bool _isCancelled = false; bool get isCancelled => _isCancelled; Stream get onCancel => _cancelController.stream; void cancel() { if (_isCancelled) return; _isCancelled = true; _cancelController.add(null); _cancelController.close(); } } class CancellationTokenSource { final CancellationToken token = CancellationToken(); Timer? _timeoutTimer; void cancelAfter(Duration duration) { _timeoutTimer?.cancel(); _timeoutTimer = Timer(duration, () => token.cancel()); } void cancel() { _timeoutTimer?.cancel(); token.cancel(); } } // 示例主函数,和C#示例逻辑完全一致 Future<void> mainAsync(List<String> args) async { // 定义工作函数:输入int、取消令牌,返回Future<bool> Future<bool> parityCheck(int number, CancellationToken token) async { await Future.delayed(const Duration(milliseconds: 50)); if (token.isCancelled) throw Exception('Task cancelled'); return number % 2 == 0; } final parallelizer = Parallelizer<int, bool>.create( workItems: List.generate(100, (index) => index + 1), // 1-100的任务 workFunction: parityCheck, degreeOfParallelism: 5, // 最大并发5 skip: 0, // 不跳过前置任务 ); // 绑定事件 parallelizer.onNewResult.listen((details) { print('Got result ${details.result} from the parity check of ${details.item}'); }); parallelizer.onCompleted.listen((_) { print('All work completed!'); }); parallelizer.onTaskError.listen((details) { print('Got error ${details.exception} while processing the item ${details.item}'); }); parallelizer.onError.listen((ex) { print('Exception: $ex'); }); await parallelizer.start(); // 设置10秒超时取消 final cts = CancellationTokenSource(); cts.cancelAfter(const Duration(seconds: 10)); await parallelizer.waitCompletion(cts.token); }
适配说明
- 如果你执行的是IO密集型任务(网络请求、文件读写等),直接使用以上代码即可,Dart事件循环会自动调度,不会阻塞UI线程
- 如果你执行的是CPU密集型任务(大数据计算、音视频编解码等),只需要把
workFunction中的逻辑放到Isolate.run中执行即可,外层调度、事件回调、并发控制逻辑不需要任何修改 - 如果你需要进度统计,只需要在Parallelizer中新增已完成任务计数字段,在单个任务完成时更新计数,对外暴露流或者值即可,实现成本极低
内容的提问来源于stack exchange,提问作者Mian Ismaili
相关产品推荐
相关产品推荐

