Health check for InfluxDB (#5637)
This commit is contained in:
parent
494c146ed5
commit
d86570a88e
|
|
@ -1,32 +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.
|
||||
*
|
||||
*/
|
||||
|
||||
package org.apache.skywalking.oap.server.library.client.elasticsearch;
|
||||
|
||||
import org.apache.skywalking.oap.server.library.client.ClientException;
|
||||
|
||||
public class ElasticSearchClientException extends ClientException {
|
||||
|
||||
public ElasticSearchClientException(String message) {
|
||||
super(message);
|
||||
}
|
||||
|
||||
public ElasticSearchClientException(String message, Throwable cause) {
|
||||
super(message, cause);
|
||||
}
|
||||
}
|
||||
|
|
@ -38,7 +38,11 @@
|
|||
<artifactId>server-core</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>org.apache.skywalking</groupId>
|
||||
<artifactId>library-client</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.influxdb</groupId>
|
||||
<artifactId>influxdb-java</artifactId>
|
||||
|
|
|
|||
|
|
@ -18,15 +18,14 @@
|
|||
|
||||
package org.apache.skywalking.oap.server.storage.plugin.influxdb;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import okhttp3.OkHttpClient;
|
||||
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
|
||||
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.HealthCheckable;
|
||||
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
|
||||
import org.apache.skywalking.oap.server.library.util.HealthChecker;
|
||||
import org.influxdb.InfluxDB;
|
||||
import org.influxdb.InfluxDBFactory;
|
||||
import org.influxdb.dto.BatchPoints;
|
||||
|
|
@ -35,16 +34,21 @@ import org.influxdb.dto.Query;
|
|||
import org.influxdb.dto.QueryResult;
|
||||
import org.influxdb.querybuilder.time.TimeInterval;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.influxdb.querybuilder.BuiltQuery.QueryBuilder.ti;
|
||||
|
||||
/**
|
||||
* InfluxDB connection maintainer, provides base data write/query API.
|
||||
*/
|
||||
@Slf4j
|
||||
public class InfluxClient implements Client {
|
||||
public class InfluxClient implements Client, HealthCheckable {
|
||||
private InfluxStorageConfig config;
|
||||
private InfluxDB influx;
|
||||
|
||||
private DelegatedHealthChecker healthChecker = new DelegatedHealthChecker();
|
||||
/**
|
||||
* A constant, the name of time field in Time-series database.
|
||||
*/
|
||||
|
|
@ -67,16 +71,22 @@ public class InfluxClient implements Client {
|
|||
|
||||
@Override
|
||||
public void connect() {
|
||||
influx = InfluxDBFactory.connect(config.getUrl(), config.getUser(), config.getPassword(),
|
||||
new OkHttpClient.Builder().readTimeout(3, TimeUnit.MINUTES)
|
||||
.writeTimeout(3, TimeUnit.MINUTES),
|
||||
InfluxDB.ResponseFormat.MSGPACK
|
||||
);
|
||||
influx.query(new Query("CREATE DATABASE " + database));
|
||||
influx.enableGzip();
|
||||
try {
|
||||
influx = InfluxDBFactory.connect(config.getUrl(), config.getUser(), config.getPassword(),
|
||||
new OkHttpClient.Builder().readTimeout(3, TimeUnit.MINUTES)
|
||||
.writeTimeout(3, TimeUnit.MINUTES),
|
||||
InfluxDB.ResponseFormat.MSGPACK
|
||||
);
|
||||
influx.query(new Query("CREATE DATABASE " + database));
|
||||
influx.enableGzip();
|
||||
|
||||
influx.enableBatch(config.getActions(), config.getDuration(), TimeUnit.MILLISECONDS);
|
||||
influx.setDatabase(database);
|
||||
influx.enableBatch(config.getActions(), config.getDuration(), TimeUnit.MILLISECONDS);
|
||||
influx.setDatabase(database);
|
||||
healthChecker.health();
|
||||
} catch (Throwable e) {
|
||||
healthChecker.unHealth(e);
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -98,14 +108,15 @@ public class InfluxClient implements Client {
|
|||
if (log.isDebugEnabled()) {
|
||||
log.debug("SQL Statement: {}", query.getCommand());
|
||||
}
|
||||
|
||||
try {
|
||||
QueryResult result = getInflux().query(new Query(query.getCommand()));
|
||||
if (result.hasError()) {
|
||||
throw new IOException(result.getError());
|
||||
}
|
||||
healthChecker.health();
|
||||
return result.getResults();
|
||||
} catch (Exception e) {
|
||||
} catch (Throwable e) {
|
||||
healthChecker.unHealth(e);
|
||||
throw new IOException(e.getMessage() + System.lineSeparator() + "SQL Statement: " + query.getCommand(), e);
|
||||
}
|
||||
}
|
||||
|
|
@ -161,11 +172,7 @@ public class InfluxClient implements Client {
|
|||
*/
|
||||
public void dropSeries(String measurement, long timeBucket) throws IOException {
|
||||
Query query = new Query("DROP SERIES FROM " + measurement + " WHERE time_bucket='" + timeBucket + "'");
|
||||
QueryResult result = getInflux().query(query);
|
||||
|
||||
if (result.hasError()) {
|
||||
throw new IOException("Statement: " + query.getCommand() + ", ErrorMsg: " + result.getError());
|
||||
}
|
||||
this.query(query);
|
||||
}
|
||||
|
||||
public void deleteByQuery(String measurement, long timestamp) throws IOException {
|
||||
|
|
@ -177,19 +184,37 @@ public class InfluxClient implements Client {
|
|||
* wait for buffer flushing.
|
||||
*/
|
||||
public void write(Point point) {
|
||||
getInflux().write(point);
|
||||
try {
|
||||
getInflux().write(point);
|
||||
this.healthChecker.health();
|
||||
} catch (Throwable e) {
|
||||
healthChecker.unHealth(e);
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* A batch operation of write. {@link Point}s flush directly.
|
||||
*/
|
||||
public void write(BatchPoints points) {
|
||||
getInflux().write(points);
|
||||
try {
|
||||
getInflux().write(points);
|
||||
this.healthChecker.health();
|
||||
} catch (Throwable e) {
|
||||
healthChecker.unHealth(e);
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() throws IOException {
|
||||
influx.close();
|
||||
try {
|
||||
getInflux().close();
|
||||
this.healthChecker.health();
|
||||
} catch (Throwable e) {
|
||||
healthChecker.unHealth(e);
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -205,4 +230,9 @@ public class InfluxClient implements Client {
|
|||
public static TimeInterval timeIntervalTB(long timeBucket) {
|
||||
return ti(TimeBucket.getTimestamp(timeBucket), "ms");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void registerChecker(HealthChecker healthChecker) {
|
||||
this.healthChecker.register(healthChecker);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -62,6 +62,10 @@ import org.apache.skywalking.oap.server.storage.plugin.influxdb.query.TopNRecord
|
|||
import org.apache.skywalking.oap.server.storage.plugin.influxdb.query.TopologyQuery;
|
||||
import org.apache.skywalking.oap.server.storage.plugin.influxdb.query.TraceQuery;
|
||||
import org.apache.skywalking.oap.server.storage.plugin.influxdb.query.UITemplateManagementDAOImpl;
|
||||
import org.apache.skywalking.oap.server.telemetry.TelemetryModule;
|
||||
import org.apache.skywalking.oap.server.telemetry.api.HealthCheckMetrics;
|
||||
import org.apache.skywalking.oap.server.telemetry.api.MetricsCreator;
|
||||
import org.apache.skywalking.oap.server.telemetry.api.MetricsTag;
|
||||
|
||||
@Slf4j
|
||||
public class InfluxStorageProvider extends ModuleProvider {
|
||||
|
|
@ -119,6 +123,9 @@ public class InfluxStorageProvider extends ModuleProvider {
|
|||
|
||||
@Override
|
||||
public void start() throws ServiceNotProvidedException, ModuleStartException {
|
||||
MetricsCreator metricCreator = getManager().find(TelemetryModule.NAME).provider().getService(MetricsCreator.class);
|
||||
HealthCheckMetrics healthChecker = metricCreator.createHealthCheckerGauge("storage_influxdb", MetricsTag.EMPTY_KEY, MetricsTag.EMPTY_VALUE);
|
||||
client.registerChecker(healthChecker);
|
||||
try {
|
||||
client.connect();
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue