27 KiB
Семинар 2. TCP/UDP, Protobuf и gRPC
Задача: положить число в Storage
Допустим, у нас есть простой сервис Storage. Он хранит одно число, умеет обновлять его через PutValue и возвращать через GetValue. Клиент хочет положить туда 42. Казалось бы, очень простая задача. Что может пойти не так?
Первый вариант — сервер вообще недоступен: процесс упал, машина выключена или ещё не начала принимать соединения. Второй — сервер жив, но запрос потерялся по дороге.
Но есть более неприятные случаи. Сервер получил запрос и начал его выполнять, а потом сломался. Успел он записать число или нет? Клиент этого не знает. Или сервер всё выполнил и отправил ответ, но уже ответ потерялся в сети. Число записано, а клиент об этом не узнает.
Клиент не может отличить «сервер ничего не сделал» от «сервер всё сделал, но не ответил». Таймаут сообщает только, что мы не дождались ответа – он не дает понять состояние сервера.
Можно попробовать ещё раз положить 42: с точки зрения самого числа повторное присваивание даст тот же результат, если между попытками его никто не менял. Это идея идемпотентности. А вот с инкрементом такой приём уже не сработает: повтор увеличит число ещё раз. Можно спросить сервер, что в нём сейчас лежит, но и на этот запрос мы можем не получить ответа.
Как с этим жить? Для начала, хотелось бы вообще не думать о том, что пакеты могут теряться, дублироваться, фрагментироваться и изменяться. Посмотрим, какие абстракции для этого дают UDP и TCP, и какие проблемы всё равно останутся.
Что дают UDP и TCP
Откуда берутся проблемы с пакетами
Сеть ненадёжна по вполне конкретным причинам:
- Переполнение буферов. В маршрутизатор приходит 10 Гбит/с, а передать дальше он может только 1 Гбит/с. Очередь маршрутизатора конечна: когда она заполнится, новые пакеты придётся отбрасывать. Очередь может переполниться и на принимающем компьютере, если он не успевает обрабатывать данные.
- Отказ линии или устройства. Кабель повредили при ремонте, оборудование выключилось, а маршрутизатор ещё отправляет пакеты по старому маршруту. Отправитель при этом не обязательно сразу узнает, что случилось.
- Повреждение данных. Помеха изменила биты. Если проверка контрольной суммы обнаружит повреждение, пакет будет отброшен.
- Повторы и разные задержки. Например, при потере подтверждения на одном из нижних уровней возможна повторная передача уже доставленных данных. Очереди и смена маршрута могут привести к тому, что более поздний пакет обгонит ранний. Сам IP не обещает ни порядка, ни защиты от дублей.
- Фрагментация. Если IP-пакет больше допустимого размера на пути — MTU, — может потребоваться разбиение на фрагменты. Если потерялся один фрагмент, исходную датаграмму собрать уже нельзя.
Приложению неудобно разбираться с каждой такой ситуацией. Поэтому поверх IP используются транспортные протоколы. UDP сохраняет целостность отдельных сообщений, а TCP гарантирует порядок и надежность передачи.
UDP
UDP — User Datagram Protocol. Мы отправляем датаграмму — отдельный массив байтов. На принимающей стороне сохраняется граница сообщений, даже вопреки фрагментации IP пакетов.
А вот доставку, порядок и отсутствие дублей UDP не гарантирует. Сам протокол не ждёт подтверждений и не повторяет потерянные датаграммы. Поэтому для нашей операции PutValue пришлось бы отдельно решать, как понять, дошёл ли запрос и что делать при потере.
Для других задач такие свойства подходят. Например, в конференции Zoom не обязательно, чтобы дошёл каждый кусочек голоса: небольшой пропуск наш мозг часто может восстановить по смыслу.
TCP
TCP — Transmission Control Protocol — даёт абстракцию надёжного упорядоченного потока байтов. Можно представить канал: мы пишем в него байты, а другая сторона читает их в том же порядке, без повторов, дублей и потерь. Ниже этой абстракции все еще остается фрагментация пакетов, потери и повторы – но TCP все это скрывает от пользователя.
Сначала стороны устанавливают TCP-соединение. Каждая сторона должна сообщить ISN и другие свои параметры и получить подтверждение. Логически это четыре действия, но ответ сервера и его собственное начало обмена объединяются в одно сообщение. Получается три пакета: SYN → SYN + ACK → ACK.
После этого обе стороны хранят состояние: какие байты уже отправлены, какие подтверждены, какие получены и чего ещё не хватает. Соединение двунаправленное: клиент и сервер могут одновременно передавать данные, у каждого направления свой независимый поток байтов.
- Sequence number (
seq) — номер первого байта в сегменте. - Acknowledgment number (
ack) — номер следующего ожидаемого байта. Он подтверждает весь непрерывный префикс перед ним. - Окно приёма показывает, сколько данных получатель готов принять. Это позволяет отправлять несколько сегментов, не ожидая подтверждения каждого отдельно, но не переполнять буфер получателя. Дополнительно отправитель ограничивает объём данных в пути с учётом перегрузки сети.
Возьмём условную нумерацию с единицы. Пришли байты 4–6, а 1–3 потерялись. Получатель может сохранить 4–6 в буфере, но продолжает подтверждать, что ждёт байт 1. Отправитель обнаружит потерю по таймеру или повторным подтверждениям и повторит передачу. Когда появятся 1–3, образуется непрерывная последовательность 1–6: её можно отдать приложению, а в ack указать 7.
Если потерялось само подтверждение, отправитель тоже может повторить данные. Получатель узнает дубликат по номерам байтов и не отдаст его приложению второй раз. На схеме также упомянут SACK: он позволяет дополнительно сообщить, какие участки уже получены после пропуска.
За порядок приходится платить ожиданием: байты 4–6 уже у нас, но приложение не увидит их, пока не придут 1–3. Это head-of-line blocking.
RPC
Теперь хочется избавиться от ручной работы: не собирать сообщения PutValue, не разбирать ответы и не писать всю сериализацию/десериализацию сообщений самостоятельно.
Посмотрим на Storage как на обычный объект. У него есть методы PutValue и GetValue: мы передаём запрос и получаем результат. Можно сделать так, чтобы клиентский код выглядел именно как вызов метода, а отправка сообщения на сервер происходила внутри, практически прозрачно для пользователя.
Такую астракцию дает RPC — Remote Procedure Call, удалённый вызов процедуры. Она находится на уровень выше TCP и UDP. Пока всё работает, нам почти не приходится думать о сети. Но при ошибке нужно вспомнить, что вызов шел через сеть, и сервер мог не ответить по разным причинам.
Операцию и её параметры можно записать в JSON. Но каждый раз передавать названия полей довольно неэффективно. В нашем примере используются Protobuf для описания и сериализации сообщений и gRPC — широко используемый фреймворк удалённых вызовов. Показанный gRPC-сервис работает поверх HTTP/2 и TCP. Protobuf определяет, как представить наши данные, а gRPC связывает вызов клиентского метода с обработчиком на сервере.
Описываем Storage в .proto
Откроем storage.proto. В начале указаны версия proto3, пакет storage и настройка go_package для генерации Go-кода. Дальше описан сам сервис:
syntax = "proto3";
package storage;
option go_package = "hsegrpc/;storagepb";
import "google/protobuf/timestamp.proto";
service Storage {
rpc PutValue(PutRequest) returns (PutResponse);
rpc GetValue(GetRequest) returns (GetResponse);
}
message Value {
uint64 payload = 1;
optional google.protobuf.Timestamp updated_at = 2;
}
message PutRequest {
Value value = 1;
}
message PutResponse {
uint64 value = 1;
}
message GetRequest {
}
message GetResponse {
Value value = 1;
}
Здесь rpc объявляет метод. PutValue принимает PutRequest и возвращает PutResponse; с GetValue всё аналогично.
PutRequest содержит Value — значение, которое хотим положить. У Value есть само число payload и время обновления updated_at. GetRequest пустой, а GetResponse содержит Value. В PutResponse возвращается записанное число.
Что означают единички и двойки после полей? Это номера, по которым поля узнаются при сериализации. Вместо строки payload можно передавать номер 1, а получатель по своей схеме знает, что это за поле. Номер относится к конкретному типу сообщения: value = 1 в PutRequest и payload = 1 в Value друг другу не мешают. Это не значение поля и не его позиция в исходном файле.
Эти номера важны для совместимости. Допустим, мы решили убрать updated_at и добавить другое поле. Нельзя просто отдать новому полю номер 2: у пользователей старой версии под этим номером всё ещё описан Timestamp. Новому полю нужно дать новый номер, например 3. Тогда это будет именно новое поле, а не другое значение на месте старого.
Из .proto получаем код
Из одного .proto можно сгенерировать код для Go, Python, C++ и других языков. На семинаре мы смотрели Go; инструкции для разных языков есть на grpc.io.
В первой строке main.go приведена команда protoc для генерации. Из каталога grpc-practice/go-server она выглядит так:
protoc --go_out=. --go_opt=paths=source_relative \
--go-grpc_out=. --go-grpc_opt=paths=source_relative \
proto/storage.proto
Для неё нужны protoc и Go-плагины protoc-gen-go, protoc-gen-go-grpc. Команда создаёт два файла:
storage.pb.go— структурыValue,PutRequest,PutResponse,GetRequest,GetResponseи вспомогательные методы;storage_grpc.pb.go— клиентский и серверный интерфейсы, а также код для работы с gRPC.
Сообщения становятся структурами
В storage.pb.go каждому message соответствует Go-структура. Вот все типы из нашего примера; для наглядности здесь опущены служебные поля генератора, теги и вспомогательные методы:
type Value struct {
Payload uint64
UpdatedAt *timestamppb.Timestamp
}
type PutRequest struct {
Value *Value
}
type PutResponse struct {
Value uint64
}
type GetRequest struct{}
type GetResponse struct {
Value *Value
}
Можно прямо сопоставить их с .proto: payload превратился в Payload, updated_at — в UpdatedAt, а вложенное сообщение Value — в поле типа *Value. Поэтому запрос «положить 42» мы собираем как обычную структуру:
request := &storagepb.PutRequest{
Value: &storagepb.Value{
Payload: 42,
},
}
Клиент
Наверное, самый сложный момент: StorageClient. Имея его экземпляр, пользователь может вызвать методы PutValue и GetValue сервиса Storage. Его интерфейс сгенерирован в storage_grpc.pb.go (упрощенно):
type StorageClient interface {
PutValue(in *PutRequest) (*PutResponse, error)
GetValue(in *GetRequest) (*GetResponse, error)
}
Получить такой клиент позволяет сгенерированный NewStorageClient. Если conn — уже созданное gRPC-подключение, клиентский код выглядит так:
client := storagepb.NewStorageClient(conn)
response, err := client.PutValue(request)
if err != nil {
log.Fatal(err)
}
fmt.Println(response.GetValue()) // 42
Сервер
StorageServer задаёт интерфейс, который серверсу предстоит реализовать. Генератор не знает, как именно наш сервис Storage хранит число. Есть заготовка UnimplementedStorageServer: её методы отвечают, что операция не реализована. Дальше мы пишем собственные PutValue и GetValue. Здесь проходит важная граница: описание сервиса и сетевой обмен можно сгенерировать, а решение, что делать с пришедшим значением, остаётся за нами.
Вот интерфейс и один из методов заготовки из storage_grpc.pb.go (упрощенно):
type StorageServer interface {
PutValue(*PutRequest) (*PutResponse, error)
GetValue(*GetRequest) (*GetResponse, error)
mustEmbedUnimplementedStorageServer()
}
type UnimplementedStorageServer struct{}
func (UnimplementedStorageServer) PutValue(
_ *PutRequest,
) (*PutResponse, error) {
return nil, status.Errorf(
codes.Unimplemented, "method PutValue not implemented",
)
}
Свою реализацию мы пишем в storage/server.go. Наш Server включает UnimplementedStorageServer, хранит значение и mutex. Клиенты могут обращаться одновременно, поэтому доступ к значению защищён блокировкой:
type Value struct {
payload uint64
updatedAt time.Time
}
type Server struct {
storagepb.UnimplementedStorageServer
valueLocker sync.RWMutex
value Value
}
func NewServer() *Server {
return &Server{}
}
Здесь Value — наша внутренняя структура для хранения числа; storagepb.Value из предыдущих листингов — сообщение для передачи по сети. Встраивание UnimplementedStorageServer даёт заготовки методов и служебный метод из интерфейса, а нужные PutValue и GetValue мы определяем у своего Server.
В PutValue приходит сгенерированный PutRequest. Сначала проверяем, передал ли клиент value: он может его не указать, и тогда request.GetValue() вернёт nil. Если поле нужно для операции, эту проверку делаем сами. Иначе записываем число, выставляем время обновления и возвращаем ответ:
func (s *Server) PutValue(
request *storagepb.PutRequest,
) (*storagepb.PutResponse, error) {
s.valueLocker.Lock()
defer s.valueLocker.Unlock()
if request.GetValue() == nil {
return nil, errors.New("missed value")
}
s.value = Value{
payload: request.GetValue().GetPayload(),
updatedAt: time.Now(),
}
return &storagepb.PutResponse{
Value: s.value.payload,
}, nil
}
GetValue берёт блокировку на чтение и возвращает текущее значение. Наш вспомогательный метод toProto перекладывает его в сгенерированное Protobuf-сообщение:
func (s *Server) GetValue(
_ *storagepb.GetRequest,
) (*storagepb.GetResponse, error) {
s.valueLocker.RLock()
defer s.valueLocker.RUnlock()
return &storagepb.GetResponse{
Value: s.value.toProto(),
}, nil
}
func (v *Value) toProto() *storagepb.Value {
return &storagepb.Value{
Payload: v.payload,
UpdatedAt: timestamppb.New(v.updatedAt),
}
}
Соединяем реализацию с gRPC-сервером
Остаётся запустить сервер, используя main.go. Ниже основной фрагмент запуска; addr — адрес, который в нашем примере по умолчанию равен 0.0.0.0:51000:
lis, err := net.Listen("tcp", addr)
if err != nil {
log.Fatalf("failed to listen: %v", err)
}
grpcServer := grpc.NewServer()
reflection.Register(grpcServer)
storageService := storage.NewServer()
storagepb.RegisterStorageServer(grpcServer, storageService)
err = grpcServer.Serve(lis)
if err != nil {
log.Fatalf("server failed")
}
net.Listen("tcp", addr) открывает TCP-listener. grpc.NewServer() создаёт сервер gRPC, а storage.NewServer() — нашу реализацию Storage. Вызов RegisterStorageServer(grpcServer, storageService) связывает их: теперь gRPC знает, какому объекту передавать вызовы PutValue и GetValue. reflection.Register позволит инструментам вроде grpcurl узнать описание сервиса. Наконец, grpcServer.Serve(lis) начинает принимать соединения.
Получается полный путь одного запроса: клиентский метод получает PutRequest → сообщение сериализуется и передаётся по сети → сервер gRPC разбирает его → вызывает наш Server.PutValue → тот меняет значение и возвращает PutResponse → ответ уходит клиенту. Сериализацию и передачу обеспечивает готовый код, а изменение значения написано нами в server.go.
Проверяем через grpcurl и Postman
Для этого примера нужны Go и установленный grpcurl. В каталоге grpc-practice/go-server запускаем сервер:
go run main.go
В другом терминале смотрим список сервисов и описание Storage:
grpcurl -plaintext localhost:51000 list
grpcurl -plaintext localhost:51000 describe storage.Storage
-plaintext отключает TLS: в этом примере работаем без шифрования. list показывает доступные сервисы, а describe — методы Storage и типы их запросов и ответов. Узнать это у работающего сервера позволяет включённый в main.go сервис reflection: grpcurl получает описание интерфейса, хотя мы не передавали ему .proto отдельным файлом.
Теперь положим 100500, как в демонстрации, и прочитаем значение обратно:
grpcurl -plaintext -d '{"value":{"payload":100500}}' \
localhost:51000 storage.Storage/PutValue
grpcurl -plaintext -d '{}' localhost:51000 storage.Storage/GetValue
В запросе повторяется структура из .proto: внутри value находится payload. PutValue возвращает записанное число, GetValue — текущее значение и время обновления. В терминале мы вводим JSON, потому что это удобно человеку; grpcurl по описанию сервиса преобразует его в Protobuf-сообщение.
У PutValue ожидается ответ {"value":"100500"}, а у GetValue — вложенный объект value с полями payload и updatedAt. Время выставляет наш сервер. Если к нему одновременно обращается другой пишущий клиент, чтение может уже показать его значение.
Если не хочется писать команды в консоли, можно использовать Postman, как в конце семинара: импортировать storage.proto, указать адрес localhost:51000, выбрать PutValue, заполнить сообщение и нажать Invoke. Затем вызвать GetValue и увидеть записанное число.
Получилось, что мы описали сервис в .proto, сгенерировали код и дописали небольшую реализацию — и у нас уже есть работающий сервис. Все файлы и команды находятся в материалах практики.




