gRPC client
The gRPC client calls remote services using a protobuf contract and the HTTP/2 transport.
In Kora, the client is built on top of generated grpc-java stub classes: the module creates a ManagedChannel, attaches interceptors, and registers ready-to-use stub instances in the application graph.
For each service, Kora makes the generated stubs (BlockingStub, FutureStub, the async Stub, and the Kotlin coroutine stub),
the raw io.grpc.Channel, and the resolved GrpcClientConfig injectable, distinguishing every client by the generated service-class @Tag
(for example @Tag(SimpleServiceGrpc.class)).
The gRPC client transport uses Netty, so common event loop and transport settings can be configured in the Netty section.
For a step-by-step walkthrough before the reference details, see gRPC Client and Advanced gRPC Client.
Dependency¶
Dependency 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"
Module:
Dependency 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")
Module:
Plugin¶
The code for the gRPC client is created with the protobuf Gradle plugin.
The plugin generates Java message classes from the protobuf contract and gRPC stub classes that are then used by Kora.
Plugin 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"
}
}
Plugin 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")
}
}
Configuration¶
A gRPC client for the SimpleService service will have the grpcClient.SimpleService configuration path.
Basic configuration parameters:
- Server
URLwhere requests will be sent (required, no default). - Maximum request execution time (default: not specified, optional). The value is applied as a
deadlineif the call does not already have its owndeadline.
Full Configuration
Example of the complete configuration described by the GrpcClientConfig class:
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"
}
}
}
}
}
- Server
URLwhere requests will be sent (required, default: not specified). - Maximum request execution time (default: not specified, optional). The value is applied as a
deadlineif the call does not already have its owndeadline. - Interval between gRPC
PINGframes (default: not specified, optional). - Timeout for acknowledging a
PINGframe (default: not specified, optional). If the acknowledgement is not received within this time, the connection is closed. - Load balancing policy for
ManagedChannelBuilder(default: not specified, optional). - Standard gRPC service configuration passed to
ManagedChannelBuilder.defaultServiceConfig(default: not specified, optional). - Enables module logging (default:
false). - Enables module metrics (default:
true). - Configures SLO for the DistributionSummary metric (default:
TelemetryConfig.MetricsConfig.DEFAULT_SLO). - Additional tags for metrics (default:
{}). - Enables module tracing (default:
true). - Additional attributes for tracing (default:
{}).
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
- Server
URLwhere requests will be sent (required, default: not specified). - Maximum request execution time (default: not specified, optional). The value is applied as a
deadlineif the call does not already have its owndeadline. - Interval between gRPC
PINGframes (default: not specified, optional). - Timeout for acknowledging a
PINGframe (default: not specified, optional). If the acknowledgement is not received within this time, the connection is closed. - Load balancing policy for
ManagedChannelBuilder(default: not specified, optional). - Standard gRPC service configuration passed to
ManagedChannelBuilder.defaultServiceConfig(default: not specified, optional). - Enables module logging (default:
false). - Enables module metrics (default:
true). - Configures SLO for the DistributionSummary metric (default:
TelemetryConfig.MetricsConfig.DEFAULT_SLO). - Additional tags for metrics (default:
{}). - Enables module tracing (default:
true). - Additional attributes for tracing (default:
{}).
Transport and TLS¶
The url scheme selects the transport when the ManagedChannel is created (ManagedChannelLifecycle):
http— plaintext transport (usePlaintext()on the builder), default port80when the port is omitted.https—TLStransport, default port443when the port is omitted.- any other scheme — the port must be specified explicitly, otherwise an
IllegalArgumentExceptionwith the messageUnknown scheme '<scheme>'is thrown at startup while resolving the default port. Plaintext is enabled only for thehttpscheme, so other schemes fall back to the default gRPC (TLS) transport unless a plaintext port is reachable.
For custom TLS (mTLS, a custom trust store, or non-Netty transport) register your own GrpcClientChannelFactory component.
It builds the ManagedChannelBuilder and can pass an io.grpc.ChannelCredentials; the default implementation is GrpcNettyClientChannelFactory,
which binds the channel to Kora's shared 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)
}
}
A production override should also bind the builder to Kora's shared Netty EventLoopGroup and NettyChannelFactory
the way GrpcNettyClientChannelFactory does, rather than letting Netty create its own event loop.
Timeouts¶
The timeout value is applied by the always-on GrpcClientConfigInterceptor as a call deadline, but only when the call has no deadline of its own.
A per-call deadline set through the stub always wins over the configured timeout:
When a deadline expires, the call fails with a StatusRuntimeException carrying Status.DEADLINE_EXCEEDED.
Channel config¶
keepAliveTime/keepAliveTimeoutmap toManagedChannelBuilderPINGsettings.keepAliveTimeis the interval betweenHTTP/2PINGframes on an idle connection;keepAliveTimeoutis how long to wait for thePINGacknowledgement before closing the connection. Both are disabled unless set.loadBalancingPolicymaps toManagedChannelBuilder.defaultLoadBalancingPolicy. The gRPC default ispick_first(a single connection to the first resolved address);round_robindistributes calls across all resolved addresses and is typically used with DNS targets that return multipleA/AAAArecords.defaultServiceConfigis passed as-is toManagedChannelBuilder.defaultServiceConfigand carries the native gRPC service config map (loadBalancingConfig, per-methodmethodConfigwith retry/hedging policy, etc.). It is described by theDefaultServiceConfigwrapper overMap<String, Object>.
Channel builder configurer¶
If file-based configuration is not enough, you can register a GrpcClientBuilderConfigurer component.
It receives an already prepared ManagedChannelBuilder and lets you configure the channel in code before it is created.
Settings from GrpcClientConfig are applied first, then GrpcClientBuilderConfigurer is called.
You can also configure Netty transport.
Module metrics are described in the Metrics Reference section.
Service¶
Created gRPC stub instances can be injected as dependencies:
A stub can also be injected straight into a @Component constructor:
Stub types¶
The protobuf plugin generates several stub classes for one service (SimpleService). Each is injectable by simply declaring the corresponding type;
no @Tag is required on the stub itself (Kora resolves the tagged Channel behind the scenes):
| Stub type | Call model | When to use |
|---|---|---|
SimpleServiceBlockingStub |
Synchronous; returns the response directly (or an Iterator for server streaming) |
Blocking code, simplest call style |
SimpleServiceFutureStub |
Asynchronous; returns a ListenableFuture<T> (unary only) |
Non-blocking code using ListenableFuture |
SimpleServiceStub (async) |
Asynchronous; delivers results through StreamObserver<T> callbacks |
Any streaming, callback-style asynchronous calls |
| Kotlin coroutine stub | suspend functions and Flow<T> |
Idiomatic Kotlin coroutines |
The BlockingStub, FutureStub, and async Stub are wired by the annotation-processor extension (GrpcClientExtension), which detects the @GrpcGenerated
stub types and calls the generated newBlockingStub / newFutureStub / newStub factory against the tagged 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)
}
For Kotlin coroutine calls, generate the coroutine stub with the gRPC Kotlin generator
(io.grpc:protoc-gen-grpc-kotlin). The generated stub extends io.grpc.kotlin.AbstractCoroutineStub and is annotated with @StubFor;
the KSP symbol processor then generates a Kora module that exposes it as a @DefaultComponent bound to the tagged Channel, so it is injected the same way:
Call styles¶
The rpc shape in the .proto contract (single vs stream request/response) determines the generated method signature.
The examples below extend the base contract with all four call styles:
service SimpleService {
rpc unary(RequestEvent) returns (ResponseEvent) {} // unary
rpc serverStream(RequestEvent) returns (stream ResponseEvent) {} // server streaming
rpc clientStream(stream RequestEvent) returns (ResponseEvent) {} // client streaming
rpc biDiStream(stream RequestEvent) returns (stream ResponseEvent) {} // bidirectional streaming
}
Requests are built with the generated message builders (RequestEvent.newBuilder()).
For Java, unary and server-streaming calls are available on the BlockingStub, while client-streaming and bidirectional calls require the async Stub
(they have no blocking variant). The Kotlin coroutine stub expresses every style with suspend functions and Flow<T>:
var request = RequestEvent.newBuilder().setName("bob").setCode("b1").build();
// unary — BlockingStub
ResponseEvent unary = blockingStub.unary(request);
// unary — FutureStub
ListenableFuture<ResponseEvent> future = futureStub.unary(request);
// server streaming — BlockingStub returns an iterator
Iterator<ResponseEvent> responses = blockingStub.serverStream(request);
// server streaming — async Stub delivers results to a StreamObserver
asyncStub.serverStream(request, new StreamObserver<>() {
@Override public void onNext(ResponseEvent value) { /* ... */ }
@Override public void onError(Throwable t) { /* ... */ }
@Override public void onCompleted() { /* ... */ }
});
// client streaming — async Stub, write requests, read one response
StreamObserver<ResponseEvent> responseObserver = new StreamObserver<>() {
@Override public void onNext(ResponseEvent value) { /* single response */ }
@Override public void onError(Throwable t) { /* ... */ }
@Override public void onCompleted() { /* ... */ }
};
StreamObserver<RequestEvent> requestObserver = asyncStub.clientStream(responseObserver);
requestObserver.onNext(request);
requestObserver.onCompleted();
// bidirectional streaming — async Stub, both sides stream
StreamObserver<RequestEvent> bidi = asyncStub.biDiStream(responseObserver);
bidi.onNext(request);
bidi.onCompleted();
val request = RequestEvent.newBuilder().setName("bob").setCode("b1").build()
// unary — suspend function
val response: ResponseEvent = coroutineStub.unary(request)
// server streaming — returns a Flow
val serverFlow: Flow<ResponseEvent> = coroutineStub.serverStream(request)
serverFlow.collect { event -> /* ... */ }
// client streaming — accepts a Flow, returns a single response
val clientResponse: ResponseEvent = coroutineStub.clientStream(flowOf(request))
// bidirectional streaming — Flow in, Flow out
val biDiFlow: Flow<ResponseEvent> = coroutineStub.biDiStream(flowOf(request))
biDiFlow.collect { event -> /* ... */ }
Injecting Channel and config¶
For advanced or manual stub construction you can inject the raw io.grpc.Channel and the resolved GrpcClientConfig
by tagging them with the generated service class. Both are provided by 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());
}
}
Interceptors¶
Interceptors allow you to intercept requests before they are passed to services.
Default¶
The following interceptors are used at client startup by default:
GrpcClientConfigInterceptor— appliestimeoutas the calldeadlinewhen the call has none.GrpcClientTelemetryInterceptor, if telemetry is available for the client.
Custom¶
Unlike the HTTP client, gRPC interceptors have no method/class/global tiers.
Every interceptor is scoped per client by tagging the component with the generated service class (@Tag(SimpleServiceGrpc.class)).
Register the interceptor as a component with that tag:
@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);
}
}
To apply one interceptor bean to several clients (a "shared" interceptor), give it multiple @Tag values — one generated service class per client:
@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)
}
}
Execution order:
ManagedChannelLifecycle collects all interceptors tagged for the service as All<ClientInterceptor> and applies them in a fixed order:
your custom interceptors first, then the telemetry interceptor (if telemetry is enabled), then the config/deadline interceptor last.
Because the deadline interceptor runs last, a deadline that a custom interceptor sets on the CallOptions is preserved, and the configured timeout
is applied only when no earlier interceptor (or per-call withDeadlineAfter) provided one.
Alternatively, you can modify the stub with GraphInterceptor.
Authorization¶
gRPC has no dedicated authorization module: authorization is done with a ClientInterceptor tagged with the service class that attaches an
Authorization (or API-key) header to the outgoing call Metadata. The interceptor wraps the call in a ForwardingClientCall.SimpleForwardingClientCall
and puts the header in start(...), before the request is sent.
Bearer¶
A Bearer interceptor reads a token from your own provider and
puts it on the Authorization header of every call. TokenProvider below is your own component:
@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¶
An API-key interceptor puts a static key on a custom metadata header (for example X-API-KEY).
The key is read from a @ConfigSource interface injected into the interceptor:
@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);
}
};
}
}
- Any
@ConfigSourceinterface exposing the API key, for exampleString 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)
}
}
- Any
@ConfigSourceinterface exposing the API key, for examplefun apiKey(): String
Error handling¶
A failed gRPC call throws an io.grpc.StatusRuntimeException. Its getStatus() carries a Status.Code
(status codes) such as UNAVAILABLE (server unreachable), DEADLINE_EXCEEDED (the timeout/deadline expired),
UNAUTHENTICATED (rejected credentials), or INVALID_ARGUMENT. Response metadata is available through getTrailers().
Causes:
UNAVAILABLE— wrongurl, plaintext/TLS mismatch, or the server is down.DEADLINE_EXCEEDED— the configuredtimeoutor a per-callwithDeadlineAfterwas exceeded.UNAUTHENTICATED/PERMISSION_DENIED— missing or invalid authorization metadata.
Recommendations:
- Branch on
e.getStatus().getCode()instead of the exception type. - Use resilient aspects (
@Retry,@CircuitBreaker,@Timeout) on the wrapping service method for transient failures.
Telemetry¶
gRPC Client uses a telemetry contract for logging, metrics, and tracing of calls.
Telemetry configuration (section telemetry { logging / metrics / tracing }) is described in the Configuration section.
Extension points are located in ru.tinkoff.kora.grpc.client.common.telemetry.
For each gRPC call, a GrpcClientTelemetry.GrpcClientTelemetryContext is created, which is closed upon call completion.
The call is described through telemetry handler parameters, including service, method, response status, and duration.
The default factory DefaultGrpcClientTelemetryFactory combines three factories:
- GrpcClientLoggerFactory builds GrpcClientLogger for logging call start/end;
- GrpcClientMetricsFactory builds GrpcClientMetrics for writing call metrics;
- GrpcClientTracerFactory builds GrpcClientTracer for distributed tracing.
Metrics and tracing are described in the Metrics Reference section.