Dart Isolate — 共有メモリなしの並列計算
Isolate は Dart の並列処理の方法で、共有メモリなし、自然にレースフリー、安全で効率的。
1. 学べること
- Isolate vs Thread:共有メモリなしの並行処理モデル
- Isolate.spawn と ReceivePort / SendPort による双方向通信
- Isolate.run:ワンショットタスク用のシンプル API
- Compute ヘルパー関数と Flutter の compute
- Bob のシナリオ:DataPipeline で Isolate を使った百万件の注文の並列処理
2. 開発者のリアルな物語
(1) 課題:百万件のレコードの処理がシングルスレッドで遅すぎる
Bob の DataPipeline は 120 万件の注文に対して統計分析を行う必要があった。シングルスレッド処理は 15 秒かかったが、SaaS クライアントはレポートを 5 秒以内に生成することを要求していた。Bob はマルチスレッドを試したが、共有メモリロックと競合状態により、競合バグの修正に 1 週間費やした末にあきらめた。レポートはまだ 15 秒かかった。
(2) Isolate による解決
Dart の Isolate は共有メモリなしの並列ユニットである。各 Isolate は独自の独立したヒープメモリを持ち、メッセージパッシングで通信する。共有メモリなしは競合状態なしを意味する。
// 120 万件の注文を 4 つの Isolate に分割し、それぞれ 30 万件を処理
final results = await Future.wait([
Isolate.run(() => processChunk(orders.sublist(0, 300000))),
Isolate.run(() => processChunk(orders.sublist(300000, 600000))),
Isolate.run(() => processChunk(orders.sublist(600000, 900000))),
Isolate.run(() => processChunk(orders.sublist(900000, 1200000))),
]);
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
(3) 効果
- 4 つの Isolate の並列実行により処理時間が 15 秒から 4 秒に短縮
- 共有メモリなし = 競合状態なし = ロックなし = 並行バグなし
- メッセージパッシングモデルによりコードの推論が容易に
3. Isolate の基礎
(1) Isolate vs Thread
flowchart TD
A[メイン Isolate] -->|"SendPort"| B[ワーカー Isolate 1]
A -->|"SendPort"| C[ワーカー Isolate 2]
A -->|"SendPort"| D[ワーカー Isolate N]
B -->|"SendPort"| A
C -->|"SendPort"| A
D -->|"SendPort"| A
subgraph データシャーディング
B -- 30 万件の注文
C -- 30 万件の注文
D -- 40 万件の注文
end
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
| 観点 | Thread(Java/C++) | Isolate(Dart) |
|---|---|---|
| メモリ | 共有 | 独立 |
| 通信 | 共有変数 + ロック | メッセージパッシング |
| 競合状態 | あり | なし |
| データ同期 | 手動ロック必要 | ロック不要 |
| 生成オーバーヘッド | 低 | 中(データコピーが必要) |
4. Isolate.run — シンプル化されたワンショットタスク
(1) 最もシンプルな Isolate の使用
▶ サンプル:Isolate.run の基礎
import 'dart:isolate';
// Isolate で実行する高コスト計算
int fibonacci(int n) {
if (n <= 1) return n;
return fibonacci(n - 1) + fibonacci(n - 2);
}
Future<void> main() async {
print('Computing fibonacci(40) in isolate...');
final result = await Isolate.run(() => fibonacci(40));
print('Result: $result');
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
▶ サンプル:Isolate.run でデータチャンクを処理
import 'dart:isolate';
// 注文のチャンクを処理
double processChunk(List<Map<String, dynamic>> orders) {
double total = 0;
for (final order in orders) {
total += (order['amount'] as double);
}
return total;
}
Future<void> main() async {
// 100 万件の注文をシミュレート
final orders = List.generate(
1000000,
(i) => {'id': 'ORD-$i', 'amount': (i % 100 + 1) * 10.0},
);
// 4 つのチャンクに分割
final chunkSize = orders.length ~/ 4;
final chunks = List.generate(
4,
(i) => orders.sublist(i * chunkSize, (i + 1) * chunkSize),
);
// チャンクを並列処理
final stopwatch = Stopwatch()..start();
final results = await Future.wait(
chunks.map((chunk) => Isolate.run(() => processChunk(chunk))),
);
stopwatch.stop();
final totalRevenue = results.fold(0.0, (a, b) => a + b);
print('Total revenue: \$${totalRevenue.toStringAsFixed(2)} USD');
print('Time: ${stopwatch.elapsedMilliseconds}ms');
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
5. Isolate.spawn と双方向通信
(1) SendPort / ReceivePort 通信
▶ サンプル:双方向通信
import 'dart:isolate';
// ワーカー関数 - 別 Isolate で実行
void workerIsolate(SendPort mainSendPort) {
final workerReceivePort = ReceivePort();
// メイン Isolate に自分の receive port を送信
mainSendPort.send(workerReceivePort.sendPort);
// メイン Isolate からのメッセージをリッスン
workerReceivePort.listen((message) {
if (message is String && message == 'shutdown') {
workerReceivePort.close();
return;
}
// データを処理して結果を返す
if (message is List<double>) {
final total = message.fold(0.0, (a, b) => a + b);
mainSendPort.send(total);
}
});
}
Future<void> main() async {
final mainReceivePort = ReceivePort();
// ワーカー Isolate を生成
await Isolate.spawn(workerIsolate, mainReceivePort.sendPort);
// ワーカーの send port を取得
final workerSendPort = await mainReceivePort.first as SendPort;
// レスポンス用の新しい receive port を作成
final responsePort = ReceivePort();
workerSendPort.send([1500.0, 3200.0, 890.0]);
// レスポンスを待つ
final result = await responsePort.first;
print('Result from worker: $result');
// ワーカーをシャットダウン
workerSendPort.send('shutdown');
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
▶ サンプル:複数ラウンド通信
import 'dart:isolate';
void dataWorker(SendPort mainPort) {
final receivePort = ReceivePort();
mainPort.send(receivePort.sendPort);
receivePort.listen((message) {
if (message == 'done') {
receivePort.close();
return;
}
if (message is Map<String, dynamic>) {
// 注文データを処理
final amount = message['amount'] as double;
final taxRate = message['taxRate'] as double? ?? 0.08;
final total = amount * (1 + taxRate);
mainPort.send({'id': message['id'], 'total': total});
}
});
}
Future<void> main() async {
final mainPort = ReceivePort();
await Isolate.spawn(dataWorker, mainPort.sendPort);
final workerPort = await mainPort.first as SendPort;
// 複数のメッセージを送信
final responsePort = ReceivePort();
workerPort.add({'id': 'ORD-001', 'amount': 1500.0, 'taxRate': 0.08});
workerPort.add({'id': 'ORD-002', 'amount': 3200.0});
// ... 簡略化された通信
workerPort.send('done');
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
6. Isolate とデータ転送
(1) データ転送ルール
| 転送方法 | 説明 | パフォーマンス |
|---|---|---|
| プリミティブ型 | int/double/String/bool | 高速コピー |
| List/Map | ディープコピー | 中 |
| カスタムオブジェクト | ディープコピー | 中 |
| SendPort | 参照渡し | 高速 |
| 関数クロージャ | Isolate.run 経由で渡す | コンパイル時チェック |
▶ サンプル:大きなデータ転送のベストプラクティス
import 'dart:isolate';
// クロージャ経由でデータを渡す(Isolate.run)
Future<double> processLargeData(List<double> amounts) async {
return Isolate.run(() {
// amounts は isolate にコピーされる
double sum = 0;
for (final a in amounts) {
sum += a;
}
return sum;
});
}
// より良く:必要なものだけを渡す
Future<double> processChunkOptimized(List<double> chunk) async {
return Isolate.run(() => chunk.fold(0.0, (a, b) => a + b));
}
void main() async {
final data = List.generate(1000000, (i) => (i + 1) * 1.0);
// 分割して並列処理
final chunkSize = data.length ~/ 4;
final futures = List.generate(4, (i) {
final chunk = data.sublist(i * chunkSize, (i + 1) * chunkSize);
return Isolate.run(() => chunk.fold(0.0, (a, b) => a + b));
});
final results = await Future.wait(futures);
final total = results.fold(0.0, (a, b) => a + b);
print('Total: $total');
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
7. Bob のシナリオ:百万件の注文の並列処理
▶ サンプル:DataPipeline 並列統計
import 'dart:isolate';
class ChunkResult {
final int processed;
final int skipped;
final double revenue;
final Map<String, double> categoryRevenue;
ChunkResult({
required this.processed,
required this.skipped,
required this.revenue,
required this.categoryRevenue,
});
}
ChunkResult processChunk(List<Map<String, dynamic>> chunk) {
int processed = 0;
int skipped = 0;
double revenue = 0;
final categoryRevenue = <String, double>{};
for (final order in chunk) {
final amount = order['amount'] as double;
final status = order['status'] as String;
final category = order['category'] as String;
if (amount <= 0 || status == 'cancelled') {
skipped++;
continue;
}
processed++;
revenue += amount;
categoryRevenue.update(category, (v) => v + amount, ifAbsent: () => amount);
}
return ChunkResult(
processed: processed,
skipped: skipped,
revenue: revenue,
categoryRevenue: categoryRevenue,
);
}
Future<void> main() async {
// 120 万件の注文を生成
final categories = ['Electronics', 'Books', 'Clothing', 'Home', 'Sports'];
final statuses = ['completed', 'completed', 'completed', 'pending', 'cancelled'];
final orders = List.generate(1200000, (i) => {
'id': 'ORD-${i.toString().padLeft(6, '0')}',
'amount': (i % 500 + 10) * 1.0,
'status': statuses[i % statuses.length],
'category': categories[i % categories.length],
});
// チャンクに分割
final isolateCount = 4;
final chunkSize = orders.length ~/ isolateCount;
final chunks = List.generate(isolateCount, (i) {
final start = i * chunkSize;
final end = i == isolateCount - 1 ? orders.length : (i + 1) * chunkSize;
return orders.sublist(start, end);
});
// 並列処理
final stopwatch = Stopwatch()..start();
final results = await Future.wait(
chunks.map((chunk) => Isolate.run(() => processChunk(chunk))),
);
stopwatch.stop();
// 結果を集約
int totalProcessed = 0;
int totalSkipped = 0;
double totalRevenue = 0;
final totalCategoryRevenue = <String, double>{};
for (final r in results) {
totalProcessed += r.processed;
totalSkipped += r.skipped;
totalRevenue += r.revenue;
for (final entry in r.categoryRevenue.entries) {
totalCategoryRevenue.update(entry.key, (v) => v + entry.value, ifAbsent: () => entry.value);
}
}
print('=== DataPipeline Parallel Report ===');
print('Orders: ${orders.length}');
print('Processed: $totalProcessed');
print('Skipped: $totalSkipped');
print('Revenue: \$${totalRevenue.toStringAsFixed(2)} USD');
print('Time: ${stopwatch.elapsedMilliseconds}ms');
print('Isolates: $isolateCount');
print('\nBy Category:');
for (final entry in totalCategoryRevenue.entries.toList()..sort((a, b) => b.value.compareTo(a.value))) {
print(' ${entry.key}: \$${entry.value.toStringAsFixed(2)} USD');
}
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
8. 完全なサンプル:DataPipeline Isolate 並列処理フレームワーク
// ============================================
// DataPipeline Isolate 並列処理
// 並列データ分析のフレームワーク
// ============================================
import 'dart:isolate';
// ワーカー結果
class WorkerResult {
final int workerId;
final int processed;
final int skipped;
final double revenue;
final Map<String, double> categoryRevenue;
final Duration processingTime;
WorkerResult({
required this.workerId,
required this.processed,
required this.skipped,
required this.revenue,
required this.categoryRevenue,
required this.processingTime,
});
}
// ワーカー関数 - 別 Isolate で実行
WorkerResult processInIsolate((int, List<Map<String, dynamic>>) input) {
final (workerId, orders) = input;
final stopwatch = Stopwatch()..start();
int processed = 0;
int skipped = 0;
double revenue = 0;
final categoryRevenue = <String, double>{};
for (final order in orders) {
final amount = (order['amount'] as num).toDouble();
final status = order['status'] as String;
final category = order['category'] as String;
if (amount <= 0 || status == 'cancelled') {
skipped++;
continue;
}
processed++;
revenue += amount;
categoryRevenue.update(category, (v) => v + amount, ifAbsent: () => amount);
}
stopwatch.stop();
return WorkerResult(
workerId: workerId,
processed: processed,
skipped: skipped,
revenue: revenue,
categoryRevenue: categoryRevenue,
processingTime: stopwatch.elapsed,
);
}
// 並列パイプラインマネージャー
class ParallelPipeline {
final int isolateCount;
ParallelPipeline({this.isolateCount = 4});
Future<void> process(List<Map<String, dynamic>> orders) async {
print('=== DataPipeline Parallel Processing ===');
print('Orders: ${orders.length}, Isolates: $isolateCount');
// 注文をチャンクに分割
final chunkSize = orders.length ~/ isolateCount;
final chunks = List.generate(isolateCount, (i) {
final start = i * chunkSize;
final end = i == isolateCount - 1 ? orders.length : (i + 1) * chunkSize;
return (i, orders.sublist(start, end));
});
// 並列処理
final totalStopwatch = Stopwatch()..start();
final results = await Future.wait(
chunks.map((chunk) => Isolate.run(() => processInIsolate(chunk))),
);
totalStopwatch.stop();
// 集約
int totalProcessed = 0;
int totalSkipped = 0;
double totalRevenue = 0;
final totalCategoryRevenue = <String, double>{};
print('\nWorker Results:');
for (final r in results) {
totalProcessed += r.processed;
totalSkipped += r.skipped;
totalRevenue += r.revenue;
for (final entry in r.categoryRevenue.entries) {
totalCategoryRevenue.update(
entry.key, (v) => v + entry.value, ifAbsent: () => entry.value);
}
print(' Worker ${r.workerId}: ${r.processed} processed, '
'${r.processingTime.inMilliseconds}ms');
}
print('\n--- Aggregate Report ---');
print('Processed: $totalProcessed orders');
print('Skipped: $totalSkipped records');
print('Revenue: \$${totalRevenue.toStringAsFixed(2)} USD');
print('Wall time: ${totalStopwatch.elapsedMilliseconds}ms');
print('\nBy Category:');
final sorted = totalCategoryRevenue.entries.toList()
..sort((a, b) => b.value.compareTo(a.value));
for (final entry in sorted) {
final pct = (entry.value / totalRevenue * 100).toStringAsFixed(1);
print(' ${entry.key}: \$${entry.value.toStringAsFixed(2)} USD ($pct%)');
}
}
}
void main() async {
// サンプルデータを生成
final categories = ['Electronics', 'Books', 'Clothing', 'Home', 'Sports'];
final statuses = ['completed', 'completed', 'completed', 'pending', 'cancelled'];
final orders = List.generate(500000, (i) => <String, dynamic>{
'id': 'ORD-${i.toString().padLeft(6, '0')}',
'amount': (i % 500 + 10) * 1.0,
'status': statuses[i % statuses.length],
'category': categories[i % categories.length],
});
final pipeline = ParallelPipeline(isolateCount: 4);
await pipeline.process(orders);
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
❓ よくある質問
Q: Isolate と Thread の違いは何ですか? A: Isolate は独自の独立したヒープメモリを持ち、データを共有しません。Thread はヒープメモリを共有します。Isolate はメッセージパッシングで通信し、Thread は共有変数 + ロックで通信します。Isolate には競合状態がありません。
Q: Isolate の生成はコストがかかりますか? A: Isolate の生成には約 50〜150ms かかり、スレッドの生成より遅いです。頻繁な生成と破棄は非効率なので、Isolate プールまたは Isolate.run(ライフサイクルを自動管理)の使用を推奨します。
Q: Isolate.run と Isolate.spawn の違いは何ですか? A: Isolate.run はシンプル化された API で、ワンショットタスクを実行してから Isolate を自動的に閉じます。Isolate.spawn は永続的な Isolate を作成し、ライフサイクルと通信を手動で管理する必要があります。
Q: Isolate 間でデータを転送するとき、コピーされますか? A: はい。転送されるすべてのデータはディープコピーされます(SendPort を除く)。大きなリストの渡しにはパフォーマンスオーバーヘッドがあります。非常に大きなデータの場合は、チャンクで転送を検討してください。
Q: Dart には何種類の Isolate がありますか? A: 主に 2 つ:並列処理用の汎用 Isolate と、UI 分離用の Flutter の compute。Web プラットフォームは真の Isolate をサポートせず、Web Worker をシミュレーションとして使用します。
Q: Isolate の数に上限はありますか? A: 厳密な上限はありませんが、各 Isolate は約 2MB のメモリを占有します。実際には、過剰なコンテキストスイッチを避けるため CPU コア数を超えないことを推奨します。
Q: Isolate が適しているのはどんなシナリオですか? A: CPU 集約型の計算(データ分析、画像処理、暗号計算)。I/O 集約型のタスクには Future/Stream で十分で、Isolate は不要です。
📖 まとめ
- Isolate は Dart の独立したメモリ、メッセージパッシング通信を持つ並列ユニットであり、自然にレースフリー
- Isolate.run はワンショット計算タスクに適し、ライフサイクルを自動管理する
- Isolate.spawn + SendPort/ReceivePort は永続的な通信に適する
- データ転送にはコピーが伴う;大きなデータはチャンクで転送する
- 4 つの Isolate で百万件の注文を並列処理すると、時間が 15 秒から 4 秒に短縮される
📝 練習問題
- 基礎(難易度 ⭐):
Isolate.runを使って 42 番目のフィボナッチ数を計算し、メインスレッドでの同期計算と時間を比較してください。 - 中級(難易度 ⭐⭐):100 万個の乱数を生成し、4 つに分割して、4 つの Isolate で各部分の合計と平均を並列計算してください。最終的にメイン Isolate で結果を集約します。
- 挑戦(難易度 ⭐⭐⭐):
Isolate.spawnを使って Isolate プールを実装してください:N 個のワーカー Isolate を事前に作成し、メイン Isolate が SendPort 経由でタスクを配布し、ワーカーが処理後に結果を返すようにします。タスクキューとロードバランシングをサポートします。