diff --git a/CHANGES.md b/CHANGES.md index 3a2ae95dc..6a81392e1 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -37,6 +37,7 @@ Release Notes. * Upgrade snake yaml caused by CVE-2017-18640. * Upgrade embed tomcat caused by CVE-2020-13935. * Upgrade commons-lang3 to avoid potential NPE in some JDK versions. +* OAL supports generating metrics from events. #### UI diff --git a/docs/en/concepts-and-designs/event.md b/docs/en/concepts-and-designs/event.md index 133ccddca..b73c8641f 100644 --- a/docs/en/concepts-and-designs/event.md +++ b/docs/en/concepts-and-designs/event.md @@ -57,7 +57,7 @@ There are also cases where you would already have both the start time and end ti ## How to Configure Alarms for Events -Events are derived from metrics, and can be the source to trigger alarms. For example, if a specific event occurs for a +Events derive from metrics, and can be the source to trigger alarms. For example, if a specific event occurs for a certain times in a period, alarms can be triggered and sent. Every event has a default `value = 1`, when `n` events with the same name are reported, they are aggregated @@ -101,6 +101,16 @@ For more alarm configuration details, please refer to the [alarm doc](../setup/b **Note** that the `Unhealthy` event above is only for demonstration, they are not detected by default in SkyWalking, however, you can use the methods in [How to Report Events](#how-to-report-events) to report this kind of events. +## Correlation between events and metrics + +SkyWalking UI visualizes the events in the dashboard when the event service / instance / endpoint matches the displayed +service / instance / endpoint. + +By default, SkyWalking also generates some metrics for events by using [OAL](oal.md). The default metrics list of event +may change over time, you can find the complete list +in [event.oal](../../../oap-server/server-bootstrap/src/main/resources/oal/event.oal). If you want to generate you +custom metrics from events, please refer to [OAL](oal.md) about how to write OAL rules. + ## Known Events | Name | Type | When | Where | diff --git a/docs/en/concepts-and-designs/scope-definitions.md b/docs/en/concepts-and-designs/scope-definitions.md index 0bdb725c4..9be43de53 100644 --- a/docs/en/concepts-and-designs/scope-definitions.md +++ b/docs/en/concepts-and-designs/scope-definitions.md @@ -181,7 +181,7 @@ This calculates the metrics data from each request of the page in the browser ap ### SCOPE `BrowserAppPagePerf` -This calculates the metrics data form each request of the page in the browser application (browser only). +This calculates the metrics data from each request of the page in the browser application (browser only). | Name | Remarks | Group Key | Type | |---|---|---|---| @@ -201,3 +201,17 @@ This calculates the metrics data form each request of the page in the browser ap | ttlTime | Time to interact. | | int(in ms) | | firstPackTime | First pack time. | | int(in ms) | | fmpTime | First Meaningful Paint. | | int(in ms) | + +### SCOPE `Event` + +This calculates the metrics data from [events](event.md). + +| Name | Remarks | Group Key | Type | +|---|---|---|---| +| name | The name of the event. | | string | +| service | The service name to which the event belongs to. | | string | +| serviceInstance | The service instance to which the event belongs to, if any. | | string| +| endpoint | The service endpoint to which the event belongs to, if any. | | string| +| type | The type of the event, `Normal` or `Error`. | | string| +| message | The message of the event. | | string | +| parameters | The parameters in the `message`, see [parameters](event.md#parameters). | | string | diff --git a/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/EventAnalyzerModuleProvider.java b/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/EventAnalyzerModuleProvider.java index d29e4cb17..16b7a2765 100644 --- a/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/EventAnalyzerModuleProvider.java +++ b/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/EventAnalyzerModuleProvider.java @@ -20,9 +20,11 @@ package org.apache.skywalking.oap.server.analyzer.event; import org.apache.skywalking.oap.server.analyzer.event.listener.EventRecordAnalyzerListener; import org.apache.skywalking.oap.server.core.CoreModule; +import org.apache.skywalking.oap.server.core.oal.rt.OALEngineLoaderService; import org.apache.skywalking.oap.server.library.module.ModuleConfig; import org.apache.skywalking.oap.server.library.module.ModuleDefine; import org.apache.skywalking.oap.server.library.module.ModuleProvider; +import org.apache.skywalking.oap.server.library.module.ModuleStartException; import org.apache.skywalking.oap.server.library.module.ServiceNotProvidedException; public class EventAnalyzerModuleProvider extends ModuleProvider { @@ -51,7 +53,12 @@ public class EventAnalyzerModuleProvider extends ModuleProvider { } @Override - public void start() { + public void start() throws ModuleStartException { + getManager().find(CoreModule.NAME) + .provider() + .getService(OALEngineLoaderService.class) + .load(EventOALDefine.INSTANCE); + analysisService.add(new EventRecordAnalyzerListener.Factory(getManager())); } diff --git a/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/EventOALDefine.java b/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/EventOALDefine.java new file mode 100644 index 000000000..8aa2ef788 --- /dev/null +++ b/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/EventOALDefine.java @@ -0,0 +1,34 @@ +/* + * 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.analyzer.event; + +import org.apache.skywalking.oap.server.core.oal.rt.OALDefine; + +/** + * OAL rules to calculate Event-specific metrics. + */ +public class EventOALDefine extends OALDefine { + public static final EventOALDefine INSTANCE = new EventOALDefine(); + + private EventOALDefine() { + super( + "oal/event.oal", + "org.apache.skywalking.oap.server.core.source" + ); + } +} diff --git a/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/listener/EventRecordAnalyzerListener.java b/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/listener/EventRecordAnalyzerListener.java index 73eba625f..b8c0d1853 100644 --- a/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/listener/EventRecordAnalyzerListener.java +++ b/oap-server/analyzer/event-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/event/listener/EventRecordAnalyzerListener.java @@ -25,7 +25,8 @@ import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.analysis.TimeBucket; import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor; import org.apache.skywalking.oap.server.core.config.NamingControl; -import org.apache.skywalking.oap.server.core.event.Event; +import org.apache.skywalking.oap.server.core.source.Event; +import org.apache.skywalking.oap.server.core.source.SourceReceiver; import org.apache.skywalking.oap.server.library.module.ModuleManager; /** @@ -37,11 +38,14 @@ public class EventRecordAnalyzerListener implements EventAnalyzerListener { private final NamingControl namingControl; + private final SourceReceiver sourceReceiver; + private final Event event = new Event(); @Override public void build() { MetricsStreamProcessor.getInstance().in(event); + sourceReceiver.receive(event); } @Override @@ -72,16 +76,20 @@ public class EventRecordAnalyzerListener implements EventAnalyzerListener { public static class Factory implements EventAnalyzerListener.Factory { private final NamingControl namingControl; + private final SourceReceiver sourceReceiver; public Factory(final ModuleManager moduleManager) { this.namingControl = moduleManager.find(CoreModule.NAME) .provider() .getService(NamingControl.class); + this.sourceReceiver = moduleManager.find(CoreModule.NAME) + .provider() + .getService(SourceReceiver.class); } @Override public EventAnalyzerListener create(final ModuleManager moduleManager) { - return new EventRecordAnalyzerListener(namingControl); + return new EventRecordAnalyzerListener(namingControl, sourceReceiver); } } } diff --git a/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALLexer.g4 b/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALLexer.g4 index 86830f20f..528676080 100644 --- a/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALLexer.g4 +++ b/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALLexer.g4 @@ -39,6 +39,7 @@ SRC_SERVICE_INSTANCE_CLR_CPU: 'ServiceInstanceCLRCPU'; SRC_SERVICE_INSTANCE_CLR_GC: 'ServiceInstanceCLRGC'; SRC_SERVICE_INSTANCE_CLR_THREAD: 'ServiceInstanceCLRThread'; SRC_ENVOY_INSTANCE_METRIC: 'EnvoyInstanceMetric'; +SRC_EVENT: 'Event'; // Browser keywords SRC_BROWSER_APP_PERF: 'BrowserAppPerf'; diff --git a/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALParser.g4 b/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALParser.g4 index 08f40087c..f8eef4246 100644 --- a/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALParser.g4 +++ b/oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALParser.g4 @@ -55,7 +55,8 @@ source SRC_SERVICE_INSTANCE_CLR_CPU | SRC_SERVICE_INSTANCE_CLR_GC | SRC_SERVICE_INSTANCE_CLR_THREAD | SRC_ENVOY_INSTANCE_METRIC | SRC_BROWSER_APP_PERF | SRC_BROWSER_APP_PAGE_PERF | SRC_BROWSER_APP_SINGLE_VERSION_PERF | - SRC_BROWSER_APP_TRAFFIC | SRC_BROWSER_APP_PAGE_TRAFFIC | SRC_BROWSER_APP_SINGLE_VERSION_TRAFFIC + SRC_BROWSER_APP_TRAFFIC | SRC_BROWSER_APP_PAGE_TRAFFIC | SRC_BROWSER_APP_SINGLE_VERSION_TRAFFIC | + SRC_EVENT ; disableSource diff --git a/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/dispatch.ftl b/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/dispatch.ftl index 4503eda61..e9d16cb16 100644 --- a/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/dispatch.ftl +++ b/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/dispatch.ftl @@ -1,6 +1,6 @@ -public void dispatch(org.apache.skywalking.oap.server.core.source.Source source) { +public void dispatch(org.apache.skywalking.oap.server.core.source.ISource source) { ${sourcePackage}${source} _source = (${sourcePackage}${source})source; <#list metrics as metrics> do${metrics.metricsName}(_source); -} \ No newline at end of file +} diff --git a/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/doMetrics.ftl b/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/doMetrics.ftl index 45ee4c405..631296035 100644 --- a/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/doMetrics.ftl +++ b/oap-server/oal-rt/src/main/resources/code-templates/dispatcher/doMetrics.ftl @@ -1,5 +1,4 @@ private void do${metricsName}(${sourcePackage}${sourceName} source) { -${metricsClassPackage}${metricsName}Metrics metrics = new ${metricsClassPackage}${metricsName}Metrics(); <#if filterExpressions??> <#list filterExpressions as filterExpression> @@ -9,6 +8,7 @@ ${metricsClassPackage}${metricsName}Metrics metrics = new ${metricsClassPackage} +${metricsClassPackage}${metricsName}Metrics metrics = new ${metricsClassPackage}${metricsName}Metrics(); metrics.setTimeBucket(source.getTimeBucket()); <#list fieldsFromSource as field> metrics.${field.fieldSetter}(source.${field.fieldGetter}()); diff --git a/oap-server/server-bootstrap/src/main/resources/oal/event.oal b/oap-server/server-bootstrap/src/main/resources/oal/event.oal new file mode 100755 index 000000000..0eb067b69 --- /dev/null +++ b/oap-server/server-bootstrap/src/main/resources/oal/event.oal @@ -0,0 +1,25 @@ +/* + * 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. + * + */ + +event_total = from(Event.*).count(); + +event_normal_count = from(Event.*).filter(type == "Normal").count(); +event_error_count = from(Event.*).filter(type == "Error").count(); + +event_start_count = from(Event.*).filter(name == "Start").count(); +event_shutdown_count = from(Event.*).filter(name == "Shutdown").count(); diff --git a/oap-server/server-bootstrap/src/main/resources/ui-initialized-templates/event.yml b/oap-server/server-bootstrap/src/main/resources/ui-initialized-templates/event.yml new file mode 100644 index 000000000..45e803225 --- /dev/null +++ b/oap-server/server-bootstrap/src/main/resources/ui-initialized-templates/event.yml @@ -0,0 +1,56 @@ +# 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. + +templates: + - name: "Event" + type: "DASHBOARD" + configuration: |- + [ + { + "name": "Event", + "type": "service", + "children": [ + { + "name": "Global", + "children": [ + { + "width": "3", + "title": "Event Count by Severity", + "height": "280", + "entityType": "All", + "independentSelector": false, + "metricType": "REGULAR_VALUE", + "metricName": "event_total,event_normal_count,event_error_count", + "queryMetricType": "readMetricsValues", + "chartType": "ChartLine" + }, + { + "width": "3", + "title": "Event Count by Lifecycle", + "height": "280", + "entityType": "All", + "independentSelector": false, + "metricType": "REGULAR_VALUE", + "metricName": "event_start_count,event_shutdown_count", + "queryMetricType": "readMetricsValues", + "chartType": "ChartLine" + } + ] + } + ] + } + ] + activated: true + disabled: false diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/DispatcherManager.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/DispatcherManager.java index 9f3f7ee58..890d1a4bd 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/DispatcherManager.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/DispatcherManager.java @@ -29,7 +29,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import org.apache.skywalking.oap.server.core.UnexpectedException; -import org.apache.skywalking.oap.server.core.source.Source; +import org.apache.skywalking.oap.server.core.source.ISource; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -43,7 +43,7 @@ public class DispatcherManager implements DispatcherDetectorListener { this.dispatcherMap = new HashMap<>(); } - public void forward(Source source) { + public void forward(ISource source) { if (source == null) { return; } @@ -96,12 +96,12 @@ public class DispatcherManager implements DispatcherDetectorListener { Object source = ((Class) argument).newInstance(); - if (!Source.class.isAssignableFrom(source.getClass())) { + if (!ISource.class.isAssignableFrom(source.getClass())) { throw new UnexpectedException( "unexpected type argument of class " + aClass.getName() + ", should be `org.apache.skywalking.oap.server.core.source.Source`. "); } - Source dispatcherSource = (Source) source; + ISource dispatcherSource = (ISource) source; SourceDispatcher dispatcher = (SourceDispatcher) aClass.newInstance(); int scopeId = dispatcherSource.scope(); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/SourceDispatcher.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/SourceDispatcher.java index 8f52e85e4..cba04555b 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/SourceDispatcher.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/SourceDispatcher.java @@ -18,7 +18,7 @@ package org.apache.skywalking.oap.server.core.analysis; -import org.apache.skywalking.oap.server.core.source.Source; +import org.apache.skywalking.oap.server.core.source.ISource; /** * SourceDispatcher implementation processes different types of the source. There are two kinds of the source @@ -29,6 +29,6 @@ import org.apache.skywalking.oap.server.core.source.Source; * * @param the data type of this dispatcher processes. */ -public interface SourceDispatcher { +public interface SourceDispatcher { void dispatch(SOURCE source); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/event/Event.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Event.java similarity index 92% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/event/Event.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Event.java index 37269be5b..17be3ae00 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/event/Event.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Event.java @@ -16,14 +16,13 @@ * */ -package org.apache.skywalking.oap.server.core.event; +package org.apache.skywalking.oap.server.core.source; import java.util.HashMap; import java.util.Map; import lombok.EqualsAndHashCode; import lombok.Getter; import lombok.Setter; -import org.apache.skywalking.apm.util.StringUtil; import org.apache.skywalking.oap.server.core.analysis.IDManager; import org.apache.skywalking.oap.server.core.analysis.MetricsExtension; import org.apache.skywalking.oap.server.core.analysis.Stream; @@ -34,12 +33,10 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.MetricsMetaInfo; import org.apache.skywalking.oap.server.core.analysis.metrics.WithMetadata; import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor; import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData; -import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine; -import org.apache.skywalking.oap.server.core.source.ScopeDeclaration; import org.apache.skywalking.oap.server.core.storage.StorageHashMapBuilder; import org.apache.skywalking.oap.server.core.storage.annotation.Column; -import org.elasticsearch.common.Strings; +import static org.apache.skywalking.apm.util.StringUtil.isNotBlank; import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EVENT; @Getter @@ -51,7 +48,7 @@ import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EV of = "uuid" ) @MetricsExtension(supportDownSampling = false, supportUpdate = true) -public class Event extends Metrics implements WithMetadata, LongValueHolder { +public class Event extends Metrics implements ISource, WithMetadata, LongValueHolder { public static final String INDEX_NAME = "events"; @@ -136,13 +133,13 @@ public class Event extends Metrics implements WithMetadata, LongValueHolder { setEndTime(event.getEndTime()); } - if (StringUtil.isNotBlank(event.getType())) { + if (isNotBlank(event.getType())) { setType(event.getType()); } - if (StringUtil.isNotBlank(event.getMessage())) { + if (isNotBlank(event.getMessage())) { setType(event.getMessage()); } - if (StringUtil.isNotBlank(event.getParameters())) { + if (isNotBlank(event.getParameters())) { setParameters(event.getParameters()); } return true; @@ -206,16 +203,31 @@ public class Event extends Metrics implements WithMetadata, LongValueHolder { @Override public MetricsMetaInfo getMeta() { int scope = DefaultScopeDefine.SERVICE; + if (isNotBlank(getServiceInstance())) { + scope = DefaultScopeDefine.SERVICE_INSTANCE; + } else if (isNotBlank(getEndpoint())) { + scope = DefaultScopeDefine.ENDPOINT; + } + + String id = getEntityId(); + return new MetricsMetaInfo(getName(), scope, id); + } + + @Override + public int scope() { + return EVENT; + } + + @Override + public String getEntityId() { final String serviceId = IDManager.ServiceID.buildId(getService(), true); String id = serviceId; - if (!Strings.isNullOrEmpty(getServiceInstance())) { - scope = DefaultScopeDefine.SERVICE_INSTANCE; + if (isNotBlank(getServiceInstance())) { id = IDManager.ServiceInstanceID.buildId(serviceId, getServiceInstance()); - } else if (!Strings.isNullOrEmpty(getEndpoint())) { - scope = DefaultScopeDefine.ENDPOINT; + } else if (isNotBlank(getEndpoint())) { id = IDManager.EndpointID.buildId(serviceId, getEndpoint()); } - return new MetricsMetaInfo(getName(), scope, id); + return id; } public static class Builder implements StorageHashMapBuilder { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/ISource.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/ISource.java new file mode 100644 index 000000000..479212d5f --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/ISource.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.core.source; + +public interface ISource { + int scope(); + + long getTimeBucket(); + + void setTimeBucket(long timeBucket); + + String getEntityId(); + + /** + * Internal data field preparation before {@link org.apache.skywalking.oap.server.core.analysis.SourceDispatcher#dispatch(ISource)} + */ + default void prepare() { + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Source.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Source.java index 25d01cd08..5c9972258 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Source.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Source.java @@ -21,19 +21,9 @@ package org.apache.skywalking.oap.server.core.source; import lombok.Getter; import lombok.Setter; -public abstract class Source { - public abstract int scope(); +public abstract class Source implements ISource { @Getter @Setter private long timeBucket; - - public abstract String getEntityId(); - - /** - * Internal data field preparation before {@link org.apache.skywalking.oap.server.core.analysis.SourceDispatcher#dispatch(Source)} - */ - public void prepare() { - - } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiver.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiver.java index c7a05a170..66cef4f48 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiver.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiver.java @@ -26,7 +26,7 @@ import org.apache.skywalking.oap.server.library.module.Service; * in order to forward source to the suitable real {@link org.apache.skywalking.oap.server.core.analysis.SourceDispatcher}. */ public interface SourceReceiver extends Service { - void receive(Source source); + void receive(ISource source); DispatcherDetectorListener getDispatcherDetectorListener(); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java index ae4564dea..1e01da819 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java @@ -32,7 +32,7 @@ public class SourceReceiverImpl implements SourceReceiver { } @Override - public void receive(Source source) { + public void receive(ISource source) { dispatcherManager.forward(source); } diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandlerTest.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandlerTest.java index c23b00457..ade713fe2 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandlerTest.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandlerTest.java @@ -23,16 +23,16 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.utils.Bytes; import org.apache.skywalking.apm.network.management.v3.InstancePingPkg; import org.apache.skywalking.apm.network.management.v3.InstanceProperties; -import org.apache.skywalking.oap.server.core.CoreModule; -import org.apache.skywalking.oap.server.core.config.NamingControl; -import org.apache.skywalking.oap.server.core.config.group.EndpointNameGrouping; -import org.apache.skywalking.oap.server.core.source.ServiceInstanceUpdate; -import org.apache.skywalking.oap.server.core.source.Source; -import org.apache.skywalking.oap.server.core.source.SourceReceiver; -import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.apache.skywalking.oap.server.analyzer.agent.kafka.mock.MockModuleManager; import org.apache.skywalking.oap.server.analyzer.agent.kafka.mock.MockModuleProvider; import org.apache.skywalking.oap.server.analyzer.agent.kafka.module.KafkaFetcherConfig; +import org.apache.skywalking.oap.server.core.CoreModule; +import org.apache.skywalking.oap.server.core.config.NamingControl; +import org.apache.skywalking.oap.server.core.config.group.EndpointNameGrouping; +import org.apache.skywalking.oap.server.core.source.ISource; +import org.apache.skywalking.oap.server.core.source.ServiceInstanceUpdate; +import org.apache.skywalking.oap.server.core.source.SourceReceiver; +import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.junit.Assert; import org.junit.Before; import org.junit.ClassRule; @@ -52,7 +52,7 @@ public class ServiceManagementHandlerTest { public static SourceReceiverRule SOURCE_RECEIVER = new SourceReceiverRule() { @Override - protected void verify(final List sourceList) throws Throwable { + protected void verify(final List sourceList) throws Throwable { ServiceInstanceUpdate instanceUpdate = (ServiceInstanceUpdate) sourceList.get(0); Assert.assertEquals(instanceUpdate.getName(), SERVICE_INSTANCE); diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/SourceReceiverRule.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/SourceReceiverRule.java index aa80b09c4..9be40b961 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/SourceReceiverRule.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/test/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/SourceReceiverRule.java @@ -21,15 +21,15 @@ package org.apache.skywalking.oap.server.analyzer.agent.kafka.provider.handler; import com.google.common.collect.Lists; import java.util.List; import org.apache.skywalking.oap.server.core.analysis.DispatcherDetectorListener; -import org.apache.skywalking.oap.server.core.source.Source; +import org.apache.skywalking.oap.server.core.source.ISource; import org.apache.skywalking.oap.server.core.source.SourceReceiver; import org.junit.rules.Verifier; public abstract class SourceReceiverRule extends Verifier implements SourceReceiver { - private final List sourceList = Lists.newArrayList(); + private final List sourceList = Lists.newArrayList(); @Override - public void receive(final Source source) { + public void receive(final ISource source) { sourceList.add(source); } @@ -38,10 +38,11 @@ public abstract class SourceReceiverRule extends Verifier implements SourceRecei return null; } + @Override protected void verify() throws Throwable { verify(sourceList); } - protected abstract void verify(List sourceList) throws Throwable; + protected abstract void verify(List sourceList) throws Throwable; } diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MockReceiver.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MockReceiver.java index 94437444b..3ba6a36e3 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MockReceiver.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MockReceiver.java @@ -22,7 +22,7 @@ import java.util.ArrayList; import java.util.List; import lombok.Getter; import org.apache.skywalking.oap.server.core.analysis.DispatcherDetectorListener; -import org.apache.skywalking.oap.server.core.source.Source; +import org.apache.skywalking.oap.server.core.source.ISource; import org.apache.skywalking.oap.server.core.source.SourceReceiver; /** @@ -30,10 +30,10 @@ import org.apache.skywalking.oap.server.core.source.SourceReceiver; */ public class MockReceiver implements SourceReceiver { @Getter - private List receivedSources = new ArrayList<>(); + private List receivedSources = new ArrayList<>(); @Override - public void receive(final Source source) { + public void receive(final ISource source) { receivedSources.add(source); } diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MultiScopesAnalysisListenerTest.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MultiScopesAnalysisListenerTest.java index 0b443ebf8..5d37e7092 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MultiScopesAnalysisListenerTest.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/MultiScopesAnalysisListenerTest.java @@ -43,12 +43,12 @@ import org.apache.skywalking.oap.server.core.source.All; import org.apache.skywalking.oap.server.core.source.DatabaseAccess; import org.apache.skywalking.oap.server.core.source.Endpoint; import org.apache.skywalking.oap.server.core.source.EndpointRelation; +import org.apache.skywalking.oap.server.core.source.ISource; import org.apache.skywalking.oap.server.core.source.Service; import org.apache.skywalking.oap.server.core.source.ServiceInstance; import org.apache.skywalking.oap.server.core.source.ServiceInstanceRelation; import org.apache.skywalking.oap.server.core.source.ServiceMeta; import org.apache.skywalking.oap.server.core.source.ServiceRelation; -import org.apache.skywalking.oap.server.core.source.Source; import org.junit.Assert; import org.junit.Before; import org.junit.Test; @@ -147,7 +147,7 @@ public class MultiScopesAnalysisListenerTest { listener.parseEntry(spanObject, segment); listener.build(); - final List receivedSources = mockReceiver.getReceivedSources(); + final List receivedSources = mockReceiver.getReceivedSources(); Assert.assertEquals(7, receivedSources.size()); final All all = (All) receivedSources.get(0); final Service service = (Service) receivedSources.get(1); @@ -212,7 +212,7 @@ public class MultiScopesAnalysisListenerTest { listener.parseEntry(spanObject, segment); listener.build(); - final List receivedSources = mockReceiver.getReceivedSources(); + final List receivedSources = mockReceiver.getReceivedSources(); Assert.assertEquals(7, receivedSources.size()); final All all = (All) receivedSources.get(0); final Service service = (Service) receivedSources.get(1); @@ -276,7 +276,7 @@ public class MultiScopesAnalysisListenerTest { listener.parseEntry(spanObject, segment); listener.build(); - final List receivedSources = mockReceiver.getReceivedSources(); + final List receivedSources = mockReceiver.getReceivedSources(); Assert.assertEquals(7, receivedSources.size()); final All all = (All) receivedSources.get(0); final Service service = (Service) receivedSources.get(1); @@ -332,7 +332,7 @@ public class MultiScopesAnalysisListenerTest { listener.parseLocal(spanObject, segment); listener.build(); - final List receivedSources = mockReceiver.getReceivedSources(); + final List receivedSources = mockReceiver.getReceivedSources(); Assert.assertEquals(1, receivedSources.size()); final Endpoint source = (Endpoint) receivedSources.get(0); Assert.assertEquals("/logic-call", source.getName()); @@ -377,7 +377,7 @@ public class MultiScopesAnalysisListenerTest { listener.parseLocal(spanObject, segment); listener.build(); - final List receivedSources = mockReceiver.getReceivedSources(); + final List receivedSources = mockReceiver.getReceivedSources(); Assert.assertEquals(1, receivedSources.size()); final Endpoint source = (Endpoint) receivedSources.get(0); Assert.assertEquals("/GraphQL-service", source.getName()); @@ -416,7 +416,7 @@ public class MultiScopesAnalysisListenerTest { listener.parseExit(spanObject, segment); listener.build(); - final List receivedSources = mockReceiver.getReceivedSources(); + final List receivedSources = mockReceiver.getReceivedSources(); Assert.assertEquals(4, receivedSources.size()); final ServiceRelation serviceRelation = (ServiceRelation) receivedSources.get(0); final ServiceInstanceRelation serviceInstanceRelation = (ServiceInstanceRelation) receivedSources.get(1); @@ -462,7 +462,7 @@ public class MultiScopesAnalysisListenerTest { listener.parseExit(spanObject, segment); listener.build(); - final List receivedSources = mockReceiver.getReceivedSources(); + final List receivedSources = mockReceiver.getReceivedSources(); Assert.assertEquals(2, receivedSources.size()); final ServiceRelation serviceRelation = (ServiceRelation) receivedSources.get(0); final ServiceInstanceRelation serviceInstanceRelation = (ServiceInstanceRelation) receivedSources.get(1); diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/ESEventQueryDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/ESEventQueryDAO.java index 749037b47..7aee3e7d8 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/ESEventQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/ESEventQueryDAO.java @@ -23,7 +23,7 @@ import java.util.List; import java.util.Objects; import java.util.stream.Collectors; import java.util.stream.Stream; -import org.apache.skywalking.oap.server.core.event.Event; +import org.apache.skywalking.oap.server.core.source.Event; import org.apache.skywalking.oap.server.core.query.enumeration.Order; import org.apache.skywalking.oap.server.core.query.input.Duration; import org.apache.skywalking.oap.server.core.query.type.event.EventQueryCondition; diff --git a/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/ES7EventQueryDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/ES7EventQueryDAO.java index 2ddfd7aae..c0df97b4e 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/ES7EventQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch7-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch7/query/ES7EventQueryDAO.java @@ -22,7 +22,7 @@ import java.io.IOException; import java.util.List; import java.util.stream.Collectors; import java.util.stream.Stream; -import org.apache.skywalking.oap.server.core.event.Event; +import org.apache.skywalking.oap.server.core.source.Event; import org.apache.skywalking.oap.server.core.query.type.event.EventQueryCondition; import org.apache.skywalking.oap.server.core.query.type.event.Events; import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; diff --git a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/EventQueryDAO.java b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/EventQueryDAO.java index ab36ca28e..1d3689227 100644 --- a/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/EventQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-influxdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/influxdb/query/EventQueryDAO.java @@ -26,7 +26,7 @@ import java.util.Map; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.apache.skywalking.oap.server.core.event.Event; +import org.apache.skywalking.oap.server.core.source.Event; import org.apache.skywalking.oap.server.core.query.enumeration.Order; import org.apache.skywalking.oap.server.core.query.input.Duration; import org.apache.skywalking.oap.server.core.query.type.event.EventQueryCondition; diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2EventQueryDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2EventQueryDAO.java index 43d8f9511..c29dfc927 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2EventQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2EventQueryDAO.java @@ -28,7 +28,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.apache.skywalking.oap.server.core.event.Event; +import org.apache.skywalking.oap.server.core.source.Event; import org.apache.skywalking.oap.server.core.query.input.Duration; import org.apache.skywalking.oap.server.core.query.type.event.EventQueryCondition; import org.apache.skywalking.oap.server.core.query.type.event.EventType; diff --git a/oap-server/server-tools/profile-exporter/tool-profile-snapshot-server-mock/src/main/java/org/apache/skywalking/oap/server/tool/profile/core/mock/MockSourceReceiver.java b/oap-server/server-tools/profile-exporter/tool-profile-snapshot-server-mock/src/main/java/org/apache/skywalking/oap/server/tool/profile/core/mock/MockSourceReceiver.java index e0f653da1..3da59d7b2 100644 --- a/oap-server/server-tools/profile-exporter/tool-profile-snapshot-server-mock/src/main/java/org/apache/skywalking/oap/server/tool/profile/core/mock/MockSourceReceiver.java +++ b/oap-server/server-tools/profile-exporter/tool-profile-snapshot-server-mock/src/main/java/org/apache/skywalking/oap/server/tool/profile/core/mock/MockSourceReceiver.java @@ -19,7 +19,7 @@ package org.apache.skywalking.oap.server.tool.profile.core.mock; import org.apache.skywalking.oap.server.core.analysis.DispatcherDetectorListener; -import org.apache.skywalking.oap.server.core.source.Source; +import org.apache.skywalking.oap.server.core.source.ISource; import org.apache.skywalking.oap.server.core.source.SourceReceiver; /** @@ -27,7 +27,7 @@ import org.apache.skywalking.oap.server.core.source.SourceReceiver; */ public class MockSourceReceiver implements SourceReceiver { @Override - public void receive(Source source) { + public void receive(ISource source) { } @Override