Channel — 생산자에서 소비자로 값을 보내는 스레드 안전 큐

Channel — 생산자에서 소비자로 값을 보내는 스레드 안전 큐

여러 스레드가 값을 주고받을 때, "한쪽에서 보내면 다른 쪽에서 받는" 안전한 통로가 필요해요. Channel은 바로 그런 역할을 하는 스레드 안전한 큐입니다.

출처: Raku 공식 문서 — Channel

본문

class Channel {}

Channel은 하나 이상의 생산자(producer)에서 하나 이상의 소비자(consumer)로 객체를 보내는 데 도움을 주는 스레드 안전한 큐예요. 각 객체는 스케줄러가 고르는 소비자 하나에게만 도착합니다. 소비자와 생산자가 각각 하나뿐이면 객체의 순서가 보장되어요. Channel로 보내는 동작은 블로킹이 아닙니다.

my $c = Channel.new;
await (^10).map: {
    start {
        my $r = rand;
        sleep $r;
        $c.send($r);
    }
}
$c.close;
say $c.list;

더 많은 예시는 동시성 페이지에서 볼 수 있어요.

메서드

method send

method send(Channel:D: \item)

Channel에 항목을 큐에 넣어요. 이미 닫힌 Channel이면 X::Channel::SendOnClosed 타입의 예외를 던져요. 이 호출은 소비자가 객체를 가져가기를 기다리며 블록되지 않아요. 큐에 넣을 수 있는 항목 수에 제한이 없으니, 큐가 폭주하지 않게 주의해야 해요.

my $c = Channel.new;
$c.send(1);
$c.send([2, 3, 4, 5]);
$c.close;
say $c.list; # OUTPUT: «(1 [2 3 4 5])␤»

method receive

method receive(Channel:D:)

Channel에서 항목 하나를 받아 제거해요. 항목이 없으면 다른 스레드의 send를 기다리며 블록됩니다.

Channel이 닫혔는데 마지막 항목까지 이미 제거된 상태라면, 혹은 receive가 항목을 기다리는 동안 close가 호출되면 X::Channel::ReceiveOnClosed 타입의 예외를 던져요.

fail 메서드로 Channel이 "비정상( erratic)"으로 표시된 상태에서 마지막 항목이 제거되면, fail에 주어진 인자를 예외로 던져요.

예외를 던지지 않는 논블로킹 버전은 poll 메서드를 참고하세요.

my $c = Channel.new;
$c.send(1);
say $c.receive; # OUTPUT: «1␤»

method poll

method poll(Channel:D:)

Channel에서 항목 하나를 받아 제거해요. 항목이 없으면 기다리는 대신 Nil을 반환해요.

my $c = Channel.new;
Promise.in(2).then: { $c.close; }
^10 .map({ $c.send($_); });
loop {
    if $c.poll -> $item { $item.say };
    if $c.closed  { last };
    sleep 0.1;
}

Channel 닫힘과 실패에 제대로 대응하는 블로킹 버전은 receive 메서드를 참고하세요.

method close

method close(Channel:D:)

Channel을 정상적으로 닫아요. 이후의 send 호출은 X::Channel::SendOnClosed로 죽습니다. 이후의 .receive 호출은 이전에 보낸 남은 항목을 계속 빼낼 수 있지만, 큐가 비어 있으면 X::Channel::ReceiveOnClosed 예외를 던져요. @()로 배열을 만들거나 .list 메서드를 호출해서 Channel에서 Seq를 만들 수 있는데, 이 메서드들은 Channel이 닫힐 때까지 끝나지 않아요. whenever 블록도 닫힌 Channel에서 제대로 종료됩니다.

my $c = Channel.new;
$c.close;
$c.send(1);
CATCH { default { put .^name, ': ', .Str } };
# OUTPUT: «X::Channel::SendOnClosed: Cannot send a message on a closed channel␤»

참고로 어떤 예외가 던져지면 .close 호출이 막혀서 받는 쪽 스레드가 멈출 수 있어요. 그럴 땐 LEAVE 페이저를 써서 .close 호출을 강제하면 됩니다.

method list

method list(Channel:D:)

큐의 항목을 순회하면서, 순회할 때마다 그 항목을 큐에서 제거하는 Seq 기반의 리스트를 반환해요. close 메서드가 호출될 때만 종료할 수 있어요.

my $c = Channel.new; $c.send(1); $c.send(2);
$c.close;
say $c.list; # OUTPUT: «(1 2)␤»

method closed

method closed(Channel:D: --> Promise:D)

close 메서드 호출로 Channel이 닫히면 지켜질(kept) 프로미스를 반환해요.

my $c = Channel.new;
$c.closed.then({ say "It's closed!" });
$c.close;
sleep 1;

method fail

method fail(Channel:D: $error)

Channel을 닫고(즉 이후의 send 호출을 죽이고), 그 오류를 Channel의 마지막 요소로 던지도록 큐에 넣어요. receive 메서드는 그 오류를 예외로 던집니다. Channel이 이미 닫혔거나 .fail이 이미 호출됐다면 아무것도 하지 않아요.

my $c = Channel.new;
$c.fail("Bad error happens!");
$c.receive;
CATCH { default { put .^name, ': ', .Str } };
# OUTPUT: «X::AdHoc: Bad error happens!␤»

method Capture

method Capture(Channel:D: --> Capture:D)

invocant에 .List.Capture를 호출하는 것과 동일해요.

method Supply

method Supply(Channel:D:)

Channel에서 받은 값마다 하나씩 내보내는 on-demand Supply를 반환해요. Channel이 닫히면 Supplydone이 호출됩니다.

my $c = Channel.new;
my Supply $s1 = $c.Supply;
my Supply $s2 = $c.Supply;
$s1.tap(-> $v { say "First $v" });
$s2.tap(-> $v { say "Second $v" });
^10 .map({ $c.send($_) });
sleep 1;

이 메서드를 여러 번 호출하면 Channel의 값을 두고 경쟁하는 여러 Supply 인스턴스가 생겨요.

sub await

multi await(Channel:D)
multi await(*@)

하나 이상의 Channel 중 하나에라도 값이 준비될 때까지 기다렸다가 그 값들을 반환해요(Channel.receive를 호출해요). Promise에도 동작해요.

my $c = Channel.new;
Promise.in(1).then({$c.send(1)});
say await $c;

6.d부터는 기다리는 동안 스레드를 블록하지 않아요.