gRPC Streaming: потоковая передача данных

gRPC Streaming: потоковая передача данных

gRPC поддерживает четыре типа RPC: Unary (один запрос - один ответ), Server streaming (один запрос - много ответов), Client streaming (много запросов - один ответ), Bidirectional (много в обе стороны). Streaming - мощный инструмент для real-time приложений, файлов и подписок на события.

Четыре типа RPC: unary, server streaming, client streaming, bidirectional

Server Streaming

Один запрос от клиента, много сообщений от сервера. Подходит для логов, прогресса, real-time-уведомлений.

service LogService {
  rpc StreamLogs(StreamLogsRequest) returns (stream LogEntry);
}
func (s *logServer) StreamLogs(req *pb.StreamLogsRequest, stream pb.LogService_StreamLogsServer) error {
    for i := 0; i < 100; i++ {
        select {
        case <-stream.Context().Done():
            return stream.Context().Err()  // клиент отменил
        default:
        }

        if err := stream.Send(&pb.LogEntry{
            Timestamp: timestamppb.Now(),
            Message:   fmt.Sprintf("Log %d", i),
        }); err != nil {
            return err
        }
        time.Sleep(100 * time.Millisecond)
    }
    return nil
}
<?php
declare(strict_types=1);

namespace App\Grpc;

use Generator;
use Google\Protobuf\Timestamp;
use Log\LogEntry;
use Log\LogServiceInterface;
use Log\StreamLogsRequest;
use Spiral\RoadRunner\GRPC\ContextInterface;

final class LogService implements LogServiceInterface
{
    public function StreamLogs(ContextInterface $ctx, StreamLogsRequest $in): Generator
    {
        for ($i = 0; $i < 100; $i++) {
            $entry = (new LogEntry())
                ->setTimestamp((new Timestamp())->setSeconds(time()))
                ->setMessage(sprintf('Log %d', $i));

            yield $entry;

            usleep(100_000); // 100ms
        }
    }
}

В PHP полноценный server-streaming доступен только через RoadRunner (spiral/roadrunner-grpc). Метод возвращает \Generator, через yield отдаются сообщения. PHP-FPM на это не подходит - его модель «запрос-ответ» несовместима с долгоживущим стримом.

Cancellation - на уровне рантайма: если клиент отвалился, RoadRunner прерывает генератор автоматически (yield бросит исключение, которое можно поймать).

Клиент:

stream, _ := client.StreamLogs(ctx, &pb.StreamLogsRequest{})
for {
    entry, err := stream.Recv()
    if err == io.EOF {
        break
    }
    if err != nil {
        return err
    }
    fmt.Printf("[%s] %s\n", entry.Level, entry.Message)
}

Server закрывает стрим, возвращая nil из метода. Клиент видит это как io.EOF на Recv.

<?php
declare(strict_types=1);

use Log\LogServiceClient;
use Log\StreamLogsRequest;

final readonly class LogConsumer
{
    public function __construct(
        private LogServiceClient $client,
    ) {}

    public function tail(): void
    {
        $call = $this->client->StreamLogs(new StreamLogsRequest());

        foreach ($call->responses() as $entry) {
            printf('[%s] %s%s', $entry->getLevel(), $entry->getMessage(), "\n");
        }

        $status = $call->getStatus();
        if ($status->code !== \Grpc\STATUS_OK) {
            throw new \RuntimeException('stream failed: ' . $status->details);
        }
    }
}

PHP-клиент работает через ext-grpc: BaseStub отдаёт объект ServerStreamingCall, у которого есть метод responses() (Generator). Конец стрима - выход из цикла.

Client Streaming

Много сообщений от клиента, один ответ. Идеально для загрузки файлов, batch-импорта.

service FileService {
  rpc Upload(stream FileChunk) returns (UploadResponse);
}
func (s *fileServer) Upload(stream pb.FileService_UploadServer) error {
    var totalBytes int64
    h := sha256.New()

    for {
        chunk, err := stream.Recv()
        if err == io.EOF {
            return stream.SendAndClose(&pb.UploadResponse{
                Size:     totalBytes,
                Checksum: hex.EncodeToString(h.Sum(nil)),
            })
        }
        if err != nil {
            return err
        }
        totalBytes += int64(len(chunk.Data))
        h.Write(chunk.Data)
    }
}

В PHP client-streaming используется реже unary и server-streaming. На клиенте BaseStub возвращает ClientStreamingCall, в который складываются чанки через write(). На сервере RoadRunner отдаёт ServerStream с методом recv(). Для крупных загрузок (видео, бэкапы) разумной альтернативой остаётся обычный multipart-upload через Symfony:

<?php
declare(strict_types=1);

use File\FileChunk;
use File\FileServiceClient;
use File\UploadResponse;

