You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 11:05:16