mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-19 02:23:24 +08:00
[ISSUE ##5354] Implement broker metrics framework (#5355)
* implement broker metrics framework * update bazel config * Fix Bazel deps * Fix unnecessary mocking issue * Fix bazel compile warning Co-authored-by: Zhanhui Li <lizhanhui@apache.org>
This commit is contained in:
@@ -87,6 +87,13 @@ maven_install(
|
||||
"io.grpc:grpc-api:1.47.0",
|
||||
"io.grpc:grpc-testing:1.47.0",
|
||||
"org.springframework:spring-core:5.3.23",
|
||||
"io.opentelemetry:opentelemetry-exporter-otlp:1.19.0",
|
||||
"io.opentelemetry:opentelemetry-exporter-prometheus:1.19.0-alpha",
|
||||
"io.opentelemetry:opentelemetry-sdk:1.19.0",
|
||||
"com.squareup.okio:okio-jvm:3.0.0",
|
||||
"io.opentelemetry:opentelemetry-api:1.19.0",
|
||||
"io.opentelemetry:opentelemetry-sdk-metrics:1.19.0",
|
||||
"io.opentelemetry:opentelemetry-sdk-common:1.19.0",
|
||||
],
|
||||
fetch_sources = True,
|
||||
repositories = [
|
||||
|
||||
+40
-34
@@ -21,55 +21,61 @@ java_library(
|
||||
srcs = glob(["src/main/java/**/*.java"]),
|
||||
visibility = ["//visibility:public"],
|
||||
deps = [
|
||||
"//remoting",
|
||||
"//logging",
|
||||
"//common",
|
||||
"//store",
|
||||
"//client",
|
||||
"//filter",
|
||||
"//srvutil",
|
||||
"//acl",
|
||||
"@maven//:io_openmessaging_storage_dledger",
|
||||
"@maven//:org_apache_commons_commons_lang3",
|
||||
"@maven//:commons_validator_commons_validator",
|
||||
"@maven//:com_github_luben_zstd_jni",
|
||||
"@maven//:org_lz4_lz4_java",
|
||||
"@maven//:com_alibaba_fastjson",
|
||||
"@maven//:io_netty_netty_all",
|
||||
"//client",
|
||||
"//common",
|
||||
"//filter",
|
||||
"//logging",
|
||||
"//remoting",
|
||||
"//srvutil",
|
||||
"//store",
|
||||
"@maven//:ch_qos_logback_logback_classic",
|
||||
"@maven//:org_slf4j_slf4j_api",
|
||||
"@maven//:commons_cli_commons_cli",
|
||||
"@maven//:com_alibaba_fastjson",
|
||||
"@maven//:com_github_luben_zstd_jni",
|
||||
"@maven//:com_google_guava_guava",
|
||||
"@maven//:com_googlecode_concurrentlinkedhashmap_concurrentlinkedhashmap_lru",
|
||||
"@maven//:commons_io_commons_io",
|
||||
"@maven//:commons_cli_commons_cli",
|
||||
"@maven//:commons_collections_commons_collections",
|
||||
"@maven//:commons_io_commons_io",
|
||||
"@maven//:commons_validator_commons_validator",
|
||||
"@maven//:io_netty_netty_all",
|
||||
"@maven//:io_openmessaging_storage_dledger",
|
||||
"@maven//:io_opentelemetry_opentelemetry_api",
|
||||
"@maven//:io_opentelemetry_opentelemetry_exporter_otlp",
|
||||
"@maven//:io_opentelemetry_opentelemetry_exporter_prometheus",
|
||||
"@maven//:io_opentelemetry_opentelemetry_sdk",
|
||||
"@maven//:io_opentelemetry_opentelemetry_sdk_common",
|
||||
"@maven//:io_opentelemetry_opentelemetry_sdk_metrics",
|
||||
"@maven//:org_apache_commons_commons_lang3",
|
||||
"@maven//:org_lz4_lz4_java",
|
||||
"@maven//:org_slf4j_slf4j_api",
|
||||
],
|
||||
)
|
||||
|
||||
java_library(
|
||||
name = "tests",
|
||||
srcs = glob(["src/test/java/**/*.java"]),
|
||||
visibility = ["//visibility:public"],
|
||||
deps = [
|
||||
":broker",
|
||||
"//acl",
|
||||
"//client",
|
||||
"//filter",
|
||||
"//logging",
|
||||
"//store",
|
||||
"//common",
|
||||
"//remoting",
|
||||
"//:test_deps",
|
||||
"@maven//:org_apache_commons_commons_lang3",
|
||||
"@maven//:io_netty_netty_all",
|
||||
"@maven//:com_google_guava_guava",
|
||||
"@maven//:com_alibaba_fastjson",
|
||||
],
|
||||
resources = [
|
||||
"src/test/resources/logback-test.xml",
|
||||
"src/test/resources/META-INF/service/org.apache.rocketmq.acl.AccessValidator",
|
||||
"src/test/resources/META-INF/service/org.apache.rocketmq.broker.transaction.AbstractTransactionalMessageCheckListener",
|
||||
"src/test/resources/META-INF/service/org.apache.rocketmq.broker.transaction.TransactionalMessageService",
|
||||
"src/test/resources/logback-test.xml",
|
||||
],
|
||||
visibility = ["//visibility:public"],
|
||||
deps = [
|
||||
":broker",
|
||||
"//:test_deps",
|
||||
"//acl",
|
||||
"//client",
|
||||
"//common",
|
||||
"//filter",
|
||||
"//logging",
|
||||
"//remoting",
|
||||
"//store",
|
||||
"@maven//:com_alibaba_fastjson",
|
||||
"@maven//:com_google_guava_guava",
|
||||
"@maven//:io_netty_netty_all",
|
||||
"@maven//:org_apache_commons_commons_lang3",
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
@@ -63,6 +63,7 @@ import org.apache.rocketmq.broker.latency.BrokerFixedThreadPoolExecutor;
|
||||
import org.apache.rocketmq.broker.longpolling.LmqPullRequestHoldService;
|
||||
import org.apache.rocketmq.broker.longpolling.NotifyMessageArrivingListener;
|
||||
import org.apache.rocketmq.broker.longpolling.PullRequestHoldService;
|
||||
import org.apache.rocketmq.broker.metrics.BrokerMetricsManager;
|
||||
import org.apache.rocketmq.broker.mqtrace.ConsumeMessageHook;
|
||||
import org.apache.rocketmq.broker.mqtrace.SendMessageHook;
|
||||
import org.apache.rocketmq.broker.offset.ConsumerOffsetManager;
|
||||
@@ -262,6 +263,7 @@ public class BrokerController {
|
||||
protected final List<ScheduledFuture<?>> scheduledFutures = new ArrayList<>();
|
||||
protected ReplicasManager replicasManager;
|
||||
private long lastSyncTimeMs = System.currentTimeMillis();
|
||||
private BrokerMetricsManager brokerMetricsManager;
|
||||
|
||||
public BrokerController(
|
||||
final BrokerConfig brokerConfig,
|
||||
@@ -775,6 +777,8 @@ public class BrokerController {
|
||||
}
|
||||
}
|
||||
|
||||
this.brokerMetricsManager = new BrokerMetricsManager(this);
|
||||
|
||||
if (result) {
|
||||
|
||||
initializeRemotingServer();
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
/*
|
||||
* 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.rocketmq.broker.metrics;
|
||||
|
||||
public class BrokerMetricsConstant {
|
||||
public static final String SLS_OTEL_PROJECT_HEADER_KEY = "x-sls-otel-project";
|
||||
public static final String SLS_OTEL_INSTANCE_ID_KEY = "x-sls-otel-instance-id";
|
||||
public static final String SLS_OTEL_AK_ID_KEY = "x-sls-otel-ak-id";
|
||||
public static final String SLS_OTEL_AK_SECRET_KEY = "x-sls-otel-ak-secret";
|
||||
public static final String OPEN_TELEMETRY_METER_NAME = "broker-meter";
|
||||
public static final String GAUGE_BROKER_PERMISSION = "rocketmq_broker_permission";
|
||||
|
||||
public static final String LABEL_CLUSTER_NAME = "cluster";
|
||||
public static final String LABEL_NODE_TYPE = "node_type";
|
||||
public static final String BROKER_NODE_TYPE = "broker";
|
||||
public static final String LABEL_NODE_ID = "node_id";
|
||||
public static final String LABEL_AGGREGATION = "aggregation";
|
||||
public static final String AGGREGATION_DELTA = "delta";
|
||||
}
|
||||
@@ -0,0 +1,209 @@
|
||||
/*
|
||||
* 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.rocketmq.broker.metrics;
|
||||
|
||||
import com.google.common.base.Splitter;
|
||||
import io.opentelemetry.api.common.Attributes;
|
||||
import io.opentelemetry.api.common.AttributesBuilder;
|
||||
import io.opentelemetry.api.metrics.Meter;
|
||||
import io.opentelemetry.api.metrics.ObservableLongGauge;
|
||||
import io.opentelemetry.exporter.otlp.metrics.OtlpGrpcMetricExporter;
|
||||
import io.opentelemetry.exporter.otlp.metrics.OtlpGrpcMetricExporterBuilder;
|
||||
import io.opentelemetry.exporter.prometheus.PrometheusHttpServer;
|
||||
import io.opentelemetry.sdk.OpenTelemetrySdk;
|
||||
import io.opentelemetry.sdk.metrics.InstrumentType;
|
||||
import io.opentelemetry.sdk.metrics.SdkMeterProvider;
|
||||
import io.opentelemetry.sdk.metrics.SdkMeterProviderBuilder;
|
||||
import io.opentelemetry.sdk.metrics.data.AggregationTemporality;
|
||||
import io.opentelemetry.sdk.metrics.export.PeriodicMetricReader;
|
||||
import io.opentelemetry.sdk.resources.Resource;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.broker.BrokerController;
|
||||
import org.apache.rocketmq.common.BrokerConfig;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.store.MessageStore;
|
||||
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.AGGREGATION_DELTA;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.BROKER_NODE_TYPE;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_BROKER_PERMISSION;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.LABEL_AGGREGATION;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.LABEL_CLUSTER_NAME;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.LABEL_NODE_ID;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.LABEL_NODE_TYPE;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.OPEN_TELEMETRY_METER_NAME;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.SLS_OTEL_AK_ID_KEY;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.SLS_OTEL_AK_SECRET_KEY;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.SLS_OTEL_INSTANCE_ID_KEY;
|
||||
import static org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.SLS_OTEL_PROJECT_HEADER_KEY;
|
||||
|
||||
public class BrokerMetricsManager {
|
||||
private static final InternalLogger LOGGER = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
|
||||
private final BrokerConfig brokerConfig;
|
||||
private final MessageStore messageStore;
|
||||
private final BrokerController brokerController;
|
||||
public final static Map<String, String> LABEL_MAP = new HashMap<>();
|
||||
private OtlpGrpcMetricExporter metricExporter;
|
||||
private PeriodicMetricReader periodicMetricReader;
|
||||
private PrometheusHttpServer prometheusHttpServer;
|
||||
|
||||
public static ObservableLongGauge brokerPermission = null;
|
||||
|
||||
public BrokerMetricsManager(BrokerController brokerController) {
|
||||
this.brokerController = brokerController;
|
||||
brokerConfig = brokerController.getBrokerConfig();
|
||||
this.messageStore = brokerController.getMessageStore();
|
||||
init();
|
||||
}
|
||||
|
||||
public static AttributesBuilder newAttributesBuilder() {
|
||||
AttributesBuilder attributesBuilder = Attributes.builder();
|
||||
LABEL_MAP.forEach(attributesBuilder::put);
|
||||
return attributesBuilder;
|
||||
}
|
||||
|
||||
private boolean checkConfig() {
|
||||
if (brokerConfig == null) {
|
||||
return false;
|
||||
}
|
||||
BrokerConfig.MetricsExporterType exporterType = brokerConfig.getMetricsExporterType();
|
||||
if (!exporterType.isEnable()) {
|
||||
return false;
|
||||
}
|
||||
|
||||
switch (exporterType) {
|
||||
case OTLP_GRPC:
|
||||
return StringUtils.isNotBlank(brokerConfig.getMetricsGrpcCollectorEndpoint());
|
||||
case OTLP_GRPC_SLS:
|
||||
return StringUtils.isNotBlank(brokerConfig.getMetricsGrpcCollectorEndpoint()) &&
|
||||
StringUtils.isNotBlank(brokerConfig.getMetricsSlsProjectName()) &&
|
||||
StringUtils.isNotBlank(brokerConfig.getMetricsSlsInstanceName()) &&
|
||||
StringUtils.isNotBlank(brokerConfig.getMetricsSlsAccessKey()) &&
|
||||
StringUtils.isNotBlank(brokerConfig.getMetricsSlsSecretKey());
|
||||
case PROM:
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
private void init() {
|
||||
if (!checkConfig()) {
|
||||
LOGGER.error("check broker metrics config failed, will not export metrics");
|
||||
return;
|
||||
}
|
||||
|
||||
String labels = brokerConfig.getBrokerMetricsLabel();
|
||||
if (StringUtils.isNotBlank(labels)) {
|
||||
List<String> kvPairs = Splitter.on(',').omitEmptyStrings().splitToList(labels);
|
||||
for (String item : kvPairs) {
|
||||
String[] split = item.split(":");
|
||||
if (split.length != 2) {
|
||||
LOGGER.warn("brokerMetricsLabel is not valid: {}", brokerConfig.getBrokerMetricsLabel());
|
||||
continue;
|
||||
}
|
||||
LABEL_MAP.put(split[0], split[1]);
|
||||
}
|
||||
}
|
||||
if (brokerConfig.isBrokerMetricsPreferDelta()) {
|
||||
LABEL_MAP.put(LABEL_AGGREGATION, AGGREGATION_DELTA);
|
||||
}
|
||||
LABEL_MAP.put(LABEL_NODE_TYPE, BROKER_NODE_TYPE);
|
||||
LABEL_MAP.put(LABEL_CLUSTER_NAME, brokerConfig.getBrokerClusterName());
|
||||
LABEL_MAP.put(LABEL_NODE_ID, brokerConfig.getBrokerName());
|
||||
|
||||
SdkMeterProviderBuilder providerBuilder = SdkMeterProvider.builder()
|
||||
.setResource(Resource.empty());
|
||||
|
||||
if (brokerConfig.getMetricsExporterType() == BrokerConfig.MetricsExporterType.OTLP_GRPC ||
|
||||
brokerConfig.getMetricsExporterType() == BrokerConfig.MetricsExporterType.OTLP_GRPC_SLS) {
|
||||
String endpoint = brokerConfig.getMetricsGrpcCollectorEndpoint();
|
||||
if (!endpoint.startsWith("https://")) {
|
||||
endpoint = "https://" + endpoint;
|
||||
}
|
||||
OtlpGrpcMetricExporterBuilder metricExporterBuilder = OtlpGrpcMetricExporter.builder()
|
||||
.setEndpoint(endpoint)
|
||||
.setTimeout(brokerConfig.getMetricGrpcExporterTimeOutInMills(), TimeUnit.MILLISECONDS)
|
||||
.setAggregationTemporalitySelector(type -> {
|
||||
if (brokerConfig.isBrokerMetricsPreferDelta() &&
|
||||
(type == InstrumentType.COUNTER || type == InstrumentType.OBSERVABLE_COUNTER || type == InstrumentType.HISTOGRAM)) {
|
||||
return AggregationTemporality.DELTA;
|
||||
}
|
||||
return AggregationTemporality.CUMULATIVE;
|
||||
});
|
||||
|
||||
if (brokerConfig.getMetricsExporterType() == BrokerConfig.MetricsExporterType.OTLP_GRPC_SLS) {
|
||||
metricExporterBuilder.addHeader(SLS_OTEL_PROJECT_HEADER_KEY, brokerConfig.getMetricsSlsProjectName())
|
||||
.addHeader(SLS_OTEL_INSTANCE_ID_KEY, brokerConfig.getMetricsSlsInstanceName())
|
||||
.addHeader(SLS_OTEL_AK_ID_KEY, brokerConfig.getMetricsSlsAccessKey())
|
||||
.addHeader(SLS_OTEL_AK_SECRET_KEY, brokerConfig.getMetricsSlsSecretKey());
|
||||
}
|
||||
|
||||
metricExporter = metricExporterBuilder.build();
|
||||
|
||||
periodicMetricReader = PeriodicMetricReader.builder(metricExporter)
|
||||
.setInterval(brokerConfig.getMetricGrpcExporterIntervalInMills(), TimeUnit.MILLISECONDS)
|
||||
.build();
|
||||
|
||||
providerBuilder.registerMetricReader(periodicMetricReader);
|
||||
}
|
||||
|
||||
if (brokerConfig.getMetricsExporterType() == BrokerConfig.MetricsExporterType.PROM) {
|
||||
String promExporterHost = brokerConfig.getMetricsPromExporterHost();
|
||||
if (StringUtils.isBlank(promExporterHost)) {
|
||||
promExporterHost = brokerConfig.getBrokerIP1();
|
||||
}
|
||||
prometheusHttpServer = PrometheusHttpServer.builder()
|
||||
.setHost(promExporterHost)
|
||||
.setPort(brokerConfig.getMetricsPromExporterPort())
|
||||
.build();
|
||||
providerBuilder.registerMetricReader(prometheusHttpServer);
|
||||
}
|
||||
|
||||
Meter brokerMeter = OpenTelemetrySdk.builder()
|
||||
.setMeterProvider(providerBuilder.build())
|
||||
.build()
|
||||
.getMeter(OPEN_TELEMETRY_METER_NAME);
|
||||
|
||||
initStatsMetrics(brokerMeter);
|
||||
}
|
||||
|
||||
private void initStatsMetrics(Meter meter) {
|
||||
brokerPermission = meter.gaugeBuilder(GAUGE_BROKER_PERMISSION)
|
||||
.setDescription("Broker permission")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> measurement.record(brokerConfig.getBrokerPermission(), newAttributesBuilder().build()));
|
||||
}
|
||||
|
||||
public void shutdown() {
|
||||
if (brokerConfig.getMetricsExporterType() == BrokerConfig.MetricsExporterType.OTLP_GRPC ||
|
||||
brokerConfig.getMetricsExporterType() == BrokerConfig.MetricsExporterType.OTLP_GRPC_SLS) {
|
||||
periodicMetricReader.forceFlush();
|
||||
periodicMetricReader.shutdown();
|
||||
metricExporter.shutdown();
|
||||
}
|
||||
if (brokerConfig.getMetricsExporterType() == BrokerConfig.MetricsExporterType.PROM) {
|
||||
prometheusHttpServer.forceFlush();
|
||||
prometheusHttpServer.shutdown();
|
||||
}
|
||||
}
|
||||
}
|
||||
+3
-3
@@ -89,7 +89,7 @@ import static org.mockito.Mockito.doReturn;
|
||||
import static org.mockito.Mockito.spy;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
@RunWith(MockitoJUnitRunner.class)
|
||||
@RunWith(MockitoJUnitRunner.Silent.class)
|
||||
public class DefaultMQConsumerWithTraceTest {
|
||||
private String consumerGroup;
|
||||
private String consumerGroupNormal;
|
||||
@@ -105,7 +105,7 @@ public class DefaultMQConsumerWithTraceTest {
|
||||
private RebalancePushImpl rebalancePushImpl;
|
||||
private DefaultMQPushConsumer pushConsumer;
|
||||
private DefaultMQPushConsumer normalPushConsumer;
|
||||
private DefaultMQPushConsumer customTraceTopicpushConsumer;
|
||||
private DefaultMQPushConsumer customTraceTopicPushConsumer;
|
||||
|
||||
private AsyncTraceDispatcher asyncTraceDispatcher;
|
||||
private MQClientInstance mQClientTraceFactory;
|
||||
@@ -126,7 +126,7 @@ public class DefaultMQConsumerWithTraceTest {
|
||||
pushConsumer = new DefaultMQPushConsumer(consumerGroup, true, "");
|
||||
consumerGroupNormal = "FooBarGroup" + System.currentTimeMillis();
|
||||
normalPushConsumer = new DefaultMQPushConsumer(consumerGroupNormal, false, "");
|
||||
customTraceTopicpushConsumer = new DefaultMQPushConsumer(consumerGroup, true, customerTraceTopic);
|
||||
customTraceTopicPushConsumer = new DefaultMQPushConsumer(consumerGroup, true, customerTraceTopic);
|
||||
pushConsumer.setNamesrvAddr("127.0.0.1:9876");
|
||||
pushConsumer.setPullInterval(60 * 1000);
|
||||
|
||||
|
||||
@@ -319,6 +319,60 @@ public class BrokerConfig extends BrokerIdentity {
|
||||
|
||||
private long syncControllerMetadataPeriod = 10 * 1000;
|
||||
|
||||
public enum MetricsExporterType {
|
||||
DISABLE(0),
|
||||
OTLP_GRPC(1),
|
||||
OTLP_GRPC_SLS(2),
|
||||
PROM(3);
|
||||
|
||||
private final int value;
|
||||
|
||||
MetricsExporterType(int value) {
|
||||
this.value = value;
|
||||
}
|
||||
|
||||
public int getValue() {
|
||||
return value;
|
||||
}
|
||||
|
||||
public static MetricsExporterType valueOf(int value) {
|
||||
switch (value) {
|
||||
case 1:
|
||||
return OTLP_GRPC;
|
||||
case 2:
|
||||
return OTLP_GRPC_SLS;
|
||||
case 3:
|
||||
return PROM;
|
||||
default:
|
||||
return DISABLE;
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isEnable() {
|
||||
return this.value > 0;
|
||||
}
|
||||
}
|
||||
|
||||
private MetricsExporterType metricsExporterType = MetricsExporterType.DISABLE;
|
||||
|
||||
private String metricsGrpcCollectorEndpoint = "";
|
||||
|
||||
private String metricsSlsProjectName = "";
|
||||
private String metricsSlsInstanceName = "";
|
||||
private String metricsSlsAccessKey = "";
|
||||
private String metricsSlsSecretKey = "";
|
||||
|
||||
private long metricGrpcExporterTimeOutInMills = 3 * 1000;
|
||||
private long metricGrpcExporterIntervalInMills = 60 * 1000;
|
||||
|
||||
private int metricsPromExporterPort = 8080;
|
||||
private String metricsPromExporterHost = "";
|
||||
|
||||
// Label pairs in CSV. Each label follows pattern of Key:Value. eg: instance_id:xxx,uid:xxx
|
||||
private String brokerMetricsLabel = "";
|
||||
|
||||
private boolean brokerMetricsPreferDelta = true;
|
||||
|
||||
public long getMaxPopPollingSize() {
|
||||
return maxPopPollingSize;
|
||||
}
|
||||
@@ -1366,4 +1420,108 @@ public class BrokerConfig extends BrokerIdentity {
|
||||
public void setUseServerSideResetOffset(boolean useServerSideResetOffset) {
|
||||
this.useServerSideResetOffset = useServerSideResetOffset;
|
||||
}
|
||||
|
||||
public MetricsExporterType getMetricsExporterType() {
|
||||
return metricsExporterType;
|
||||
}
|
||||
|
||||
public void setMetricsExporterType(MetricsExporterType metricsExporterType) {
|
||||
this.metricsExporterType = metricsExporterType;
|
||||
}
|
||||
|
||||
public void setMetricsExporterType(int metricsExporterType) {
|
||||
this.metricsExporterType = MetricsExporterType.valueOf(metricsExporterType);
|
||||
}
|
||||
|
||||
public void setMetricsExporterType(String metricsExporterType) {
|
||||
this.metricsExporterType = MetricsExporterType.valueOf(metricsExporterType);
|
||||
}
|
||||
|
||||
public String getMetricsGrpcCollectorEndpoint() {
|
||||
return metricsGrpcCollectorEndpoint;
|
||||
}
|
||||
|
||||
public void setMetricsGrpcCollectorEndpoint(String metricsGrpcCollectorEndpoint) {
|
||||
this.metricsGrpcCollectorEndpoint = metricsGrpcCollectorEndpoint;
|
||||
}
|
||||
|
||||
public String getMetricsSlsProjectName() {
|
||||
return metricsSlsProjectName;
|
||||
}
|
||||
|
||||
public void setMetricsSlsProjectName(String metricsSlsProjectName) {
|
||||
this.metricsSlsProjectName = metricsSlsProjectName;
|
||||
}
|
||||
|
||||
public String getMetricsSlsInstanceName() {
|
||||
return metricsSlsInstanceName;
|
||||
}
|
||||
|
||||
public void setMetricsSlsInstanceName(String metricsSlsInstanceName) {
|
||||
this.metricsSlsInstanceName = metricsSlsInstanceName;
|
||||
}
|
||||
|
||||
public String getMetricsSlsAccessKey() {
|
||||
return metricsSlsAccessKey;
|
||||
}
|
||||
|
||||
public void setMetricsSlsAccessKey(String metricsSlsAccessKey) {
|
||||
this.metricsSlsAccessKey = metricsSlsAccessKey;
|
||||
}
|
||||
|
||||
public String getMetricsSlsSecretKey() {
|
||||
return metricsSlsSecretKey;
|
||||
}
|
||||
|
||||
public void setMetricsSlsSecretKey(String metricsSlsSecretKey) {
|
||||
this.metricsSlsSecretKey = metricsSlsSecretKey;
|
||||
}
|
||||
|
||||
public long getMetricGrpcExporterTimeOutInMills() {
|
||||
return metricGrpcExporterTimeOutInMills;
|
||||
}
|
||||
|
||||
public void setMetricGrpcExporterTimeOutInMills(long metricGrpcExporterTimeOutInMills) {
|
||||
this.metricGrpcExporterTimeOutInMills = metricGrpcExporterTimeOutInMills;
|
||||
}
|
||||
|
||||
public long getMetricGrpcExporterIntervalInMills() {
|
||||
return metricGrpcExporterIntervalInMills;
|
||||
}
|
||||
|
||||
public void setMetricGrpcExporterIntervalInMills(long metricGrpcExporterIntervalInMills) {
|
||||
this.metricGrpcExporterIntervalInMills = metricGrpcExporterIntervalInMills;
|
||||
}
|
||||
|
||||
public String getBrokerMetricsLabel() {
|
||||
return brokerMetricsLabel;
|
||||
}
|
||||
|
||||
public void setBrokerMetricsLabel(String brokerMetricsLabel) {
|
||||
this.brokerMetricsLabel = brokerMetricsLabel;
|
||||
}
|
||||
|
||||
public boolean isBrokerMetricsPreferDelta() {
|
||||
return brokerMetricsPreferDelta;
|
||||
}
|
||||
|
||||
public void setBrokerMetricsPreferDelta(boolean brokerMetricsPreferDelta) {
|
||||
this.brokerMetricsPreferDelta = brokerMetricsPreferDelta;
|
||||
}
|
||||
|
||||
public int getMetricsPromExporterPort() {
|
||||
return metricsPromExporterPort;
|
||||
}
|
||||
|
||||
public void setMetricsPromExporterPort(int metricsPromExporterPort) {
|
||||
this.metricsPromExporterPort = metricsPromExporterPort;
|
||||
}
|
||||
|
||||
public String getMetricsPromExporterHost() {
|
||||
return metricsPromExporterHost;
|
||||
}
|
||||
|
||||
public void setMetricsPromExporterHost(String metricsPromExporterHost) {
|
||||
this.metricsPromExporterHost = metricsPromExporterHost;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -127,7 +127,7 @@
|
||||
<extra-enforcer-rules.version>1.0-beta-4</extra-enforcer-rules.version>
|
||||
<concurrentlinkedhashmap-lru.version>1.4.2</concurrentlinkedhashmap-lru.version>
|
||||
<rocketmq-proto.version>2.0.1</rocketmq-proto.version>
|
||||
<grpc.version>1.45.0</grpc.version>
|
||||
<grpc.version>1.50.0</grpc.version>
|
||||
<protobuf.version>3.20.1</protobuf.version>
|
||||
<disruptor.version>1.2.10</disruptor.version>
|
||||
<org.relection.version>0.9.11</org.relection.version>
|
||||
@@ -755,6 +755,17 @@
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.jaegertracing</groupId>
|
||||
<artifactId>jaeger-thrift</artifactId>
|
||||
<version>${jaeger.version}</version>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>com.squareup.okhttp3</groupId>
|
||||
<artifactId>okhttp</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.jaegertracing</groupId>
|
||||
<artifactId>jaeger-client</artifactId>
|
||||
@@ -857,6 +868,10 @@
|
||||
<groupId>com.google.errorprone</groupId>
|
||||
<artifactId>error_prone_annotations</artifactId>
|
||||
</exclusion>
|
||||
<exclusion>
|
||||
<groupId>com.google.code.gson</groupId>
|
||||
<artifactId>gson</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
<dependency>
|
||||
@@ -868,6 +883,12 @@
|
||||
<groupId>com.github.ben-manes.caffeine</groupId>
|
||||
<artifactId>caffeine</artifactId>
|
||||
<version>${caffeine.version}</version>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>com.google.errorprone</groupId>
|
||||
<artifactId>error_prone_annotations</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
@@ -876,6 +897,39 @@
|
||||
<version>${spring.version}</version>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>com.squareup.okio</groupId>
|
||||
<artifactId>okio-jvm</artifactId>
|
||||
<version>3.0.0</version>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>org.jetbrains.kotlin</groupId>
|
||||
<artifactId>kotlin-stdlib</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.opentelemetry</groupId>
|
||||
<artifactId>opentelemetry-exporter-otlp</artifactId>
|
||||
<version>1.19.0</version>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>com.squareup.okio</groupId>
|
||||
<artifactId>okio-jvm</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.opentelemetry</groupId>
|
||||
<artifactId>opentelemetry-exporter-prometheus</artifactId>
|
||||
<version>1.19.0-alpha</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.opentelemetry</groupId>
|
||||
<artifactId>opentelemetry-sdk</artifactId>
|
||||
<version>1.19.0</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</dependencyManagement>
|
||||
|
||||
@@ -904,4 +958,4 @@
|
||||
<version>${awaitility.version}</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</project>
|
||||
</project>
|
||||
|
||||
@@ -74,6 +74,7 @@ java_library(
|
||||
"//remoting",
|
||||
"@maven//:ch_qos_logback_logback_core",
|
||||
"@maven//:com_alibaba_fastjson",
|
||||
"@maven//:com_github_ben_manes_caffeine_caffeine",
|
||||
"@maven//:com_google_guava_guava",
|
||||
"@maven//:com_google_protobuf_protobuf_java",
|
||||
"@maven//:com_google_protobuf_protobuf_java_util",
|
||||
@@ -84,6 +85,7 @@ java_library(
|
||||
"@maven//:io_netty_netty_all",
|
||||
"@maven//:org_apache_commons_commons_lang3",
|
||||
"@maven//:org_apache_rocketmq_rocketmq_proto",
|
||||
"@maven//:org_checkerframework_checker_qual",
|
||||
"@maven//:org_slf4j_slf4j_api",
|
||||
"@maven//:org_springframework_spring_core",
|
||||
],
|
||||
|
||||
+10
-2
@@ -25,6 +25,10 @@ java_library(
|
||||
"@maven//:com_alibaba_fastjson",
|
||||
"@maven//:io_netty_netty_all",
|
||||
"@maven//:org_apache_commons_commons_lang3",
|
||||
"@maven//:io_opentelemetry_opentelemetry_exporter_otlp",
|
||||
"@maven//:io_opentelemetry_opentelemetry_exporter_prometheus",
|
||||
"@maven//:io_opentelemetry_opentelemetry_sdk",
|
||||
"@maven//:com_squareup_okio_okio_jvm",
|
||||
],
|
||||
)
|
||||
|
||||
@@ -35,10 +39,14 @@ java_library(
|
||||
deps = [
|
||||
":remoting",
|
||||
"//:test_deps",
|
||||
"@maven//:io_netty_netty_all",
|
||||
"@maven//:com_google_code_gson_gson",
|
||||
"@maven//:io_netty_netty_all",
|
||||
"@maven//:com_google_code_gson_gson",
|
||||
"@maven//:com_alibaba_fastjson",
|
||||
"@maven//:org_apache_commons_commons_lang3",
|
||||
"@maven//:io_opentelemetry_opentelemetry_exporter_otlp",
|
||||
"@maven//:io_opentelemetry_opentelemetry_exporter_prometheus",
|
||||
"@maven//:io_opentelemetry_opentelemetry_sdk",
|
||||
"@maven//:com_squareup_okio_okio_jvm",
|
||||
],
|
||||
resources = glob(["src/test/resources/certs/*.pem"]) + glob(["src/test/resources/certs/*.key"])
|
||||
)
|
||||
|
||||
@@ -54,5 +54,25 @@
|
||||
<version>2.9.0</version>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.opentelemetry</groupId>
|
||||
<artifactId>opentelemetry-exporter-otlp</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.opentelemetry</groupId>
|
||||
<artifactId>opentelemetry-exporter-prometheus</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.opentelemetry</groupId>
|
||||
<artifactId>opentelemetry-sdk</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.grpc</groupId>
|
||||
<artifactId>grpc-stub</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.grpc</groupId>
|
||||
<artifactId>grpc-netty-shaded</artifactId>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</project>
|
||||
|
||||
Reference in New Issue
Block a user