final readonly class FileUploader
{
    public function __construct(
        private FileServiceClient $client,
    ) {}

    /** @param iterable<string> $chunks */
    public function upload(iterable $chunks): UploadResponse
    {
        $call = $this->client->Upload();

        foreach ($chunks as $bytes) {
            $call->write((new FileChunk())->setData($bytes));
        }

        [$response, $status] = $call->wait();
        if ($status->code !== \Grpc\STATUS_OK) {
            throw new \RuntimeException('upload failed: ' . $status->details);
        }
        return $response;
    }
}

Сервер на RoadRunner:

<?php
declare(strict_types=1);

use File\FileChunk;
use File\UploadResponse;
use Spiral\RoadRunner\GRPC\ContextInterface;
use Spiral\RoadRunner\GRPC\ServerStream;

final class FileService implements FileServiceInterface
{
    public function Upload(ContextInterface $ctx, ServerStream $stream): UploadResponse
    {
        $totalBytes = 0;
        $hash = hash_init('sha256');

        while (($chunk = $stream->recv(FileChunk::class)) !== null) {
            $data = $chunk->getData();
            $totalBytes += strlen($data);
            hash_update($hash, $data);
        }

        return (new UploadResponse())
            ->setSize($totalBytes)
            ->setChecksum(hash_final($hash));
    }
}

SendAndClose отправляет ответ и закрывает стрим. До этого момента сервер не отвечает - он накапливает.

Bidirectional Streaming

Полный дуплекс: обе стороны шлют независимо. Подходит для чата, real-time-игр, торгов.

service ChatService {
  rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}
func (s *chatServer) Chat(stream pb.ChatService_ChatServer) error {
    for {
        msg, err := stream.Recv()
        if err == io.EOF {
            return nil
        }
        if err != nil {
            return err
        }

        // Эхо всем подписчикам или просто обратно
        if err := stream.Send(&pb.ChatMessage{
            User:    "bot",
            Content: "Echo: " + msg.Content,
        }); err != nil {
            return err
        }
    }
}
<?php
declare(strict_types=1);

use Chat\ChatMessage;
use Spiral\RoadRunner\GRPC\BidiStream;
use Spiral\RoadRunner\GRPC\ContextInterface;

final class ChatService implements ChatServiceInterface
{
    public function Chat(ContextInterface $ctx, BidiStream $stream): void
    {
        while (($msg = $stream->recv(ChatMessage::class)) !== null) {
            $reply = (new ChatMessage())
                ->setUser('bot')
                ->setContent('Echo: ' . $msg->getContent());

            $stream->send($reply);
        }
    }
}

Bidi - самый сложный в PHP сценарий. На RoadRunner работает: BidiStream даёт recv() и send(). Для настоящего полного дуплекса (параллельное чтение и запись) нужен Swoole c корутинами - RoadRunner обрабатывает один поток последовательно.

Send и Recv работают независимо: можно одновременно читать в одной горутине и писать в другой. В PHP-RoadRunner это последовательно в рамках одного воркера; для параллельного reader/writer пайплайна нужен Swoole с Coroutine\go().

Flow control и backpressure

HTTP/2 имеет встроенный flow control: каждый поток имеет окно (window size). Sender не может отправить больше, чем позволяет окно. Это даёт автоматическое backpressure: если клиент медленно читает, сервер замедляется тоже.

// Большое окно для streaming с большим throughput
grpcServer := grpc.NewServer(
    grpc.InitialWindowSize(1024*1024),       // 1MB stream window
    grpc.InitialConnWindowSize(2*1024*1024), // 2MB connection window
)

В PHP размеры окон HTTP/2 не настраиваются из приложения. У RoadRunner есть параметры в .rr.yaml, для тонкой настройки придётся компилировать с другим Go-конфигом. Обычно дефолтов хватает:

# .rr.yaml
grpc:
    listen: tcp://0.0.0.0:50051
    proto:
        - proto/log.proto
    max_send_msg_size: 4   # MB
    max_recv_msg_size: 16  # MB
    max_connection_age: 0s
    max_connection_idle: 0s
    pool:
        num_workers: 8
        max_jobs: 1000

При типовом backpressure-сценарии (быстрый sender, медленный receiver) sender естественным образом затормозится. Не нужно вручную управлять очередями - HTTP/2 делает это в транспорте.

Stream lifecycle и errors

Стрим может закрыться четырьмя способами:

  • Нормальное закрытие (return nil из server-метода или EOF на client side)
  • Ошибка с status (return status.Error с кодом - клиент получит ошибку)
  • Cancellation (клиент отменил context - сервер видит через ctx.Done())
  • Network failure (потеря соединения - обе стороны получают ошибку)
