構造化並行性
💡 お知らせ: このドキュメントはAIによって翻訳されています。表現に違和感がある場合は、原文(英語)を参照するか、翻訳にご協力ください。
Flix は、Go と Rust に着想を得た、チャネルとプロセスによる CSP スタイルの並行処理をサポートしています。
プロセスの生成
spawn キーワードを使って、プロセス(Process)を生成できます:
def main(): Unit \ IO = region rc {
spawn println("Hello from thread") @ rc;
println("Hello from main")
}
生成されたプロセスは、必ずリージョンに関連付けられます。リージョンは、それに関連付けられたすべてのプロセスが完了するまで終了しません:
def main(): Unit \ IO =
region r1 {
region r2 {
spawn println("Hello from r1") @ r1;
spawn println("Hello from r2") @ r2
};
println("r2 is now complete")
};
println("r1 is now complete")
これはつまり、Flix が構造化並行性をサポートしているということです。生成されたプロセスには、明確に定義された開始点と終了点があります。
チャネルによる通信
プロセス間で通信するには、チャネルを使用します。*チャネル(Channel)*を使うと、2つ以上のプロセスが互いにイミュータブルなメッセージを送り合うことで、データを交換できます。
チャネルには、*バッファ付き(Buffered)とバッファなし(Unbuffered)*の2つの種類があります。チャネルは必ずリージョンに関連付けられます。
バッファ付きチャネルは、作成時に設定されるサイズを持ち、その数だけメッセージを保持できます。満杯のバッファ付きチャネルにプロセスがメッセージを入れようとすると、そのプロセスは空きができるまでブロックされます。逆に、空のチャネルからプロセスがメッセージを取り出そうとすると、そのプロセスはチャネルにメッセージが入れられるまでブロックされます。
バッファなしチャネルは、サイズ0のバッファ付きチャネルのように動作します。取り出し(get)と書き込み(put)が成立するためには、送信側から受信側へメッセージが渡されるまで、両方のプロセスがランデブー(ブロック)しなければなりません。
チャネルを介してメッセージを送受信する例を示します:
def main(): Int32 \ {Chan, NonDet, IO} = region rc {
let (tx, rx) = Channel.unbuffered();
spawn Channel.send(42, tx) @ rc;
Channel.recv(rx)
}
ここで main 関数は、Sender チャネル tx と Receiver チャネル rx を返すバッファなしチャネルを作成し、send 関数を spawn して、チャネルからのメッセージを待ちます。
この例が示すように、チャネルは Sender(センダー) と Receiver(レシーバー) という2つのエンドポイントで構成されます。予想される通り、メッセージの送信は Sender からのみ、受信は Receiver からのみ行えます。
チャネルに対する select
select 式を使うと、複数のチャネルの集まりからメッセージを受信できます。例えば:
def meow(tx: Sender[String]): Unit \ Chan =
Channel.send("Meow!", tx)
def woof(tx: Sender[String]): Unit \ Chan =
Channel.send("Woof!", tx)
def main(): Unit \ {Chan, NonDet, IO} = region rc {
let (tx1, rx1) = Channel.buffered(1);
let (tx2, rx2) = Channel.buffered(1);
spawn meow(tx1) @ rc;
spawn woof(tx2) @ rc;
select {
case m <- recv(rx1) => m
case m <- recv(rx2) => m
} |> println
}
生産者・消費者(producer-consumer)やロードバランサーなど、多くの重要な並行処理パターンを select 式で表現できます。
デフォルトケース付きの select
場合によっては、メッセージが届くまでブロックして、永遠に待ち続けるかもしれない状況を避けたいことがあります。そのような場合は、メッセージがすぐに利用できないときに代わりのアクションを取りたくなります。これは、以下に示すように*デフォルトケース(Default case)*で実現できます:
def main(): String \ {Chan, NonDet} = region rc {
let (_, rx1) = Channel.buffered(1);
let (_, rx2) = Channel.buffered(1);
select {
case _ <- recv(rx1) => "one"
case _ <- recv(rx2) => "two"
case _ => "default"
}
}
ここでは、r1 にも r2 にもメッセージが送信されることはありません。select 式はすべてのケースを試し、どのチャネルも準備できていなければ、直ちにデフォルトケースを選択します。したがって、デフォルトケースを使うことで、select 式が永遠にブロックすることを防げます。
タイムアウト付きの select
デフォルトケースの代わりに、*ティッカー(Ticker)とタイマー(Timer)*を使って、select 式の中であらかじめ定められた時間だけ待つこともできます。
例えば、次のプログラムには、チャネルにメッセージを送るまでに1分かかる遅い関数がありますが、select 式は Channel.timeout を利用して、5 秒だけ待って諦めるようになっています:
def slow(tx: Sender[String]): Unit \ {Chan, NonDet, IO} =
let delay = Channel.timeout(60, Time.TimeUnit.Seconds);
Channel.recv(delay);
Channel.send("I am very slow", tx)
def main(): Unit \ {Chan, NonDet, IO} = region rc {
let (tx, rx) = Channel.buffered(1);
spawn slow(tx) @ rc;
let timeout = Channel.timeout(5, Time.TimeUnit.Seconds);
select {
case m <- recv(rx) => m
case _ <- recv(timeout) => "timeout"
} |> println
}
このプログラムは、5秒後に文字列 "timeout" を出力します。
チャネルのエフェクト
お気づきかもしれませんが、チャネルを使うと Chan と NonDet というエフェクトが現れます。
チャネルに対するあらゆる操作は Chan エフェクトを持ちます。このエフェクトは、プログラムがチャネルのグローバルな状態を変更または参照していることを表します。
チャネルに対する recv 操作は NonDet エフェクトを持ちます。これは、受け取る値が一般には非決定的であり、スレッドスケジューラの選択に依存するためです。2つのスレッドが同時にチャネルへ値を送信する準備ができていることがあり、どちらが先に送信できるかはスケジューラ次第です。