From 2e846c60781ac90128f7d447a6eae48aa0b7bed7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=90=B4=E6=99=9F=20Wu=20Sheng?= Date: Mon, 3 Jun 2019 23:11:52 +0800 Subject: [PATCH] Dynamic Configuration and Dynamic Configuration Service (#2810) * Implement the base of config API module. * Finish configuration module codebases w/o tests, revert some changes. * Finish mock tests. * Add new gRPC service for configuration. * Make gRPC configuration sync provider works. * Add database slow statement dynamic setting. * Finish the dynamic configuration service. * Fix version. --- docs/en/guides/How-to-build.md | 1 + docs/en/setup/backend/backend-setup.md | 2 + docs/en/setup/backend/dynamic-config.md | 39 +++++ oap-server/pom.xml | 1 + .../configuration-api/pom.xml | 41 +++++ .../api/AbstractConfigurationProvider.java | 55 +++++++ .../api/ConfigChangeWatcher.java | 77 ++++++++++ .../configuration/api/ConfigTable.java} | 39 +++-- .../api/ConfigWatcherRegister.java | 145 ++++++++++++++++++ .../api/ConfigurationModule.java | 42 +++++ .../api/DynamicConfigurationService.java | 35 +++++ .../api/NoneConfigurationProvider.java | 60 ++++++++ ...ing.oap.server.library.module.ModuleDefine | 19 +++ ...g.oap.server.library.module.ModuleProvider | 20 +++ .../api/ConfigWatcherRegisterTest.java | 135 ++++++++++++++++ .../grpc-configuration-sync/pom.xml | 105 +++++++++++++ .../grpc/GRPCConfigWatcherRegister.java | 53 +++++++ .../grpc/GRPCConfigurationProvider.java | 57 +++++++ .../grpc/RemoteEndpointSettings.java | 43 ++++++ .../main/proto/configuration-service.proto | 50 ++++++ ...g.oap.server.library.module.ModuleProvider | 19 +++ oap-server/server-configuration/pom.xml | 37 +++++ oap-server/server-core/pom.xml | 5 + .../oap/server/core/CoreModuleProvider.java | 3 +- .../library/module/ModuleStartException.java | 3 + .../DBLatencyThresholdsAndWatcher.java | 78 ++++++++++ .../trace/provider/TraceModuleProvider.java | 10 +- .../provider/TraceServiceModuleConfig.java | 2 +- .../endpoint/MultiScopesSpanListener.java | 2 +- oap-server/server-starter/pom.xml | 7 + .../src/main/assembly/application.yml | 2 + .../src/main/resources/application.yml | 2 + pom.xml | 3 +- 33 files changed, 1165 insertions(+), 27 deletions(-) create mode 100644 docs/en/setup/backend/dynamic-config.md create mode 100644 oap-server/server-configuration/configuration-api/pom.xml create mode 100644 oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/AbstractConfigurationProvider.java create mode 100644 oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigChangeWatcher.java rename oap-server/{server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholds.java => server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigTable.java} (50%) create mode 100644 oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegister.java create mode 100644 oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigurationModule.java create mode 100644 oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/DynamicConfigurationService.java create mode 100644 oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/NoneConfigurationProvider.java create mode 100644 oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine create mode 100644 oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider create mode 100644 oap-server/server-configuration/configuration-api/src/test/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegisterTest.java create mode 100644 oap-server/server-configuration/grpc-configuration-sync/pom.xml create mode 100644 oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigWatcherRegister.java create mode 100644 oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigurationProvider.java create mode 100644 oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/RemoteEndpointSettings.java create mode 100644 oap-server/server-configuration/grpc-configuration-sync/src/main/proto/configuration-service.proto create mode 100644 oap-server/server-configuration/grpc-configuration-sync/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider create mode 100644 oap-server/server-configuration/pom.xml create mode 100644 oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholdsAndWatcher.java diff --git a/docs/en/guides/How-to-build.md b/docs/en/guides/How-to-build.md index f403b0520..893e7ad76 100644 --- a/docs/en/guides/How-to-build.md +++ b/docs/en/guides/How-to-build.md @@ -72,6 +72,7 @@ a oap image with name `bar/oap:foo`, run the following command * `grpc-java` and `java` folders in **oap-server/server-core/target/generated-sources/protobuf** * `grpc-java` and `java` folders in **oap-server/server-receiver-plugin/receiver-proto/target/generated-sources/protobuf** * `grpc-java` and `java` folders in **oap-server/exporter/target/generated-sources/protobuf** + * `grpc-java` and `java` folders in **oap-server/server-configuration/grpc-configuration-sync/target/generated-sources/protobuf** * `antlr4` folder in **oap-server/generate-tool-grammar/target/generated-sources** * `oal` folder in **oap-server/generated-analysis/target/generated-sources** diff --git a/docs/en/setup/backend/backend-setup.md b/docs/en/setup/backend/backend-setup.md index 0092b5c1a..6ec0798c3 100644 --- a/docs/en/setup/backend/backend-setup.md +++ b/docs/en/setup/backend/backend-setup.md @@ -81,6 +81,8 @@ scale and support high payload, you may need this. system. 1. [Time To Live (TTL)](ttl.md). Metrics and trace are time series data, they would be saved forever, you could set the expired time for each dimension. +1. [Dynamic Configuration](dynamic-config.md). Make configuration of OAP changed dynamic, from remote service +or 3rd party configuration management system. ## Telemetry for backend OAP backend cluster itself underlying is a distributed streaming process system. For helping the Ops team, diff --git a/docs/en/setup/backend/dynamic-config.md b/docs/en/setup/backend/dynamic-config.md new file mode 100644 index 000000000..de13691fc --- /dev/null +++ b/docs/en/setup/backend/dynamic-config.md @@ -0,0 +1,39 @@ +# Dynamic Configuration +SkyWalking Configurations mostly are set through `application.yml` and OS system environment variables. +But some of them are supporting dynamic settings from upstream management system. + +Right now, SkyWalking supports following dynamic configurations. + +| Config Key | Value Description | Value Format Example | +|:----:|:----:|:----:| +|receiver-trace.default.slowDBAccessThreshold| Thresholds of slow Database statement, override `receiver-trace/default/slowDBAccessThreshold` of `applciation.yml`. | default:200,mongodb:50| + + +This feature depends on upstream service, so it is **OFF** as default. +```yaml +configuration: + none: +``` + +## Dynamic Configuration Service, DCS +[Dynamic Configuration Service](../../../../oap-server/server-configuration/grpc-configuration-sync/src/main/proto/configuration-service.proto) +is a gRPC service, which requires the upstream system implemented. +The SkyWalking OAP fetches the configuration from the implementation(any system), after you open this implementation like this. + +```yaml +configuration: + grpc: + # Upstream system hostname + host: 127.0.0.1 + # Upstream system port + port: 9555 + #period : 60 # Unit seconds, sync period. Default fetch every 60 seconds. + #clusterName: "default" # the name of current cluster, set the name if you want to upstream system known. +``` + + +## 3rd party Configuration Center +We are welcome contributions to implement this module provider to support popular configuration center, +such as Zookeeper, etcd, Consul, Nacos. Submit issue to discuss. + + diff --git a/oap-server/pom.xml b/oap-server/pom.xml index 6c4948c75..f6e475234 100644 --- a/oap-server/pom.xml +++ b/oap-server/pom.xml @@ -42,6 +42,7 @@ server-telemetry generate-tool-grammar exporter + server-configuration diff --git a/oap-server/server-configuration/configuration-api/pom.xml b/oap-server/server-configuration/configuration-api/pom.xml new file mode 100644 index 000000000..093800079 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/pom.xml @@ -0,0 +1,41 @@ + + + + + + server-configuration + org.apache.skywalking + 6.2.0-SNAPSHOT + + 4.0.0 + + configuration-api + + + org.apache.skywalking + library-module + ${project.version} + compile + + + + + \ No newline at end of file diff --git a/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/AbstractConfigurationProvider.java b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/AbstractConfigurationProvider.java new file mode 100644 index 000000000..19459af69 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/AbstractConfigurationProvider.java @@ -0,0 +1,55 @@ +/* + * 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.configuration.api; + +import org.apache.skywalking.oap.server.library.module.*; + +/** + * The recommendation default base implementor of Configuration module. The real implementor could extend this provider + * to make a new one, easily. + * + * @author wusheng + */ +public abstract class AbstractConfigurationProvider extends ModuleProvider { + private ConfigWatcherRegister configWatcherRegister; + + @Override public Class module() { + return ConfigurationModule.class; + } + + @Override public void prepare() throws ServiceNotProvidedException, ModuleStartException { + configWatcherRegister = initConfigReader(); + this.registerServiceImplementation(DynamicConfigurationService.class, configWatcherRegister); + } + + protected abstract ConfigWatcherRegister initConfigReader() throws ModuleStartException; + + @Override public void start() throws ServiceNotProvidedException, ModuleStartException { + + } + + @Override public void notifyAfterCompleted() throws ServiceNotProvidedException, ModuleStartException { + configWatcherRegister.start(); + } + + @Override public String[] requiredModules() { + return new String[0]; + } + +} diff --git a/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigChangeWatcher.java b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigChangeWatcher.java new file mode 100644 index 000000000..1623332f4 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigChangeWatcher.java @@ -0,0 +1,77 @@ +/* + * 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.configuration.api; + +import lombok.*; +import org.apache.skywalking.oap.server.library.module.*; + +/** + * ConfigChangeWatcher represents a watcher implementor, it will be called when the target value changed. + * + * @author wusheng + */ +@Getter +public abstract class ConfigChangeWatcher { + private final String module; + private final ModuleProvider provider; + private final String itemName; + + public ConfigChangeWatcher(String module, ModuleProvider provider, String itemName) { + this.module = module; + this.provider = provider; + this.itemName = itemName; + } + + /** + * Notify the watcher, the new value received. + * + * @param value of new. + */ + public abstract void notify(ConfigChangeEvent value); + + /** + * @return current value of current config. + */ + public abstract String value(); + + @Override public String toString() { + return "ConfigChangeWatcher{" + + "module=" + module + + ", provider=" + provider + + ", itemName='" + itemName + '\'' + + '}'; + } + + @Setter(AccessLevel.PACKAGE) + @Getter + public static class ConfigChangeEvent { + private String newValue; + private EventType eventType; + + public ConfigChangeEvent(String newValue, + EventType eventType) { + this.newValue = newValue; + this.eventType = eventType; + } + } + + public enum EventType { + ADD, MODIFY, DELETE + } +} diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholds.java b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigTable.java similarity index 50% rename from oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholds.java rename to oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigTable.java index a0fbbecc8..742f44f4a 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholds.java +++ b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigTable.java @@ -16,36 +16,33 @@ * */ -package org.apache.skywalking.oap.server.receiver.trace.provider; +package org.apache.skywalking.oap.server.configuration.api; import java.util.*; +import lombok.*; /** + * ConfigTable contains all config. + * * @author wusheng */ -public class DBLatencyThresholds { - private Map thresholds; +public class ConfigTable { + @Getter + private List items = new ArrayList<>(); - DBLatencyThresholds(String config) { - thresholds = new HashMap<>(); - String[] settings = config.split(","); - for (String setting : settings) { - String[] typeValue = setting.split(":"); - if (typeValue.length == 2) { - thresholds.put(typeValue[0].trim().toLowerCase(), Integer.parseInt(typeValue[1].trim())); - } - } - if (!thresholds.containsKey("default")) { - thresholds.put("default", 10000); - } + public void add(ConfigItem item) { + items.add(item); } - public int getThreshold(String type) { - type = type.toLowerCase(); - if (thresholds.containsKey(type)) { - return thresholds.get(type); - } else { - return thresholds.get("default"); + @Getter + @Setter + public static class ConfigItem { + private String name; + private String value; + + public ConfigItem(String name, String value) { + this.name = name; + this.value = value; } } } diff --git a/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegister.java b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegister.java new file mode 100644 index 000000000..7bc1e092e --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegister.java @@ -0,0 +1,145 @@ +/* + * 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.configuration.api; + +import java.util.*; +import java.util.concurrent.*; +import lombok.Getter; +import org.apache.skywalking.apm.util.RunnableWithExceptionProtection; +import org.slf4j.*; + +/** + * The default implementor of Config Watcher register. + * + * @author wusheng + */ +public abstract class ConfigWatcherRegister implements DynamicConfigurationService { + private static final Logger logger = LoggerFactory.getLogger(ConfigWatcherRegister.class); + public static final String LINE_SEPARATOR = System.getProperty("line.separator", "\n"); + + private Register register = new Register(); + private volatile boolean isStarted = false; + private final long syncPeriod; + + public ConfigWatcherRegister() { + this(60); + } + + public ConfigWatcherRegister(long syncPeriod) { + this.syncPeriod = syncPeriod; + } + + @Override synchronized public void registerConfigChangeWatcher(ConfigChangeWatcher watcher) { + if (isStarted) { + throw new IllegalStateException("Config Register has been started. Can't register new watcher."); + } + + WatcherHolder holder = new WatcherHolder(watcher); + if (register.containsKey(holder.getKey())) { + throw new IllegalStateException("Duplicate register, watcher=" + watcher); + } + register.put(holder.getKey(), holder); + } + + public void start() { + isStarted = true; + + configSync(); + logger.info("Current configurations after the bootstrap sync." + LINE_SEPARATOR + register.toString()); + + Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate( + new RunnableWithExceptionProtection(() -> configSync(), + t -> logger.error("Sync config center error.", t)), syncPeriod, syncPeriod, TimeUnit.SECONDS); + } + + void configSync() { + ConfigTable configTable = readConfig(); + + configTable.getItems().forEach(item -> { + String itemName = item.getName(); + WatcherHolder holder = register.get(itemName); + if (holder != null) { + ConfigChangeWatcher watcher = holder.getWatcher(); + String newItemValue = item.getValue(); + if (newItemValue == null) { + if (watcher.value() != null) { + // Notify watcher, the new value is null with delete event type. + watcher.notify(new ConfigChangeWatcher.ConfigChangeEvent(null, ConfigChangeWatcher.EventType.DELETE)); + } else { + // Don't need to notify, stay in null. + } + } else { + if (!newItemValue.equals(watcher.value())) { + watcher.notify(new ConfigChangeWatcher.ConfigChangeEvent(newItemValue, ConfigChangeWatcher.EventType.MODIFY)); + } else { + // Don't need to notify, stay in the same config value. + } + } + } else { + logger.warn("Config {} from configuration center, doesn't match any watcher, ignore.", itemName); + } + }); + + logger.trace("Current configurations after the sync." + LINE_SEPARATOR + register.toString()); + } + + public abstract ConfigTable readConfig(); + + public class Register { + private Map register = new HashMap<>(); + + private boolean containsKey(String key) { + return register.containsKey(key); + } + + private void put(String key, WatcherHolder holder) { + register.put(key, holder); + } + + public WatcherHolder get(String name) { + return register.get(name); + } + + @Override public String toString() { + StringBuilder registerTableDescription = new StringBuilder(); + registerTableDescription.append("Following dynamic config items are available.").append(LINE_SEPARATOR); + registerTableDescription.append("---------------------------------------------").append(LINE_SEPARATOR); + register.forEach((key, holder) -> { + ConfigChangeWatcher watcher = holder.getWatcher(); + registerTableDescription.append("key:").append(key) + .append(" module:").append(watcher.getModule()) + .append(" provider:").append(watcher.getProvider().name()) + .append(" value(current):").append(watcher.value()) + .append(LINE_SEPARATOR); + }); + return registerTableDescription.toString(); + } + } + + @Getter + private class WatcherHolder { + private ConfigChangeWatcher watcher; + private final String key; + + public WatcherHolder(ConfigChangeWatcher watcher) { + this.watcher = watcher; + this.key = String.join(".", watcher.getModule(), watcher.getProvider().name(), watcher.getItemName()); + } + } +} diff --git a/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigurationModule.java b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigurationModule.java new file mode 100644 index 000000000..f0f55a426 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigurationModule.java @@ -0,0 +1,42 @@ +/* + * 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.configuration.api; + +import org.apache.skywalking.oap.server.library.module.ModuleDefine; + +/** + * Configuration Module sync the settings from remote service, the remote service could be implemented by this module + * provider. + * + * Any configuration item in the whole OAP backend could register a watcher to configuration module, the item change + * watcher will be called, if the value changed. + * + * @author wusheng + */ +public class ConfigurationModule extends ModuleDefine { + public static final String NAME = "configuration"; + + public ConfigurationModule() { + super(NAME); + } + + @Override public Class[] services() { + return new Class[] {DynamicConfigurationService.class}; + } +} diff --git a/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/DynamicConfigurationService.java b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/DynamicConfigurationService.java new file mode 100644 index 000000000..a70c27048 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/DynamicConfigurationService.java @@ -0,0 +1,35 @@ +/* + * 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.configuration.api; + +import org.apache.skywalking.oap.server.library.module.*; + +/** + * DynamicConfigurationService provides API to register config change watcher. + * + * @author wusheng + */ +public interface DynamicConfigurationService extends Service { + /** + * Register a watcher to the target value + * + * @param watcher to register + */ + void registerConfigChangeWatcher(ConfigChangeWatcher watcher); +} diff --git a/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/NoneConfigurationProvider.java b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/NoneConfigurationProvider.java new file mode 100644 index 000000000..633edcba1 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/NoneConfigurationProvider.java @@ -0,0 +1,60 @@ +/* + * 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.configuration.api; + +import org.apache.skywalking.oap.server.library.module.*; + +/** + * A nutshell configuration implementor. + * + * @author wusheng + */ +public class NoneConfigurationProvider extends ModuleProvider { + @Override public String name() { + return "none"; + } + + @Override public Class module() { + return ConfigurationModule.class; + } + + @Override public ModuleConfig createConfigBeanIfAbsent() { + return null; + } + + @Override public void prepare() throws ServiceNotProvidedException, ModuleStartException { + this.registerServiceImplementation(DynamicConfigurationService.class, new DynamicConfigurationService() { + @Override public void registerConfigChangeWatcher(ConfigChangeWatcher watcher) { + + } + }); + } + + @Override public void start() throws ServiceNotProvidedException, ModuleStartException { + + } + + @Override public void notifyAfterCompleted() throws ServiceNotProvidedException, ModuleStartException { + + } + + @Override public String[] requiredModules() { + return new String[0]; + } +} diff --git a/oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine b/oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine new file mode 100644 index 000000000..0353b08d1 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine @@ -0,0 +1,19 @@ +# +# 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. +# +# + +org.apache.skywalking.oap.server.configuration.api.ConfigurationModule \ No newline at end of file diff --git a/oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider b/oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider new file mode 100644 index 000000000..48e42ddd6 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider @@ -0,0 +1,20 @@ +# +# 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. +# +# +# + +org.apache.skywalking.oap.server.configuration.api.NoneConfigurationProvider \ No newline at end of file diff --git a/oap-server/server-configuration/configuration-api/src/test/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegisterTest.java b/oap-server/server-configuration/configuration-api/src/test/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegisterTest.java new file mode 100644 index 000000000..159837e34 --- /dev/null +++ b/oap-server/server-configuration/configuration-api/src/test/java/org/apache/skywalking/oap/server/configuration/api/ConfigWatcherRegisterTest.java @@ -0,0 +1,135 @@ +/* + * 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.configuration.api; + +import org.apache.skywalking.oap.server.library.module.*; +import org.junit.*; +import org.powermock.reflect.Whitebox; + +/** + * @author wusheng + */ +public class ConfigWatcherRegisterTest { + private ConfigWatcherRegister register; + + @Before + public void setup() { + register = new MockConfigWatcherRegister(); + } + + @After + public void tearDown() { + register = null; + } + + @Test + public void testInit() { + final String[] newValue = new String[1]; + + register.registerConfigChangeWatcher(new ConfigChangeWatcher("MockModule", new MockProvider(), "prop2") { + @Override public void notify(ConfigChangeEvent value) { + newValue[0] = value.getNewValue(); + } + + @Override public String value() { + return null; + } + }); + + register.configSync(); + + Assert.assertEquals("abc2", newValue[0]); + } + + @Test + public void testRegisterTableLog() { + register.registerConfigChangeWatcher(new ConfigChangeWatcher("MockModule", new MockProvider(), "prop2") { + @Override public void notify(ConfigChangeEvent value) { + } + + @Override public String value() { + return null; + } + }); + + register.configSync(); + ConfigWatcherRegister.Register registerTable = Whitebox.getInternalState(this.register, "register"); + + String expected = "Following dynamic config items are available." + ConfigWatcherRegister.LINE_SEPARATOR + + "---------------------------------------------" + ConfigWatcherRegister.LINE_SEPARATOR + + "key:MockModule.provider.prop2 module:MockModule provider:provider value(current):null" + ConfigWatcherRegister.LINE_SEPARATOR; + + Assert.assertEquals(expected, registerTable.toString()); + } + + public static class MockConfigWatcherRegister extends ConfigWatcherRegister { + + @Override public ConfigTable readConfig() { + ConfigTable.ConfigItem item1 = new ConfigTable.ConfigItem("module.provider.prop1", "abc"); + ConfigTable.ConfigItem item2 = new ConfigTable.ConfigItem("MockModule.provider.prop2", "abc2"); + + ConfigTable table = new ConfigTable(); + table.add(item1); + table.add(item2); + return table; + } + } + + public static class MockModule extends ModuleDefine { + + public MockModule() { + super("MockModule"); + } + + @Override public Class[] services() { + return new Class[0]; + } + } + + public static class MockProvider extends ModuleProvider { + + @Override public String name() { + return "provider"; + } + + @Override public Class module() { + return MockModule.class; + } + + @Override public ModuleConfig createConfigBeanIfAbsent() { + return null; + } + + @Override public void prepare() throws ServiceNotProvidedException, ModuleStartException { + + } + + @Override public void start() throws ServiceNotProvidedException, ModuleStartException { + + } + + @Override public void notifyAfterCompleted() throws ServiceNotProvidedException, ModuleStartException { + + } + + @Override public String[] requiredModules() { + return new String[0]; + } + } +} diff --git a/oap-server/server-configuration/grpc-configuration-sync/pom.xml b/oap-server/server-configuration/grpc-configuration-sync/pom.xml new file mode 100644 index 000000000..10d784014 --- /dev/null +++ b/oap-server/server-configuration/grpc-configuration-sync/pom.xml @@ -0,0 +1,105 @@ + + + + + + server-configuration + org.apache.skywalking + 6.2.0-SNAPSHOT + + 4.0.0 + + grpc-configuration-sync + + + + org.apache.skywalking + configuration-api + ${project.version} + + + org.apache.skywalking + library-client + ${project.version} + + + io.grpc + grpc-core + + + io.grpc + grpc-netty + + + io.grpc + grpc-protobuf + + + io.grpc + grpc-stub + + + + + + + kr.motd.maven + os-maven-plugin + 1.4.1.Final + + + + + org.apache.maven.plugins + maven-resources-plugin + 2.4.3 + + ${project.build.sourceEncoding} + + + + org.xolstice.maven.plugins + protobuf-maven-plugin + 0.5.0 + + + com.google.protobuf:protoc:3.3.0:exe:${os.detected.classifier} + + grpc-java + io.grpc:protoc-gen-grpc-java:1.4.0:exe:${os.detected.classifier} + + + + + + compile + compile-custom + + + + + + + \ No newline at end of file diff --git a/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigWatcherRegister.java b/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigWatcherRegister.java new file mode 100644 index 000000000..e883ec35c --- /dev/null +++ b/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigWatcherRegister.java @@ -0,0 +1,53 @@ +/* + * 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.configuration.grpc; + +import io.grpc.netty.NettyChannelBuilder; +import org.apache.skywalking.oap.server.configuration.api.*; +import org.apache.skywalking.oap.server.configuration.service.*; +import org.slf4j.*; + +/** + * @author wusheng + */ +public class GRPCConfigWatcherRegister extends ConfigWatcherRegister { + private static final Logger logger = LoggerFactory.getLogger(GRPCConfigWatcherRegister.class); + + private RemoteEndpointSettings settings; + private ConfigurationServiceGrpc.ConfigurationServiceBlockingStub stub; + + public GRPCConfigWatcherRegister(RemoteEndpointSettings settings) { + super(settings.getPeriod()); + this.settings = settings; + stub = ConfigurationServiceGrpc.newBlockingStub(NettyChannelBuilder.forAddress(settings.getHost(), settings.getPort()).usePlaintext().build()); + } + + @Override public ConfigTable readConfig() { + ConfigTable table = new ConfigTable(); + try { + ConfigurationResponse response = stub.call(ConfigurationRequest.newBuilder().setClusterName(settings.getClusterName()).build()); + response.getConfigTableList().forEach(config -> { + table.add(new ConfigTable.ConfigItem(config.getName(), config.getValue())); + }); + } catch (Exception e) { + logger.error("Remote config center [" + settings + "] is not available.", e); + } + return table; + } +} diff --git a/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigurationProvider.java b/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigurationProvider.java new file mode 100644 index 000000000..0dc4ccf70 --- /dev/null +++ b/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/GRPCConfigurationProvider.java @@ -0,0 +1,57 @@ +/* + * 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.configuration.grpc; + +import com.google.common.base.Strings; +import org.apache.skywalking.oap.server.configuration.api.*; +import org.apache.skywalking.oap.server.library.module.*; + +/** + * Get configuration from remote through gRPC protocol. + * + * Read configuration-service.proto for more details. + * + * @author wusheng + */ +public class GRPCConfigurationProvider extends AbstractConfigurationProvider { + private RemoteEndpointSettings settings; + + public GRPCConfigurationProvider() { + settings = new RemoteEndpointSettings(); + } + + @Override public String name() { + return "grpc"; + } + + @Override public ModuleConfig createConfigBeanIfAbsent() { + return settings; + } + + @Override protected ConfigWatcherRegister initConfigReader() throws ModuleStartException { + if (Strings.isNullOrEmpty(settings.getHost())) { + throw new ModuleStartException("No host setting."); + } + if (settings.getPort() < 1) { + throw new ModuleStartException("No port setting."); + } + + return new GRPCConfigWatcherRegister(settings); + } +} diff --git a/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/RemoteEndpointSettings.java b/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/RemoteEndpointSettings.java new file mode 100644 index 000000000..8d041fb24 --- /dev/null +++ b/oap-server/server-configuration/grpc-configuration-sync/src/main/java/org/apache/skywalking/oap/server/configuration/grpc/RemoteEndpointSettings.java @@ -0,0 +1,43 @@ +/* + * 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.configuration.grpc; + +import lombok.*; +import org.apache.skywalking.oap.server.library.module.ModuleConfig; + +/** + * @author wusheng + */ +@Setter +@Getter +public class RemoteEndpointSettings extends ModuleConfig { + private String host; + private int port; + private String clusterName = "default"; + // Sync configuration per 60 seconds. + private int period = 60; + + @Override public String toString() { + return "RemoteEndpointSettings{" + + "host='" + host + '\'' + + ", port=" + port + + ", clusterName='" + clusterName + '\'' + + '}'; + } +} diff --git a/oap-server/server-configuration/grpc-configuration-sync/src/main/proto/configuration-service.proto b/oap-server/server-configuration/grpc-configuration-sync/src/main/proto/configuration-service.proto new file mode 100644 index 000000000..be15aef1d --- /dev/null +++ b/oap-server/server-configuration/grpc-configuration-sync/src/main/proto/configuration-service.proto @@ -0,0 +1,50 @@ +/* + * 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. + * + */ + +syntax = "proto3"; + +option java_multiple_files = true; +option java_package = "org.apache.skywalking.oap.server.configuration.service"; + +service ConfigurationService { + rpc call (ConfigurationRequest) returns (ConfigurationResponse) { + } +} + +message ConfigurationRequest { + // Logic name of this cluster, + // in case the remote configuration center implementation support + // configuration management for multiple clusters. + string clusterName = 1; +} + +message ConfigurationResponse { + // Include all config items. + // All config name should be not empty, + // the name is composed by "module name"."provider name"."item name". + // Each watcher implementor provides this, and it will be notified when the value changed. + // + // If the config center wants to set the value to NULL or empty, + // must set the name with empty value explicitly. + repeated Config configTable = 1; +} + +message Config { + string name = 1; + string value = 2; +} \ No newline at end of file diff --git a/oap-server/server-configuration/grpc-configuration-sync/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider b/oap-server/server-configuration/grpc-configuration-sync/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider new file mode 100644 index 000000000..1a9644e83 --- /dev/null +++ b/oap-server/server-configuration/grpc-configuration-sync/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider @@ -0,0 +1,19 @@ +# +# 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. +# +# + +org.apache.skywalking.oap.server.configuration.grpc.GRPCConfigurationProvider \ No newline at end of file diff --git a/oap-server/server-configuration/pom.xml b/oap-server/server-configuration/pom.xml new file mode 100644 index 000000000..b161e5185 --- /dev/null +++ b/oap-server/server-configuration/pom.xml @@ -0,0 +1,37 @@ + + + + + + oap-server + org.apache.skywalking + 6.2.0-SNAPSHOT + + 4.0.0 + + server-configuration + pom + + configuration-api + grpc-configuration-sync + + + \ No newline at end of file diff --git a/oap-server/server-core/pom.xml b/oap-server/server-core/pom.xml index a66eba2d8..0eb0a8843 100644 --- a/oap-server/server-core/pom.xml +++ b/oap-server/server-core/pom.xml @@ -44,6 +44,11 @@ telemetry-api ${project.version} + + org.apache.skywalking + configuration-api + ${project.version} + org.apache.skywalking library-util diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java index 1545a2b78..9430020c5 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.core; import java.io.IOException; +import org.apache.skywalking.oap.server.configuration.api.ConfigurationModule; import org.apache.skywalking.oap.server.core.analysis.*; import org.apache.skywalking.oap.server.core.annotation.AnnotationScan; import org.apache.skywalking.oap.server.core.cache.*; @@ -188,6 +189,6 @@ public class CoreModuleProvider extends ModuleProvider { @Override public String[] requiredModules() { - return new String[] {TelemetryModule.NAME}; + return new String[] {TelemetryModule.NAME, ConfigurationModule.NAME}; } } diff --git a/oap-server/server-library/library-module/src/main/java/org/apache/skywalking/oap/server/library/module/ModuleStartException.java b/oap-server/server-library/library-module/src/main/java/org/apache/skywalking/oap/server/library/module/ModuleStartException.java index 010cb65da..01be5957c 100644 --- a/oap-server/server-library/library-module/src/main/java/org/apache/skywalking/oap/server/library/module/ModuleStartException.java +++ b/oap-server/server-library/library-module/src/main/java/org/apache/skywalking/oap/server/library/module/ModuleStartException.java @@ -22,6 +22,9 @@ package org.apache.skywalking.oap.server.library.module; * @author peng-yongsheng */ public class ModuleStartException extends Exception { + public ModuleStartException(String message) { + super(message); + } public ModuleStartException(String message, Throwable cause) { super(message, cause); diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholdsAndWatcher.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholdsAndWatcher.java new file mode 100644 index 000000000..a2b2dc3f9 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/DBLatencyThresholdsAndWatcher.java @@ -0,0 +1,78 @@ +/* + * 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.receiver.trace.provider; + +import java.util.*; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.skywalking.oap.server.configuration.api.ConfigChangeWatcher; +import org.apache.skywalking.oap.server.receiver.trace.module.TraceModule; + +/** + * @author wusheng + */ +public class DBLatencyThresholdsAndWatcher extends ConfigChangeWatcher { + private AtomicReference> thresholds; + private AtomicReference settingsString; + + DBLatencyThresholdsAndWatcher(String config, TraceModuleProvider provider) { + super(TraceModule.NAME, provider, "slowDBAccessThreshold"); + thresholds = new AtomicReference(new HashMap<>()); + settingsString = new AtomicReference<>(""); + + activeSetting(config); + } + + private void activeSetting(String config) { + Map newThresholds = new HashMap<>(); + String[] settings = config.split(","); + for (String setting : settings) { + String[] typeValue = setting.split(":"); + if (typeValue.length == 2) { + newThresholds.put(typeValue[0].trim().toLowerCase(), Integer.parseInt(typeValue[1].trim())); + } + } + if (!newThresholds.containsKey("default")) { + newThresholds.put("default", 10000); + } + + thresholds.set(newThresholds); + settingsString.set(config); + } + + public int getThreshold(String type) { + type = type.toLowerCase(); + if (thresholds.get().containsKey(type)) { + return thresholds.get().get(type); + } else { + return thresholds.get().get("default"); + } + } + + @Override public void notify(ConfigChangeEvent value) { + if (EventType.DELETE.equals(value.getEventType())) { + activeSetting(""); + } else { + activeSetting(value.getNewValue()); + } + } + + @Override public String value() { + return settingsString.get(); + } +} diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java index a47ec32ea..c54840836 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.receiver.trace.provider; import java.io.IOException; +import org.apache.skywalking.oap.server.configuration.api.*; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.server.*; import org.apache.skywalking.oap.server.library.module.*; @@ -42,6 +43,7 @@ public class TraceModuleProvider extends ModuleProvider { private final TraceServiceModuleConfig moduleConfig; private SegmentParse.Producer segmentProducer; private SegmentParseV2.Producer segmentProducerV2; + private DBLatencyThresholdsAndWatcher thresholds; public TraceModuleProvider() { this.moduleConfig = new TraceServiceModuleConfig(); @@ -60,7 +62,9 @@ public class TraceModuleProvider extends ModuleProvider { } @Override public void prepare() throws ServiceNotProvidedException { - moduleConfig.setDbLatencyThresholds(new DBLatencyThresholds(moduleConfig.getSlowDBAccessThreshold())); + thresholds = new DBLatencyThresholdsAndWatcher(moduleConfig.getSlowDBAccessThreshold(), this); + + moduleConfig.setDbLatencyThresholdsAndWatcher(thresholds); SegmentParserListenerManager listenerManager = new SegmentParserListenerManager(); if (moduleConfig.isTraceAnalysis()) { @@ -84,9 +88,11 @@ public class TraceModuleProvider extends ModuleProvider { } @Override public void start() throws ModuleStartException { + DynamicConfigurationService dynamicConfigurationService = getManager().find(ConfigurationModule.NAME).provider().getService(DynamicConfigurationService.class); GRPCHandlerRegister grpcHandlerRegister = getManager().find(SharingServerModule.NAME).provider().getService(GRPCHandlerRegister.class); JettyHandlerRegister jettyHandlerRegister = getManager().find(SharingServerModule.NAME).provider().getService(JettyHandlerRegister.class); try { + dynamicConfigurationService.registerConfigChangeWatcher(thresholds); grpcHandlerRegister.addHandler(new TraceSegmentServiceHandler(segmentProducer)); grpcHandlerRegister.addHandler(new TraceSegmentReportServiceHandler(segmentProducerV2, getManager())); @@ -111,6 +117,6 @@ public class TraceModuleProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[] {TelemetryModule.NAME, CoreModule.NAME, SharingServerModule.NAME}; + return new String[] {TelemetryModule.NAME, CoreModule.NAME, SharingServerModule.NAME, ConfigurationModule.NAME}; } } diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceServiceModuleConfig.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceServiceModuleConfig.java index 13ceabeff..0c322b2f1 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceServiceModuleConfig.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceServiceModuleConfig.java @@ -38,7 +38,7 @@ public class TraceServiceModuleConfig extends ModuleConfig { * The threshold used to check the slow database access. Unit, millisecond. */ @Setter @Getter private String slowDBAccessThreshold = "default:200"; - @Setter @Getter private DBLatencyThresholds dbLatencyThresholds; + @Setter @Getter private DBLatencyThresholdsAndWatcher dbLatencyThresholdsAndWatcher; /** * Analysis trace status. * diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java index b1c92c8cd..20f3b9b7e 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java @@ -165,7 +165,7 @@ public class MultiScopesSpanListener implements EntrySpanListener, ExitSpanListe } else if (SpanTags.DB_TYPE.equals(tag.getKey())) { String dbType = tag.getValue(); - DBLatencyThresholds thresholds = config.getDbLatencyThresholds(); + DBLatencyThresholdsAndWatcher thresholds = config.getDbLatencyThresholdsAndWatcher(); int threshold = thresholds.getThreshold(dbType); if (sourceBuilder.getLatency() > threshold) { isSlowDBAccess = true; diff --git a/oap-server/server-starter/pom.xml b/oap-server/server-starter/pom.xml index 44653f9ef..c80be5b5f 100644 --- a/oap-server/server-starter/pom.xml +++ b/oap-server/server-starter/pom.xml @@ -171,6 +171,13 @@ exporter ${project.version} + + + + org.apache.skywalking + grpc-configuration-sync + ${project.version} + skywalking-oap diff --git a/oap-server/server-starter/src/main/assembly/application.yml b/oap-server/server-starter/src/main/assembly/application.yml index 97d12b113..e0c0c16cc 100644 --- a/oap-server/server-starter/src/main/assembly/application.yml +++ b/oap-server/server-starter/src/main/assembly/application.yml @@ -123,6 +123,8 @@ alarm: default: telemetry: none: +configuration: + none: #exporter: # grpc: # targetHost: ${SW_EXPORTER_GRPC_HOST:127.0.0.1} diff --git a/oap-server/server-starter/src/main/resources/application.yml b/oap-server/server-starter/src/main/resources/application.yml index bac8a8496..e1940f6cb 100644 --- a/oap-server/server-starter/src/main/resources/application.yml +++ b/oap-server/server-starter/src/main/resources/application.yml @@ -129,6 +129,8 @@ telemetry: prometheus: host: ${SW_TELEMETRY_PROMETHEUS_HOST:0.0.0.0} port: ${SW_TELEMETRY_PROMETHEUS_PORT:1234} +configuration: + none: #exporter: # grpc: # targetHost: ${SW_EXPORTER_GRPC_HOST:127.0.0.1} diff --git a/pom.xml b/pom.xml index 4b1ce2b7d..a4dce435c 100644 --- a/pom.xml +++ b/pom.xml @@ -485,7 +485,8 @@ org/apache/skywalking/apm/network/language/**/*.java, org/apache/skywalking/oap/server/core/remote/grpc/proto/*.java, org/apache/skywalking/oal/tool/grammar/*.java, - org/apache/skywalking/oap/server/exporter/grpc/*.java + org/apache/skywalking/oap/server/exporter/grpc/*.java, + org/apache/skywalking/oap/server/configuration/service/*.java