for {
    msg, err := stream.Recv()
    switch {
    case err == nil:
        // обработка
    case err == io.EOF:
        return nil
    case status.Code(err) == codes.Canceled:
        // клиент отменил
        return err
    default:
        slog.Error("stream recv", "err", err)
        return err
    }
}
<?php
declare(strict_types=1);

use Spiral\RoadRunner\GRPC\Exception\GRPCException;
use Spiral\RoadRunner\GRPC\StatusCode;

while (true) {
    try {
        $msg = $stream->recv(ChatMessage::class);
        if ($msg === null) {
            return; // EOF, клиент закрыл send-side
        }
        // обработка $msg
    } catch (GRPCException $e) {
        match ($e->getCode()) {
            StatusCode::CANCELLED => $this->logger->info('client cancelled'),
            default => $this->logger->error('stream recv', ['err' => $e->getMessage()]),
        };
        return;
    }
}

В RoadRunner recv() возвращает null при EOF, бросает GRPCException при отмене или сетевой ошибке.

Concurrent send/receive

В bidirectional часто запускают goroutines для send и receive параллельно:

go func() {
    for {
        msg := <-outbound
        if err := stream.Send(msg); err != nil {
            return
        }
    }
}()

for {
    msg, err := stream.Recv()
    if err != nil {
        return err
    }
    inbound <- msg
}

stream.Send и stream.Recv thread-safe раздельно, но нельзя делать одновременные Send из двух горутин - нужен mutex или единый канал для исходящих.

<?php
declare(strict_types=1);

use Swoole\Coroutine;
use Swoole\Coroutine\Channel;

$outbound = new Channel(64);
$inbound  = new Channel(64);

// «Goroutine» на send
Coroutine::create(function () use ($stream, $outbound): void {
    while (($msg = $outbound->pop()) !== false) {
        try {
            $stream->send($msg);
        } catch (\Throwable $e) {
            return;
        }
    }
});

// Главная корутина читает
while (true) {
    try {
        $msg = $stream->recv(ChatMessage::class);
        if ($msg === null) {
            break;
        }
        $inbound->push($msg);
    } catch (\Throwable $e) {
        break;
    }
}

PHP-эквивалент - через Swoole-корутины (RoadRunner это сделать не сможет, там нет shared-памяти между «горутинами»).

Когда использовать какой тип

**Unary** - большинство CRUD-операций. **Server streaming** - логи, real-time-уведомления, прогресс долгой операции. **Client streaming** - загрузка файлов, batch-импорт. **Bidirectional** - чат, real-time-игры, IoT-телеметрия с командами.

Decision matrix

Сценарий                                  Тип RPC               Почему
──────────────────────────                ───────────────       ───────────────────────────
CRUD: создать, получить, обновить         Unary                 один запрос → один ответ;
                                                                 классика, не усложняй
Список из 10K записей с фильтрами          Unary с pagination   проще, чем server stream
                                                                 (offset/cursor)
Список без верхней границы (логи          Server streaming      сервер сам решает, когда
realtime, тикер цен, прогресс долгой                            закончить; backpressure
операции)                                                       работает автоматически
Загрузка файла, batch-импорт              Client streaming      клиент шлёт чанки; сервер
ETL-задания (1K событий → 1 ответ)                              отвечает в конце
Чат, multiplayer-игры, биржевой          Bidirectional         обе стороны независимо
order book, IoT-телеметрия                                      шлют - синхронизация на
с обратными командами                                            уровне HTTP/2 streams
Webhook-уведомления извне                 НЕ gRPC: HTTP POST   gRPC требует HTTP/2-клиента,
                                                                 partner-системы редко его держат
SSE-подобный broadcast на тысячи          Server streaming      работает, но для веба часто
веб-клиентов                              + gRPC-Web ИЛИ SSE   проще: SSE/WebSocket в браузере

Антипаттерны streaming

  • Эмулировать pub/sub через bidirectional: HTTP/2-стрим - это одно соединение, не очередь сообщений. При падении соединения вы теряете состояние. Нужен pub/sub - берите Kafka/NATS/Redis Streams.
  • Стримить «навсегда»: HTTP/2 ping/keepalive поможет, но прокси и LB могут разрывать долгие соединения (idle timeout 60s-5min). Делайте reconnect-логику или периодический keepalive-сообщения.
  • Client streaming для маленьких пакетов: если у вас 10 объектов - отправьте repeated Item в одном Unary-запросе. Streaming имеет смысл при сотнях/тысячах чанков или неизвестной длине.

Мини-практика

Реализуй сервис загрузки файлов: клиент отправляет файл чанками через client streaming, сервер возвращает размер и SHA256-checksum. Проверь поведение при отмене контекста на стороне клиента - сервер должен корректно прерваться. Затем сделай bidirectional chat-сервис с broadcast-ом сообщений всем подключённым клиентам через map[string]grpc.ServerStream.

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.