Support custom GRPCClient health checker logic. (#13353)

[NOTICE] Roll back score meaning in GraphQL health check API.
This commit is contained in:
Wan Kai 2025-07-03 15:28:26 +08:00 committed by GitHub
parent 4430cf5e0d
commit 4b508b930d
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
14 changed files with 116 additions and 97 deletions

View File

@ -226,7 +226,7 @@ The text of each license is the standard Apache 2.0 license.
https://mvnrepository.com/artifact/com.fasterxml.jackson.datatype/jackson-datatype-jsr310/2.18.2 Apache-2.0 https://mvnrepository.com/artifact/com.fasterxml.jackson.datatype/jackson-datatype-jsr310/2.18.2 Apache-2.0
https://mvnrepository.com/artifact/com.fasterxml.jackson.module/jackson-module-kotlin/2.13.4 Apache-2.0 https://mvnrepository.com/artifact/com.fasterxml.jackson.module/jackson-module-kotlin/2.13.4 Apache-2.0
https://mvnrepository.com/artifact/com.fasterxml/classmate/1.5.1 Apache-2.0 https://mvnrepository.com/artifact/com.fasterxml/classmate/1.5.1 Apache-2.0
https://mvnrepository.com/artifact/com.google.api.grpc/proto-google-common-protos/2.41.0 Apache-2.0 https://mvnrepository.com/artifact/com.google.api.grpc/proto-google-common-protos/2.48.0 Apache-2.0
https://mvnrepository.com/artifact/com.google.auto.service/auto-service-annotations/1.0.1 Apache-2.0 https://mvnrepository.com/artifact/com.google.auto.service/auto-service-annotations/1.0.1 Apache-2.0
https://mvnrepository.com/artifact/com.google.code.findbugs/jsr305/3.0.2 Apache-2.0 https://mvnrepository.com/artifact/com.google.code.findbugs/jsr305/3.0.2 Apache-2.0
https://mvnrepository.com/artifact/com.google.code.gson/gson/2.9.0 Apache-2.0 https://mvnrepository.com/artifact/com.google.code.gson/gson/2.9.0 Apache-2.0
@ -290,16 +290,16 @@ The text of each license is the standard Apache 2.0 license.
https://mvnrepository.com/artifact/io.fabric8/kubernetes-model-scheduling/6.7.1 Apache-2.0 https://mvnrepository.com/artifact/io.fabric8/kubernetes-model-scheduling/6.7.1 Apache-2.0
https://mvnrepository.com/artifact/io.fabric8/kubernetes-model-storageclass/6.7.1 Apache-2.0 https://mvnrepository.com/artifact/io.fabric8/kubernetes-model-storageclass/6.7.1 Apache-2.0
https://mvnrepository.com/artifact/io.fabric8/zjsonpatch/0.3.0 Apache-2.0 https://mvnrepository.com/artifact/io.fabric8/zjsonpatch/0.3.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-api/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-api/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-context/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-context/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-core/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-core/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-grpclb/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-grpclb/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-netty/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-netty/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-protobuf/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-protobuf/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-protobuf-lite/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-protobuf-lite/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-services/1.70.0 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-services/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-stub/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-stub/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.grpc/grpc-util/1.68.1 Apache-2.0 https://mvnrepository.com/artifact/io.grpc/grpc-util/1.70.0 Apache-2.0
https://mvnrepository.com/artifact/io.micrometer/micrometer-commons/1.14.4 Apache-2.0 https://mvnrepository.com/artifact/io.micrometer/micrometer-commons/1.14.4 Apache-2.0
https://mvnrepository.com/artifact/io.micrometer/micrometer-core/1.14.4 Apache-2.0 https://mvnrepository.com/artifact/io.micrometer/micrometer-core/1.14.4 Apache-2.0
https://mvnrepository.com/artifact/io.micrometer/micrometer-observation/1.14.4 Apache-2.0 https://mvnrepository.com/artifact/io.micrometer/micrometer-observation/1.14.4 Apache-2.0

View File

@ -37,7 +37,8 @@
* chore: add a warning log when connecting to ES takes too long. * chore: add a warning log when connecting to ES takes too long.
* Fix the query time range in the metadata API. * Fix the query time range in the metadata API.
* OAP gRPC-Client support `Health Check`. * OAP gRPC-Client support `Health Check`.
* [Break Change] `Health Check` make response 1 represents healthy, 0 represents unhealthy. * [Break Change] `health_check_xx` metrics make response 1 represents healthy, 0 represents unhealthy.
* Bump up grpc to 1.70.0.
#### UI #### UI

View File

@ -36,7 +36,7 @@ If the OAP server is healthy, the response should be
{ {
"data": { "data": {
"checkHealth": { "checkHealth": {
"score": 1, "score": 0,
"details": "" "details": ""
} }
} }
@ -49,7 +49,7 @@ If some modules are unhealthy (e.g. storage H2 is down), then the result may loo
{ {
"data": { "data": {
"checkHealth": { "checkHealth": {
"score": 0, "score": 1,
"details": "storage_h2," "details": "storage_h2,"
} }
} }

View File

@ -253,6 +253,11 @@
<artifactId>grpc-stub</artifactId> <artifactId>grpc-stub</artifactId>
<version>${grpc.version}</version> <version>${grpc.version}</version>
</dependency> </dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-services</artifactId>
<version>${grpc.version}</version>
</dependency>
<dependency> <dependency>
<groupId>io.netty</groupId> <groupId>io.netty</groupId>
<artifactId>netty-tcnative-boringssl-static</artifactId> <artifactId>netty-tcnative-boringssl-static</artifactId>

View File

@ -26,7 +26,7 @@ import lombok.ToString;
@Setter @Setter
@ToString @ToString
public class HealthStatus { public class HealthStatus {
// score == 1 means healthy, otherwise it's unhealthy. // score == 0 means healthy and no unhealthy component or connection, otherwise it's unhealthy.
private int score; private int score;
private String details; private String details;
} }

View File

@ -18,17 +18,15 @@
package org.apache.skywalking.oap.server.core.remote.health; package org.apache.skywalking.oap.server.core.remote.health;
import grpc.health.v1.HealthCheckService; import io.grpc.health.v1.HealthCheckRequest;
import grpc.health.v1.HealthGrpc; import io.grpc.health.v1.HealthCheckResponse;
import io.grpc.health.v1.HealthGrpc;
import io.grpc.stub.StreamObserver; import io.grpc.stub.StreamObserver;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.oap.server.library.server.grpc.GRPCHandler; import org.apache.skywalking.oap.server.library.server.grpc.GRPCHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@Slf4j
public class HealthCheckServiceHandler extends HealthGrpc.HealthImplBase implements GRPCHandler { public class HealthCheckServiceHandler extends HealthGrpc.HealthImplBase implements GRPCHandler {
private static final Logger LOGGER = LoggerFactory.getLogger(HealthCheckServiceHandler.class);
/** /**
* By my test, consul didn't send the service. * By my test, consul didn't send the service.
* *
@ -36,15 +34,13 @@ public class HealthCheckServiceHandler extends HealthGrpc.HealthImplBase impleme
* @param responseObserver status * @param responseObserver status
*/ */
@Override @Override
public void check(HealthCheckService.HealthCheckRequest request, public void check(HealthCheckRequest request, StreamObserver<HealthCheckResponse> responseObserver) {
StreamObserver<HealthCheckService.HealthCheckResponse> responseObserver) { if (log.isDebugEnabled()) {
log.debug("Received the gRPC server health check with the service name of {}", request.getService());
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Received the gRPC server health check with the service name of {}", request.getService());
} }
HealthCheckService.HealthCheckResponse.Builder response = HealthCheckService.HealthCheckResponse.newBuilder(); HealthCheckResponse.Builder response = HealthCheckResponse.newBuilder();
response.setStatus(HealthCheckService.HealthCheckResponse.ServingStatus.SERVING); response.setStatus(HealthCheckResponse.ServingStatus.SERVING);
responseObserver.onNext(response.build()); responseObserver.onNext(response.build());
responseObserver.onCompleted(); responseObserver.onCompleted();

View File

@ -1,40 +0,0 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
//This is a health check proto, provided by gRPC team. Please don't change it.
//https://github.com/grpc/grpc/blob/master/doc/health-checking.md
syntax = "proto3";
package grpc.health.v1;
message HealthCheckRequest {
string service = 1;
}
message HealthCheckResponse {
enum ServingStatus {
UNKNOWN = 0;
SERVING = 1;
NOT_SERVING = 2;
}
ServingStatus status = 1;
}
service Health {
rpc Check (HealthCheckRequest) returns (HealthCheckResponse);
}

View File

@ -37,7 +37,7 @@ public class HealthCheckerHttpService {
final var status = healthQueryService.checkHealth(); final var status = healthQueryService.checkHealth();
log.info("Health status: {}", status); log.info("Health status: {}", status);
if (status.getScore() == 1) { if (status.getScore() == 0) {
return HttpResponse.of(HttpStatus.OK); return HttpResponse.of(HttpStatus.OK);
} }
return HttpResponse.of(HttpStatus.SERVICE_UNAVAILABLE); return HttpResponse.of(HttpStatus.SERVICE_UNAVAILABLE);

View File

@ -25,7 +25,6 @@ import java.util.Arrays;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.CoreModule;
@ -100,18 +99,22 @@ public class HealthCheckerProvider extends ModuleProvider {
@Override public void notifyAfterCompleted() throws ServiceNotProvidedException, ModuleStartException { @Override public void notifyAfterCompleted() throws ServiceNotProvidedException, ModuleStartException {
ses.scheduleAtFixedRate(() -> { ses.scheduleAtFixedRate(() -> {
StringBuilder unhealthyModules = new StringBuilder(); StringBuilder unhealthyModules = new StringBuilder();
AtomicBoolean hasUnhealthyModule = new AtomicBoolean(false); AtomicDouble unhealthyModule = new AtomicDouble(0);
Stream.ofAll(collector.collect()) Stream.ofAll(collector.collect())
.flatMap(metricFamily -> metricFamily.samples) .flatMap(metricFamily -> metricFamily.samples)
.filter(sample -> metricsCreator.isHealthCheckerMetrics(sample.name)) .filter(sample -> metricsCreator.isHealthCheckerMetrics(sample.name))
.forEach(sample -> { .forEach(sample -> {
if (sample.value < 1) { if (sample.value < 1) {
unhealthyModules.append(metricsCreator.extractModuleName(sample.name)).append(","); unhealthyModules.append(metricsCreator.extractModuleName(sample.name)).append(",");
hasUnhealthyModule.set(true); unhealthyModule.updateAndGet(v -> v + 1);
} }
}); });
score.set(hasUnhealthyModule.get() ? 0 : 1); if (unhealthyModule.get() > 0) {
score.set(unhealthyModule.get());
} else {
score.set(0);
}
details.set(unhealthyModules.toString()); details.set(unhealthyModules.toString());
}, },
2, config.getCheckIntervalSeconds(), TimeUnit.SECONDS); 2, config.getCheckIntervalSeconds(), TimeUnit.SECONDS);

View File

@ -53,6 +53,10 @@
<groupId>io.grpc</groupId> <groupId>io.grpc</groupId>
<artifactId>grpc-netty</artifactId> <artifactId>grpc-netty</artifactId>
</dependency> </dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-services</artifactId>
</dependency>
<dependency> <dependency>
<groupId>io.netty</groupId> <groupId>io.netty</groupId>
<artifactId>netty-codec-http2</artifactId> <artifactId>netty-codec-http2</artifactId>

View File

@ -18,26 +18,27 @@
package org.apache.skywalking.oap.server.library.client.grpc; package org.apache.skywalking.oap.server.library.client.grpc;
import io.grpc.ConnectivityState;
import io.grpc.ManagedChannel; import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder; import io.grpc.ManagedChannelBuilder;
import io.grpc.Status;
import io.grpc.StatusRuntimeException;
import io.grpc.health.v1.HealthCheckRequest;
import io.grpc.health.v1.HealthCheckResponse;
import io.grpc.health.v1.HealthGrpc;
import io.grpc.netty.NettyChannelBuilder; import io.grpc.netty.NettyChannelBuilder;
import io.netty.handler.ssl.SslContext; import io.netty.handler.ssl.SslContext;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import lombok.Getter; import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.oap.server.library.client.Client; import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.client.healthcheck.DelegatedHealthChecker; import org.apache.skywalking.oap.server.library.client.healthcheck.DelegatedHealthChecker;
import org.apache.skywalking.oap.server.library.client.healthcheck.HealthCheckable; import org.apache.skywalking.oap.server.library.client.healthcheck.HealthCheckable;
import org.apache.skywalking.oap.server.library.util.HealthChecker; import org.apache.skywalking.oap.server.library.util.HealthChecker;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@Slf4j
public class GRPCClient implements Client, HealthCheckable { public class GRPCClient implements Client, HealthCheckable {
private static final Logger LOGGER = LoggerFactory.getLogger(GRPCClient.class);
@Getter @Getter
private final String host; private final String host;
@ -54,6 +55,35 @@ public class GRPCClient implements Client, HealthCheckable {
private boolean enableHealthCheck = false; private boolean enableHealthCheck = false;
private long initialDelay = 5; // Initial delay for health check in seconds
private long period = 20; // Period for health check in seconds
// The default health check runnable that checks the health of the gRPC channel.
private Runnable healthCheckRunnable = () -> {
if (getChannel() != null && !getChannel().isShutdown()) {
HealthGrpc.HealthBlockingStub healthStub = HealthGrpc.newBlockingStub(getChannel());
HealthCheckRequest request = HealthCheckRequest.newBuilder().setService("").build();
try {
HealthCheckResponse response = healthStub.check(request);
handleStateChange(response);
} catch (StatusRuntimeException s) {
if (s.getStatus().getCode() == Status.Code.UNIMPLEMENTED) {
log.warn("Health check is not implemented on the remote gRPC server, regard as healthy. Host: {}, Port: {}", getHost(), getPort());
healthChecker.health();
} else {
log.warn("Health check failed for gRPC channel. Host: {}, Port: {}", getHost(), getPort(), s);
healthChecker.unHealth(s);
}
} catch (Throwable t) {
log.warn("Health check failed for gRPC channel. Host: {}, Port: {}", getHost(), getPort(), t);
healthChecker.unHealth(t);
}
} else {
healthChecker.unHealth("gRPC channel is not available or shutting down. Host: " + getHost() + ", Port: " + getPort());
}
};
public GRPCClient(String host, int port) { public GRPCClient(String host, int port) {
this.host = host; this.host = host;
this.port = port; this.port = port;
@ -81,7 +111,7 @@ public class GRPCClient implements Client, HealthCheckable {
try { try {
channel.shutdownNow(); channel.shutdownNow();
} catch (Throwable t) { } catch (Throwable t) {
LOGGER.error(t.getMessage(), t); log.error(t.getMessage(), t);
} finally { } finally {
if (healthCheckExecutor != null) { if (healthCheckExecutor != null) {
healthCheckExecutor.shutdownNow(); healthCheckExecutor.shutdownNow();
@ -114,32 +144,51 @@ public class GRPCClient implements Client, HealthCheckable {
this.enableHealthCheck = true; this.enableHealthCheck = true;
} }
/**
* Override the default health check runnable with a custom one.
* Must override before calling connect()
* This can be used to provide a different health check logic.
*
* @param healthCheckRunnable The custom health check runnable.
* @param initialDelay Initial delay before the first health check.
* @param period Period between subsequent health checks.
*/
public void overrideCheckerRunnable(final Runnable healthCheckRunnable, final long initialDelay, final long period) {
this.healthCheckRunnable = healthCheckRunnable;
if (initialDelay < 0) {
throw new IllegalArgumentException("initialDelay must be non-negative. Provided value: " + initialDelay);
}
if (period < 0) {
throw new IllegalArgumentException("period must be non-negative. Provided value: " + period);
}
this.initialDelay = initialDelay;
this.period = period;
}
private void checkHealth() { private void checkHealth() {
if (healthCheckExecutor == null) { if (healthCheckExecutor == null) {
healthCheckExecutor = Executors.newSingleThreadScheduledExecutor(); healthCheckExecutor = Executors.newSingleThreadScheduledExecutor();
healthCheckExecutor.scheduleAtFixedRate( healthCheckExecutor.scheduleAtFixedRate(healthCheckRunnable, initialDelay, period, TimeUnit.SECONDS
() -> {
ConnectivityState currentState = channel.getState(true); // true means try to connect
handleStateChange(currentState);
}, 5, 10, TimeUnit.SECONDS
); );
} }
} }
private void handleStateChange(ConnectivityState newState) { private void handleStateChange(HealthCheckResponse response) {
switch (newState) { switch (response.getStatus()) {
case READY: case SERVING:
case IDLE:
this.healthChecker.health(); this.healthChecker.health();
break; break;
case CONNECTING: case NOT_SERVING:
this.healthChecker.unHealth("gRPC connecting, waiting for ready. Host: " + host + ", Port: " + port); this.healthChecker.unHealth("Remote gRPC Server NOT_SERVING. Host: " + host + ", Port: " + port);
break; break;
case TRANSIENT_FAILURE: case SERVICE_UNKNOWN:
this.healthChecker.unHealth("gRPC connection failed, will retry. Host: " + host + ", Port: " + port); this.healthChecker.unHealth("Remote gRPC Server SERVICE_UNKNOWN. Host: " + host + ", Port: " + port);
break; break;
case SHUTDOWN: case UNKNOWN:
this.healthChecker.unHealth("gRPC channel is shutting down. Host: " + host + ", Port: " + port); this.healthChecker.unHealth("Remote gRPC Server UNKNOWN. Host: " + host + ", Port: " + port);
break;
case UNRECOGNIZED:
this.healthChecker.unHealth("Remote gRPC Server UNRECOGNIZED. Host: " + host + ", Port: " + port);
break; break;
} }
} }

@ -1 +1 @@
Subproject commit 021c0ad768f8f6f64dceead9d79a3dd7e9ad8dd9 Subproject commit e6d2b99597a6489ab58f8a2ff32f212c8ef9b388

View File

@ -317,6 +317,7 @@
<exclude>hierarchy-definition.yml</exclude> <exclude>hierarchy-definition.yml</exclude>
<exclude>bydb.dependencies.properties</exclude> <exclude>bydb.dependencies.properties</exclude>
<exclude>bydb.yml</exclude> <exclude>bydb.yml</exclude>
<exclude>bydb-topn.yml</exclude>
<exclude>oal/</exclude> <exclude>oal/</exclude>
<exclude>fetcher-prom-rules/</exclude> <exclude>fetcher-prom-rules/</exclude>
<exclude>envoy-metrics-rules/</exclude> <exclude>envoy-metrics-rules/</exclude>

View File

@ -165,7 +165,7 @@
<byte-buddy.version>1.14.9</byte-buddy.version> <byte-buddy.version>1.14.9</byte-buddy.version>
<!-- core lib dependency --> <!-- core lib dependency -->
<grpc.version>1.68.1</grpc.version> <grpc.version>1.70.0</grpc.version>
<netty.version>4.1.118.Final</netty.version> <netty.version>4.1.118.Final</netty.version>
<netty-tcnative-boringssl-static.version>2.0.69.Final</netty-tcnative-boringssl-static.version> <netty-tcnative-boringssl-static.version>2.0.69.Final</netty-tcnative-boringssl-static.version>
<gson.version>2.9.0</gson.version> <gson.version>2.9.0</gson.version>