You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flutter如何实现并行多线程任务 复刻C# Parallelizer并行能力

Flutter 并行任务实现方案

Flutter 完全支持实现你需要的限定并发数的并行处理能力,不需要引入第三方依赖,基于Dart自带的异步、Stream、Isolate能力即可1:1对齐参考的功能逻辑。

概念对应关系

C#实现里的核心能力,在Dart中都有原生对应:

  • C# Task 对应Dart Future
  • 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 21:42:20