组件ThreadChannel::recv()

ThreadChannel::recv

(PHP 8.6+, True Async 1.0)

php
public ThreadChannel::recv(?Completable $cancellationToken = null): mixed

从通道接收下一个值。这是一个阻塞操作 — 如果通道中当前没有可用值,调用线程将被阻塞。

  • 对于有缓冲通道,如果缓冲区中至少有一个值,recv() 立即返回。 如果缓冲区为空,线程阻塞直到发送方放入值。
  • 对于无缓冲通道capacity = 0),recv() 阻塞直到另一个线程调用 send()

如果通道已关闭且缓冲区仍有值,这些值将正常返回。 一旦缓冲区排空且通道已关闭,recv() 将抛出 ThreadChannelException

接收到的值是原始值的深拷贝 — 对返回值的修改不会影响发送方的副本。

返回值

通道中的下一个值(mixed)。

错误

  • 如果通道已关闭且缓冲区为空,抛出 Async\ThreadChannelException

示例

示例 #1 接收工作线程产生的值

php
<?php

use Async\ThreadChannel;
use function Async\spawn;
use function Async\spawn_thread;
use function Async\await;

spawn(function() {
    $channel = new ThreadChannel(5);

    $worker = spawn_thread(function() use ($channel) {
        for ($i = 1; $i <= 5; $i++) {
            $channel->send($i * 10);
        }
        $channel->close();
    });

    // 接收所有值 — 缓冲区为空时阻塞
    try {
        while (true) {
            echo $channel->recv(), "\n";
        }
    } catch (\Async\ThreadChannelException) {
        echo "All values received\n";
    }

    await($worker);
});

示例 #2 消费者线程排空共享通道

php
<?php

use Async\ThreadChannel;
use function Async\spawn;
use function Async\spawn_thread;
use function Async\await;

spawn(function() {
    $channel = new ThreadChannel(20);

    // 生产者:从一个线程填充通道
    $producer = spawn_thread(function() use ($channel) {
        foreach (range('a', 'e') as $letter) {
            $channel->send($letter);
        }
        $channel->close();
    });

    // 消费者:从另一个线程排空通道
    $consumer = spawn_thread(function() use ($channel) {
        $collected = [];
        try {
            while (true) {
                $collected[] = $channel->recv();
            }
        } catch (\Async\ThreadChannelException) {
            // 缓冲区已排空且通道已关闭
        }
        return $collected;
    });

    await($producer);
    $result = await($consumer);
    echo implode(', ', $result), "\n"; // "a, b, c, d, e"
});

示例 #3 从无缓冲通道接收

php
<?php

use Async\ThreadChannel;
use function Async\spawn;
use function Async\spawn_thread;
use function Async\await;

spawn(function() {
    $channel = new ThreadChannel(); // 无缓冲

    $sender = spawn_thread(function() use ($channel) {
        // 在这里阻塞,直到主线程调用 recv()
        $channel->send(['task' => 'compress', 'file' => '/tmp/data.bin']);
    });

    // 主协程(线程)调用 recv() — 解除发送方的阻塞
    $task = $channel->recv();
    echo "Got task: {$task['task']} on {$task['file']}\n";

    await($sender);
});

参见