From 6e27b543499358634a1ee363e2d3b6bce4422555 Mon Sep 17 00:00:00 2001 From: SSpirits Date: Wed, 19 Oct 2022 23:14:45 +0800 Subject: [PATCH] [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 --- WORKSPACE | 7 + broker/BUILD.bazel | 74 ++++--- .../rocketmq/broker/BrokerController.java | 4 + .../broker/metrics/BrokerMetricsConstant.java | 33 +++ .../broker/metrics/BrokerMetricsManager.java | 209 ++++++++++++++++++ .../trace/DefaultMQConsumerWithTraceTest.java | 6 +- .../apache/rocketmq/common/BrokerConfig.java | 158 +++++++++++++ pom.xml | 58 ++++- proxy/BUILD.bazel | 2 + remoting/BUILD.bazel | 12 +- remoting/pom.xml | 20 ++ 11 files changed, 542 insertions(+), 41 deletions(-) create mode 100644 broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java create mode 100644 broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java diff --git a/WORKSPACE b/WORKSPACE index e2971459b1..6fec91abe6 100644 --- a/WORKSPACE +++ b/WORKSPACE @@ -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 = [ diff --git a/broker/BUILD.bazel b/broker/BUILD.bazel index 1c7403a447..a537520f8f 100644 --- a/broker/BUILD.bazel +++ b/broker/BUILD.bazel @@ -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", ], ) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java index 9e4ee83eb2..717a08021b 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java @@ -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> 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(); diff --git a/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java new file mode 100644 index 0000000000..14ace640b9 --- /dev/null +++ b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java @@ -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"; +} diff --git a/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java new file mode 100644 index 0000000000..2e2ff94b27 --- /dev/null +++ b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java @@ -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 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 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(); + } + } +} diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java index ff82610246..e7abd570d8 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java @@ -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); diff --git a/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java b/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java index 854ef6334c..f39741e26b 100644 --- a/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java +++ b/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java @@ -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; + } } diff --git a/pom.xml b/pom.xml index 19311d60d0..b2176e6cef 100644 --- a/pom.xml +++ b/pom.xml @@ -127,7 +127,7 @@ 1.0-beta-4 1.4.2 2.0.1 - 1.45.0 + 1.50.0 3.20.1 1.2.10 0.9.11 @@ -755,6 +755,17 @@ + + io.jaegertracing + jaeger-thrift + ${jaeger.version} + + + com.squareup.okhttp3 + okhttp + + + io.jaegertracing jaeger-client @@ -857,6 +868,10 @@ com.google.errorprone error_prone_annotations + + com.google.code.gson + gson + @@ -868,6 +883,12 @@ com.github.ben-manes.caffeine caffeine ${caffeine.version} + + + com.google.errorprone + error_prone_annotations + + @@ -876,6 +897,39 @@ ${spring.version} test + + + com.squareup.okio + okio-jvm + 3.0.0 + + + org.jetbrains.kotlin + kotlin-stdlib + + + + + io.opentelemetry + opentelemetry-exporter-otlp + 1.19.0 + + + com.squareup.okio + okio-jvm + + + + + io.opentelemetry + opentelemetry-exporter-prometheus + 1.19.0-alpha + + + io.opentelemetry + opentelemetry-sdk + 1.19.0 + @@ -904,4 +958,4 @@ ${awaitility.version} - \ No newline at end of file + diff --git a/proxy/BUILD.bazel b/proxy/BUILD.bazel index e25273cc8d..fa67fc0186 100644 --- a/proxy/BUILD.bazel +++ b/proxy/BUILD.bazel @@ -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", ], diff --git a/remoting/BUILD.bazel b/remoting/BUILD.bazel index a7ab0df8f3..580b950809 100644 --- a/remoting/BUILD.bazel +++ b/remoting/BUILD.bazel @@ -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"]) ) diff --git a/remoting/pom.xml b/remoting/pom.xml index a61764319f..403e527e3c 100644 --- a/remoting/pom.xml +++ b/remoting/pom.xml @@ -54,5 +54,25 @@ 2.9.0 test + + io.opentelemetry + opentelemetry-exporter-otlp + + + io.opentelemetry + opentelemetry-exporter-prometheus + + + io.opentelemetry + opentelemetry-sdk + + + io.grpc + grpc-stub + + + io.grpc + grpc-netty-shaded +