gRPC клиент
gRPC-клиент вызывает удалённые службы, используя контракт protobuf и транспорт HTTP/2.
В Kora клиент строится поверх сгенерированных классов stub библиотеки grpc-java: модуль создаёт ManagedChannel, подключает перехватчики и регистрирует готовые к использованию экземпляры stub в графе приложения.
Для каждой службы Kora делает доступными для внедрения сгенерированные stub-классы (BlockingStub, FutureStub, асинхронный Stub и корутинный stub для Kotlin),
исходный io.grpc.Channel и итоговый GrpcClientConfig, различая каждый клиент по @Tag сгенерированного класса службы
(например, @Tag(SimpleServiceGrpc.class)).
Транспорт gRPC-клиента использует Netty, поэтому общие настройки event loop и транспорта можно задать в разделе Netty.
Если нужен пошаговый разбор перед справочным описанием, смотрите gRPC-клиент и продвинутый gRPC-клиент.
Подключение¶
Зависимость build.gradle:
implementation "ru.tinkoff.kora:grpc-client"
implementation "io.grpc:grpc-protobuf:1.74.0"
implementation "javax.annotation:javax.annotation-api:1.3.2"
Модуль:
Зависимость build.gradle.kts:
implementation("ru.tinkoff.kora:grpc-client")
implementation("io.grpc:grpc-protobuf:1.74.0")
implementation("javax.annotation:javax.annotation-api:1.3.2")
Модуль:
Плагин¶
Код gRPC-клиента создаётся с помощью Gradle-плагина protobuf.
Плагин генерирует классы Java-сообщений из контракта protobuf и gRPC-классы stub, которые затем используются Kora.
Плагин build.gradle:
plugins {
id "com.google.protobuf" version "0.9.4"
}
protobuf {
protoc { artifact = "com.google.protobuf:protoc:3.25.3" }
plugins {
grpc { artifact = "io.grpc:protoc-gen-grpc-java:1.74.0" }
}
generateProtoTasks {
all()*.plugins { grpc {} }
}
}
sourceSets {
main.java {
srcDirs "build/generated/source/proto/main/grpc"
srcDirs "build/generated/source/proto/main/java"
}
}
Плагин build.gradle.kts:
import com.google.protobuf.gradle.id
plugins {
id("com.google.protobuf") version ("0.9.4")
}
protobuf {
protoc { artifact = "com.google.protobuf:protoc:3.25.3" }
plugins {
id("grpc") { artifact = "io.grpc:protoc-gen-grpc-java:1.74.0" }
}
generateProtoTasks {
ofSourceSet("main").forEach { it.plugins { id("grpc") { } } }
}
}
kotlin {
sourceSets.main {
kotlin.srcDir("build/generated/source/proto/main/grpc")
kotlin.srcDir("build/generated/source/proto/main/java")
}
}
Конфигурация¶
gRPC-клиент для службы SimpleService будет иметь путь конфигурации grpcClient.SimpleService.
Основные параметры конфигурации:
URLсервера, куда будут отправляться запросы (обязательный, по умолчанию не указано).- Максимальное время выполнения запроса (по умолчанию не указано, необязательно). Значение применяется как
deadline, если у вызова ещё нет собственногоdeadline.
URLсервера, куда будут отправляться запросы (обязательный, по умолчанию не указано).- Максимальное время выполнения запроса (по умолчанию не указано, необязательно). Значение применяется как
deadline, если у вызова ещё нет собственногоdeadline.
Полная конфигурация
Пример полной конфигурации, описанной в классе GrpcClientConfig:
grpcClient {
SimpleService {
url = "http://localhost:8090" //(1)!
timeout = "10s" //(2)!
keepAliveTime = "0s" //(3)!
keepAliveTimeout = "0s" //(4)!
loadBalancingPolicy = "pick_first" //(5)!
defaultServiceConfig { //(6)!
loadBalancingConfig = [
{
round_robin = {}
}
]
}
telemetry {
logging {
enabled = false //(7)!
}
metrics {
enabled = true //(8)!
slo = [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] //(9)!
tags = { // (10)!
"key1" = "value1"
"key2" = "value2"
}
}
tracing {
enabled = true //(11)!
attributes = { // (12)!
"key1" = "value1"
"key2" = "value2"
}
}
}
}
}
URLсервера, куда будут отправляться запросы (обязательная, по умолчанию не указано).- Максимальное время выполнения запроса (по умолчанию не указано, необязательно). Значение применяется как
deadline, если у вызова ещё нет собственногоdeadline. - Интервал между gRPC-фреймами
PING(по умолчанию не указано, необязательно). - Время ожидания подтверждения фрейма
PING(по умолчанию не указано, необязательно). Если подтверждение не получено за это время, соединение закрывается. - Политика балансировки нагрузки для
ManagedChannelBuilder(по умолчанию не указано, необязательно). - Стандартная конфигурация службы gRPC, передаваемая в
ManagedChannelBuilder.defaultServiceConfig(по умолчанию не указано, необязательно). - Включает логирование модуля (по умолчанию:
false). - Включает метрики модуля (по умолчанию:
true). - Настраивает SLO для метрики DistributionSummary (по умолчанию:
TelemetryConfig.MetricsConfig.DEFAULT_SLO). - Дополнительные теги для метрик (по умолчанию:
{}). - Включает трассировку модуля (по умолчанию:
true). - Дополнительные атрибуты для трассировки (по умолчанию:
{}).
grpcClient:
SimpleService:
url: "http://localhost:8090" #(1)!
timeout: "10s" #(2)!
keepAliveTime: "0s" #(3)!
keepAliveTimeout: "0s" #(4)!
loadBalancingPolicy: "pick_first" #(5)!
defaultServiceConfig: #(6)!
loadBalancingConfig:
- round_robin: {}
telemetry:
logging:
enabled: false #(7)!
metrics:
enabled: true #(8)!
slo: [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] #(9)!
tags: #(10)!
key1: value1
key2: value2
tracing:
enabled: true #(11)!
attributes: #(12)!
key1: value1
key2: value2
URLсервера, куда будут отправляться запросы (обязательная, по умолчанию не указано).- Максимальное время выполнения запроса (по умолчанию не указано, необязательно). Значение применяется как
deadline, если у вызова ещё нет собственногоdeadline. - Интервал между gRPC-фреймами
PING(по умолчанию не указано, необязательно). - Время ожидания подтверждения фрейма
PING(по умолчанию не указано, необязательно). Если подтверждение не получено за это время, соединение закрывается. - Политика балансировки нагрузки для
ManagedChannelBuilder(по умолчанию не указано, необязательно). - Стандартная конфигурация службы gRPC, передаваемая в
ManagedChannelBuilder.defaultServiceConfig(по умолчанию не указано, необязательно). - Включает логирование модуля (по умолчанию:
false). - Включает метрики модуля (по умолчанию:
true). - Настраивает SLO для метрики DistributionSummary (по умолчанию:
TelemetryConfig.MetricsConfig.DEFAULT_SLO). - Дополнительные теги для метрик (по умолчанию:
{}). - Включает трассировку модуля (по умолчанию:
true). - Дополнительные атрибуты для трассировки (по умолчанию:
{}).
Транспорт и TLS¶
Схема в url выбирает транспорт при создании ManagedChannel (ManagedChannelLifecycle):
http— незашифрованный транспортplaintext(usePlaintext()у построителя), порт по умолчанию80, если порт не указан.https— транспортTLS, порт по умолчанию443, если порт не указан.- любая другая схема — порт должен быть указан явно, иначе при определении порта по умолчанию во время запуска будет выброшено исключение
IllegalArgumentExceptionс сообщениемUnknown scheme '<scheme>'. Режимplaintextвключается только для схемыhttp, поэтому остальные схемы используют транспорт gRPC по умолчанию (TLS), если недоступен портplaintext.
Для нестандартного TLS (mTLS, собственное хранилище доверенных сертификатов или транспорт, отличный от Netty) зарегистрируйте свой компонент GrpcClientChannelFactory.
Он создаёт ManagedChannelBuilder и может передавать io.grpc.ChannelCredentials; реализация по умолчанию — GrpcNettyClientChannelFactory,
которая привязывает канал к общей для Kora группе Netty EventLoopGroup.
@Component
public final class TlsGrpcClientChannelFactory implements GrpcClientChannelFactory {
@Override
public ManagedChannelBuilder<?> forAddress(SocketAddress serverAddress) {
return NettyChannelBuilder.forAddress(serverAddress);
}
@Override
public ManagedChannelBuilder<?> forAddress(SocketAddress serverAddress, ChannelCredentials creds) {
return NettyChannelBuilder.forAddress(serverAddress, creds);
}
@Override
public ManagedChannelBuilder<?> forTarget(String target) {
return NettyChannelBuilder.forTarget(target);
}
@Override
public ManagedChannelBuilder<?> forTarget(String target, ChannelCredentials creds) {
return NettyChannelBuilder.forTarget(target, creds);
}
}
@Component
class TlsGrpcClientChannelFactory : GrpcClientChannelFactory {
override fun forAddress(serverAddress: SocketAddress): ManagedChannelBuilder<*> {
return NettyChannelBuilder.forAddress(serverAddress)
}
override fun forAddress(serverAddress: SocketAddress, creds: ChannelCredentials): ManagedChannelBuilder<*> {
return NettyChannelBuilder.forAddress(serverAddress, creds)
}
override fun forTarget(target: String): ManagedChannelBuilder<*> {
return NettyChannelBuilder.forTarget(target)
}
override fun forTarget(target: String, creds: ChannelCredentials): ManagedChannelBuilder<*> {
return NettyChannelBuilder.forTarget(target, creds)
}
}
Реализация для промышленного окружения также должна привязывать построитель к общей для Kora группе Netty EventLoopGroup и NettyChannelFactory
так же, как это делает GrpcNettyClientChannelFactory, вместо того чтобы позволять Netty создавать собственный event loop.
Ограничение по времени¶
Значение timeout применяется всегда включённым перехватчиком GrpcClientConfigInterceptor как deadline вызова, но только если у вызова нет собственного deadline.
Заданный для конкретного вызова deadline через stub всегда имеет приоритет над настроенным timeout:
Когда истекает deadline, вызов завершается с ошибкой StatusRuntimeException, несущей Status.DEADLINE_EXCEEDED.
Конфигурация канала¶
keepAliveTime/keepAliveTimeoutсоответствуют настройкамPINGуManagedChannelBuilder.keepAliveTime— интервал между фреймамиPINGпротоколаHTTP/2на простаивающем соединении;keepAliveTimeout— как долго ждать подтвержденияPINGперед закрытием соединения. Оба отключены, если не заданы.loadBalancingPolicyсоответствуетManagedChannelBuilder.defaultLoadBalancingPolicy. Значение по умолчанию в gRPC —pick_first(единственное соединение с первым разрешённым адресом);round_robinраспределяет вызовы по всем разрешённым адресам и обычно используется с DNS-адресами, возвращающими несколько записейA/AAAA.defaultServiceConfigпередаётся как есть вManagedChannelBuilder.defaultServiceConfigи содержит нативную карту конфигурации службы gRPC (loadBalancingConfig,methodConfigдля отдельных методов с политикой повторов/хеджирования и т. д.). Она описана обёрткойDefaultServiceConfigнадMap<String, Object>.
Настройщик построителя канала¶
Если конфигурации через файл недостаточно, можно зарегистрировать компонент GrpcClientBuilderConfigurer.
Он получает уже подготовленный ManagedChannelBuilder и позволяет настроить канал в коде до его создания.
Сначала применяются настройки из GrpcClientConfig, затем вызывается GrpcClientBuilderConfigurer.
Также можно настроить транспорт Netty.
Предоставляемые метрики модуля описаны в разделе Справочник метрик.
Служба¶
Созданные экземпляры gRPC stub можно внедрять как зависимости:
stub также можно внедрить напрямую в конструктор @Component:
Типы реализаций¶
Плагин protobuf генерирует для одной службы (SimpleService) несколько stub-классов. Каждый доступен для внедрения простым объявлением соответствующего типа;
на самом stub указывать @Tag не нужно (Kora самостоятельно разрешает помеченный тегом Channel):
| Тип stub | Модель вызова | Когда использовать |
|---|---|---|
SimpleServiceBlockingStub |
Синхронная; возвращает ответ напрямую (или Iterator для серверной потоковой передачи) |
Блокирующий код, простейший стиль вызова |
SimpleServiceFutureStub |
Асинхронная; возвращает ListenableFuture<T> (только унарные вызовы) |
Неблокирующий код с использованием ListenableFuture |
SimpleServiceStub (асинхронный) |
Асинхронная; передаёт результаты через обратные вызовы StreamObserver<T> |
Любая потоковая передача, асинхронные вызовы на обратных вызовах |
| Корутинный stub для Kotlin | suspend-функции и Flow<T> |
Идиоматичные корутины Kotlin |
BlockingStub, FutureStub и асинхронный Stub связываются расширением обработчика аннотаций (GrpcClientExtension), которое обнаруживает stub-типы @GrpcGenerated
и вызывает сгенерированную фабрику newBlockingStub / newFutureStub / newStub для помеченного тегом Channel.
@KoraApp
public interface Application extends HoconConfigModule, GrpcClientModule {
default BlockingCaller blockingCaller(SimpleServiceGrpc.SimpleServiceBlockingStub stub) {
return new BlockingCaller(stub);
}
default FutureCaller futureCaller(SimpleServiceGrpc.SimpleServiceFutureStub stub) {
return new FutureCaller(stub);
}
default AsyncCaller asyncCaller(SimpleServiceGrpc.SimpleServiceStub stub) {
return new AsyncCaller(stub);
}
}
@KoraApp
interface Application : HoconConfigModule, GrpcClientModule {
fun blockingCaller(stub: SimpleServiceGrpc.SimpleServiceBlockingStub) = BlockingCaller(stub)
fun futureCaller(stub: SimpleServiceGrpc.SimpleServiceFutureStub) = FutureCaller(stub)
fun asyncCaller(stub: SimpleServiceGrpc.SimpleServiceStub) = AsyncCaller(stub)
}
Для вызовов на корутинах Kotlin сгенерируйте корутинный stub с помощью генератора gRPC Kotlin
(io.grpc:protoc-gen-grpc-kotlin). Сгенерированный stub наследуется от io.grpc.kotlin.AbstractCoroutineStub и помечен аннотацией @StubFor;
далее обработчик символов KSP генерирует модуль Kora, который предоставляет его как @DefaultComponent, привязанный к помеченному тегом Channel, поэтому он внедряется таким же образом:
Стили вызова¶
Форма rpc в контракте .proto (одиночный или stream запрос/ответ) определяет сигнатуру сгенерированного метода.
В примерах ниже базовый контракт дополнен всеми четырьмя стилями вызова:
service SimpleService {
rpc unary(RequestEvent) returns (ResponseEvent) {} // унарный
rpc serverStream(RequestEvent) returns (stream ResponseEvent) {} // серверная потоковая передача
rpc clientStream(stream RequestEvent) returns (ResponseEvent) {} // клиентская потоковая передача
rpc biDiStream(stream RequestEvent) returns (stream ResponseEvent) {} // двунаправленная потоковая передача
}
Запросы создаются с помощью сгенерированных построителей сообщений (RequestEvent.newBuilder()).
В Java унарные вызовы и серверная потоковая передача доступны у BlockingStub, тогда как клиентская потоковая передача и двунаправленные вызовы требуют асинхронного Stub
(у них нет блокирующего варианта). Корутинный stub Kotlin выражает каждый стиль через suspend-функции и Flow<T>:
var request = RequestEvent.newBuilder().setName("bob").setCode("b1").build();
// унарный — BlockingStub
ResponseEvent unary = blockingStub.unary(request);
// унарный — FutureStub
ListenableFuture<ResponseEvent> future = futureStub.unary(request);
// серверная потоковая передача — BlockingStub возвращает итератор
Iterator<ResponseEvent> responses = blockingStub.serverStream(request);
// серверная потоковая передача — асинхронный Stub передаёт результаты в StreamObserver
asyncStub.serverStream(request, new StreamObserver<>() {
@Override public void onNext(ResponseEvent value) { /* ... */ }
@Override public void onError(Throwable t) { /* ... */ }
@Override public void onCompleted() { /* ... */ }
});
// клиентская потоковая передача — асинхронный Stub, пишем запросы, читаем один ответ
StreamObserver<ResponseEvent> responseObserver = new StreamObserver<>() {
@Override public void onNext(ResponseEvent value) { /* единственный ответ */ }
@Override public void onError(Throwable t) { /* ... */ }
@Override public void onCompleted() { /* ... */ }
};
StreamObserver<RequestEvent> requestObserver = asyncStub.clientStream(responseObserver);
requestObserver.onNext(request);
requestObserver.onCompleted();
// двунаправленная потоковая передача — асинхронный Stub, потоки с обеих сторон
StreamObserver<RequestEvent> bidi = asyncStub.biDiStream(responseObserver);
bidi.onNext(request);
bidi.onCompleted();
val request = RequestEvent.newBuilder().setName("bob").setCode("b1").build()
// унарный — suspend-функция
val response: ResponseEvent = coroutineStub.unary(request)
// серверная потоковая передача — возвращает Flow
val serverFlow: Flow<ResponseEvent> = coroutineStub.serverStream(request)
serverFlow.collect { event -> /* ... */ }
// клиентская потоковая передача — принимает Flow, возвращает один ответ
val clientResponse: ResponseEvent = coroutineStub.clientStream(flowOf(request))
// двунаправленная потоковая передача — Flow на вход, Flow на выход
val biDiFlow: Flow<ResponseEvent> = coroutineStub.biDiStream(flowOf(request))
biDiFlow.collect { event -> /* ... */ }
Внедрение Channel и конфигурации¶
Для продвинутого или ручного создания stub можно внедрить исходный io.grpc.Channel и итоговый GrpcClientConfig,
пометив их сгенерированным классом службы. Оба предоставляются GrpcClientExtension:
@KoraApp
public interface Application extends HoconConfigModule, GrpcClientModule {
default SomeService someService(@Tag(SimpleServiceGrpc.class) Channel channel,
@Tag(SimpleServiceGrpc.class) GrpcClientConfig config) {
return new SomeService(SimpleServiceGrpc.newBlockingStub(channel), config.url());
}
}
Перехватчики¶
Перехватчики позволяют перехватывать запросы до того, как они будут переданы службам.
По умолчанию¶
По умолчанию при запуске клиента используются следующие перехватчики:
GrpcClientConfigInterceptor— применяетtimeoutкакdeadlineвызова, если у вызова его нет.GrpcClientTelemetryInterceptor, если для клиента доступна телеметрия.
Собственные¶
В отличие от HTTP-клиента, у перехватчиков gRPC нет уровней метода/класса/глобального уровня.
Каждый перехватчик действует в пределах одного клиента — за счёт пометки компонента сгенерированным классом службы (@Tag(SimpleServiceGrpc.class)).
Зарегистрируйте перехватчик как компонент с этим тегом:
@Tag(SimpleServiceGrpc.class)
@Component
public final class MyClientInterceptor implements ClientInterceptor {
@Override
public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, Channel next) {
LoggerFactory.getLogger(Application.class).info("INTERCEPTED");
return next.newCall(method, callOptions);
}
}
Чтобы применить один компонент-перехватчик к нескольким клиентам («общий» перехватчик), укажите ему несколько значений @Tag — по одному сгенерированному классу службы на каждый клиент:
@Tag({SimpleServiceGrpc.class, OtherServiceGrpc.class})
@Component
public final class SharedInterceptor implements ClientInterceptor {
@Override
public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, Channel next) {
return next.newCall(method, callOptions);
}
}
@Tag(SimpleServiceGrpc::class, OtherServiceGrpc::class)
@Component
class SharedInterceptor : ClientInterceptor {
override fun <ReqT : Any, RespT : Any> interceptCall(
method: MethodDescriptor<ReqT, RespT>,
callOptions: CallOptions,
next: Channel
): ClientCall<ReqT, RespT> {
return next.newCall(method, callOptions)
}
}
Порядок выполнения:
ManagedChannelLifecycle собирает все перехватчики, помеченные для службы, как All<ClientInterceptor> и применяет их в фиксированном порядке:
сначала ваши собственные перехватчики, затем перехватчик телеметрии (если телеметрия включена), и последним — перехватчик конфигурации/deadline.
Запрос → Собственные перехватчики → Перехватчик телеметрии → Перехватчик конфигурации (deadline) → gRPC-сервер
Поскольку перехватчик deadline выполняется последним, deadline, установленный собственным перехватчиком в CallOptions, сохраняется, а настроенный timeout
применяется только тогда, когда его не задал ни один более ранний перехватчик (или withDeadlineAfter для конкретного вызова).
В качестве альтернативы можно изменить stub с помощью GraphInterceptor.
Авторизация¶
В gRPC нет отдельного модуля авторизации: авторизация выполняется с помощью ClientInterceptor, помеченного классом службы, который добавляет заголовок
Authorization (или API-ключа) в Metadata исходящего вызова. Перехватчик оборачивает вызов в ForwardingClientCall.SimpleForwardingClientCall
и помещает заголовок в start(...), до отправки запроса.
Bearer¶
Перехватчик Bearer читает токен из вашего собственного поставщика и
помещает его в заголовок Authorization каждого вызова. TokenProvider ниже — ваш собственный компонент:
@Tag(SimpleServiceGrpc.class)
@Component
public final class BearerAuthInterceptor implements ClientInterceptor {
private static final Metadata.Key<String> AUTHORIZATION =
Metadata.Key.of("Authorization", Metadata.ASCII_STRING_MARSHALLER);
private final TokenProvider tokenProvider;
public BearerAuthInterceptor(TokenProvider tokenProvider) {
this.tokenProvider = tokenProvider;
}
@Override
public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, Channel next) {
return new ForwardingClientCall.SimpleForwardingClientCall<>(next.newCall(method, callOptions)) {
@Override
public void start(Listener<RespT> responseListener, Metadata headers) {
headers.put(AUTHORIZATION, "Bearer " + tokenProvider.getToken());
super.start(responseListener, headers);
}
};
}
}
@Tag(SimpleServiceGrpc::class)
@Component
class BearerAuthInterceptor(private val tokenProvider: TokenProvider) : ClientInterceptor {
override fun <ReqT : Any, RespT : Any> interceptCall(
method: MethodDescriptor<ReqT, RespT>,
callOptions: CallOptions,
next: Channel
): ClientCall<ReqT, RespT> {
return object : ForwardingClientCall.SimpleForwardingClientCall<ReqT, RespT>(next.newCall(method, callOptions)) {
override fun start(responseListener: Listener<RespT>, headers: Metadata) {
headers.put(AUTHORIZATION, "Bearer " + tokenProvider.getToken())
super.start(responseListener, headers)
}
}
}
companion object {
private val AUTHORIZATION: Metadata.Key<String> =
Metadata.Key.of("Authorization", Metadata.ASCII_STRING_MARSHALLER)
}
}
ApiKey¶
Перехватчик API-ключа помещает статический ключ в собственный заголовок метаданных (например, X-API-KEY).
Ключ читается из интерфейса @ConfigSource, внедрённого в перехватчик:
@Tag(SimpleServiceGrpc.class)
@Component
public final class ApiKeyInterceptor implements ClientInterceptor {
private static final Metadata.Key<String> API_KEY =
Metadata.Key.of("X-API-KEY", Metadata.ASCII_STRING_MARSHALLER);
private final String apiKey;
public ApiKeyInterceptor(ApiKeyConfig config) { //(1)!
this.apiKey = config.apiKey();
}
@Override
public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, Channel next) {
return new ForwardingClientCall.SimpleForwardingClientCall<>(next.newCall(method, callOptions)) {
@Override
public void start(Listener<RespT> responseListener, Metadata headers) {
headers.put(API_KEY, apiKey);
super.start(responseListener, headers);
}
};
}
}
- Любой интерфейс
@ConfigSource, предоставляющий API-ключ, напримерString apiKey();
@Tag(SimpleServiceGrpc::class)
@Component
class ApiKeyInterceptor(config: ApiKeyConfig) : ClientInterceptor { //(1)!
private val apiKey: String = config.apiKey()
override fun <ReqT : Any, RespT : Any> interceptCall(
method: MethodDescriptor<ReqT, RespT>,
callOptions: CallOptions,
next: Channel
): ClientCall<ReqT, RespT> {
return object : ForwardingClientCall.SimpleForwardingClientCall<ReqT, RespT>(next.newCall(method, callOptions)) {
override fun start(responseListener: Listener<RespT>, headers: Metadata) {
headers.put(API_KEY, apiKey)
super.start(responseListener, headers)
}
}
}
companion object {
private val API_KEY: Metadata.Key<String> =
Metadata.Key.of("X-API-KEY", Metadata.ASCII_STRING_MARSHALLER)
}
}
- Любой интерфейс
@ConfigSource, предоставляющий API-ключ, напримерfun apiKey(): String
Обработка ошибок¶
Неуспешный вызов gRPC выбрасывает io.grpc.StatusRuntimeException. Его getStatus() несёт Status.Code
(коды статусов), например UNAVAILABLE (сервер недоступен), DEADLINE_EXCEEDED (истёк timeout/deadline),
UNAUTHENTICATED (отклонённые учётные данные) или INVALID_ARGUMENT. Метаданные ответа доступны через getTrailers().
Причины:
UNAVAILABLE— неверныйurl, несоответствие plaintext/TLS или сервер не работает.DEADLINE_EXCEEDED— превышен настроенныйtimeoutилиwithDeadlineAfterдля конкретного вызова.UNAUTHENTICATED/PERMISSION_DENIED— отсутствуют или недействительны метаданные авторизации.
Рекомендации:
- Ветвитесь по
e.getStatus().getCode(), а не по типу исключения. - Используйте аспекты отказоустойчивости (
@Retry,@CircuitBreaker,@Timeout) на оборачивающем методе службы для временных сбоев.
Тестирование¶
gRPC-клиент тестируется как любой другой компонент Kora с помощью @KoraAppTest.
Реализуйте KoraAppTestConfigModifier, чтобы задать url (например, через подстановку переменной окружения GRPC_URL, используемую в примере),
внедрите службу на основе stub через @TestComponent, соберите запрос сгенерированным построителем и проверьте StatusRuntimeException:
@KoraAppTest(Application.class)
class GrpcClientTests implements KoraAppTestConfigModifier {
@TestComponent
private RootService service;
@Override
public KoraConfigModification config() {
return KoraConfigModification.ofSystemProperty("GRPC_URL", "grpc://localhost:8090");
}
@Test
void createUser() {
var event = Message.RequestEvent.newBuilder()
.setName("bob")
.setCode("b1")
.build();
var stub = service.service();
assertThrows(StatusRuntimeException.class, () -> stub.createUser(event));
}
}
@KoraAppTest(Application::class)
class GrpcClientTests : KoraAppTestConfigModifier {
@TestComponent
lateinit var service: RootService
override fun config(): KoraConfigModification =
KoraConfigModification.ofSystemProperty("GRPC_URL", "grpc://localhost:8090")
@Test
fun createUser() {
val event = Message.RequestEvent.newBuilder()
.setName("bob")
.setCode("b1")
.build()
val stub = service.service()
assertThrows(StatusRuntimeException::class.java) { stub.createUser(event) }
}
}
Телеметрия¶
gRPC Client использует контракт телеметрии для логирования, метрик и трассировки вызовов.
Конфигурация телеметрии (секция telemetry { logging / metrics / tracing }) описана в разделе Конфигурация.
Точки расширения находятся в ru.tinkoff.kora.grpc.client.common.telemetry.
Для каждого gRPC-вызова создаётся GrpcClientTelemetry.GrpcClientTelemetryContext, который закрывается по завершении вызова.
Вызов описывается через параметры обработчика телеметрии, включая сервис, метод, статус ответа и длительность.
Фабрика по умолчанию DefaultGrpcClientTelemetryFactory объединяет три фабрики:
- GrpcClientLoggerFactory строит GrpcClientLogger для логирования начала/конца вызова;
- GrpcClientMetricsFactory строит GrpcClientMetrics для записи метрик вызовов;
- GrpcClientTracerFactory строит GrpcClientTracer для распределённой трассировки.
Метрики и трассировка описаны в разделе Справочник метрик.