gRPC Streaming: потоковая передача данных
gRPC Streaming: потоковая передача данных
gRPC поддерживает четыре типа RPC: Unary (один запрос - один ответ), Server streaming (один запрос - много ответов), Client streaming (много запросов - один ответ), Bidirectional (много в обе стороны). Streaming - мощный инструмент для real-time приложений, файлов и подписок на события.
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-памяти между «горутинами»).
Когда использовать какой тип
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.