amphp Pipeline流式返回结果:自定义系统间通信实现求助
问题描述
我需要实现一个能处理迭代器中数万条数据的Pipeline,由于处理耗时较长,希望结果一生成就流式返回给客户端。这不是Web服务器场景,而是我控制的两个系统之间的通信。
我尝试用Pipeline::tap()来处理输出,但目前最大的问题是流的构建:原本想返回一个以WriteableStream为响应体的Response,持续写入直到所有数据处理完成,但Response仅接受ReadableStream作为响应体。如果用sockets实现会丢失Request/Response的实用抽象,我不想这么做,但http-server本身基于sockets构建,应该有可行的实现方式。
以下是我目前的代码,用了tap()但不确定是否是其预期用法,还没找到在Pipeline处理完成后发送结束标记的方法:
$pipeline = Pipeline::fromIterable($bigList) ->concurrent(Router::MAX_CONCURRENT_PROCESSES_PER_REQUEST) ->unordered() ->map(fn (ListItem $item) => $this->getResultForItem($item)) ->filter(fn (?Result $r) => $r instanceof Result); // if there's a stream, don't do anything here that will block until the pipeline is finished. if (!is_null($stream)) { $streamStarted = false; $pipeline->tap(function (Result $r) use ($stream, &$streamStarted) { $stream->write($streamStarted ? ',' : '['); $streamStarted = true; $stream->write(json_encode($r)); }); // something here to send the `end()` message when the pipeline has been completed for all items //$future = async() } else { foreach($pipeline as $profile) { $indexed[$profile->dsid] = $profile; } }
解决方案
核心思路:用ReadableStream包装Pipeline输出
既然Response只接受ReadableStream,我们可以把Pipeline的处理过程转换成可读流,既保留Request/Response的抽象,又实现流式返回。
具体实现
方法一:将Pipeline转为可读流(推荐)
通过自定义可读流,把Pipeline的输出作为流的chunk发送,完美适配Response的要求:
use React\Stream\ReadableStream; // 构建Pipeline $pipeline = Pipeline::fromIterable($bigList) ->concurrent(Router::MAX_CONCURRENT_PROCESSES_PER_REQUEST) ->unordered() ->map(fn (ListItem $item) => $this->getResultForItem($item)) ->filter(fn (?Result $r) => $r instanceof Result); if (!is_null($stream)) { // 自定义可读流,绑定Pipeline作为数据源 $readableStream = new class($pipeline) extends ReadableStream { private $pipeline; private $isFirstChunk = true; public function __construct($pipeline) { $this->pipeline = $pipeline; } public function resume() { parent::resume(); $this->processPipeline(); } private function processPipeline() { foreach ($this->pipeline as $result) { if ($this->isFirstChunk) { // 发送JSON数组开头 $this->emit('data', ['[']); $this->isFirstChunk = false; } else { // 发送结果分隔符 $this->emit('data', [',']); } // 发送单个结果的JSON $this->emit('data', [json_encode($result)]); } // 处理完成后关闭流并补全JSON数组 if (!$this->isFirstChunk) { $this->emit('data', [']']); } $this->close(); } }; // 用可读流构建响应并返回 $response = new Response(200, ['Content-Type' => 'application/json'], $readableStream); // 此处根据你的通信框架逻辑发送response给客户端 } else { $indexed = []; foreach($pipeline as $profile) { $indexed[$profile->dsid] = $profile; } }
方法二:异步驱动可写流写入
如果必须使用可写流,可将Pipeline处理放入异步任务,避免阻塞响应返回:
if (!is_null($stream)) { // 异步执行Pipeline处理,后台写入流 async(function() use ($pipeline, $stream) { $streamStarted = false; foreach ($pipeline as $r) { $stream->write($streamStarted ? ',' : '['); $streamStarted = true; $stream->write(json_encode($r)); } if ($streamStarted) { $stream->write(']'); } $stream->end(); }); // 立即返回绑定了可写流的响应(需框架支持异步响应模式) $response = new Response(200, ['Content-Type' => 'application/json'], $stream); }
关键说明
- 方法一完全符合
Response的接口要求,无需额外框架支持,是最稳妥的实现方式。 - 两种方法都实现了“边处理边返回”的流式效果,避免客户端等待全部处理完成。
- 处理JSON格式时,确保输出是合法的JSON数组,避免客户端解析报错。
内容的提问来源于stack exchange,提问作者Jerry
相关产品推荐
相关产品推荐

