From 0103d8ad2892f888cb55a01a7efd39a44a22726f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=90=B4=E6=99=9F=20Wu=20Sheng?= Date: Fri, 14 Jun 2019 11:31:51 +0800 Subject: [PATCH] Fix no stream register. (#2873) * Fix no stream register. * Make sure list order. * Refactor codes. * Fix wrong revert. --- .../oap/server/core/CoreModuleProvider.java | 1 + .../worker/MetricsStreamProcessor.java | 4 +++ .../worker/InventoryStreamProcessor.java | 4 +++ .../core/remote/define/StreamDataMapping.java | 27 +++++++++++++++---- 4 files changed, 31 insertions(+), 5 deletions(-) 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 9430020c5..01ddd9f83 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 @@ -160,6 +160,7 @@ public class CoreModuleProvider extends ModuleProvider { annotationScan.scan(() -> { }); + streamDataMapping.init(); } catch (IOException | IllegalAccessException | InstantiationException e) { throw new ModuleStartException(e.getMessage(), e); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsStreamProcessor.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsStreamProcessor.java index e2505068e..20d09bf7f 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsStreamProcessor.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsStreamProcessor.java @@ -24,6 +24,7 @@ import org.apache.skywalking.oap.server.core.*; import org.apache.skywalking.oap.server.core.analysis.*; import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics; import org.apache.skywalking.oap.server.core.config.DownsamplingConfigService; +import org.apache.skywalking.oap.server.core.remote.define.StreamDataMappingSetter; import org.apache.skywalking.oap.server.core.storage.*; import org.apache.skywalking.oap.server.core.storage.annotation.Storage; import org.apache.skywalking.oap.server.core.storage.model.*; @@ -66,6 +67,9 @@ public class MetricsStreamProcessor implements StreamProcessor { IModelSetter modelSetter = moduleDefineHolder.find(CoreModule.NAME).provider().getService(IModelSetter.class); DownsamplingConfigService configService = moduleDefineHolder.find(CoreModule.NAME).provider().getService(DownsamplingConfigService.class); + StreamDataMappingSetter streamDataMappingSetter = moduleDefineHolder.find(CoreModule.NAME).provider().getService(StreamDataMappingSetter.class); + streamDataMappingSetter.putIfAbsent(metricsClass); + MetricsPersistentWorker hourPersistentWorker = null; MetricsPersistentWorker dayPersistentWorker = null; MetricsPersistentWorker monthPersistentWorker = null; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/InventoryStreamProcessor.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/InventoryStreamProcessor.java index 0e4b3dd76..596ab8171 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/InventoryStreamProcessor.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/InventoryStreamProcessor.java @@ -22,6 +22,7 @@ import java.util.*; import org.apache.skywalking.oap.server.core.*; import org.apache.skywalking.oap.server.core.analysis.*; import org.apache.skywalking.oap.server.core.register.RegisterSource; +import org.apache.skywalking.oap.server.core.remote.define.StreamDataMappingSetter; import org.apache.skywalking.oap.server.core.storage.*; import org.apache.skywalking.oap.server.core.storage.annotation.Storage; import org.apache.skywalking.oap.server.core.storage.model.*; @@ -57,6 +58,9 @@ public class InventoryStreamProcessor implements StreamProcessor IModelSetter modelSetter = moduleDefineHolder.find(CoreModule.NAME).provider().getService(IModelSetter.class); Model model = modelSetter.putIfAbsent(inventoryClass, stream.scopeId(), new Storage(stream.name(), false, false, Downsampling.None)); + StreamDataMappingSetter streamDataMappingSetter = moduleDefineHolder.find(CoreModule.NAME).provider().getService(StreamDataMappingSetter.class); + streamDataMappingSetter.putIfAbsent(inventoryClass); + RegisterPersistentWorker persistentWorker = new RegisterPersistentWorker(moduleDefineHolder, model.getName(), registerDAO, stream.scopeId()); RegisterRemoteWorker remoteWorker = new RegisterRemoteWorker(moduleDefineHolder, persistentWorker); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/define/StreamDataMapping.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/define/StreamDataMapping.java index 31eb34aad..00ef25b45 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/define/StreamDataMapping.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/define/StreamDataMapping.java @@ -25,12 +25,12 @@ import org.apache.skywalking.oap.server.core.remote.data.StreamData; * @author peng-yongsheng */ public class StreamDataMapping implements StreamDataMappingGetter, StreamDataMappingSetter { - - private int id = 0; + private List> streamClassList; private final Map, Integer> classMap; private final Map> idMap; public StreamDataMapping() { + streamClassList = new ArrayList<>(); this.classMap = new HashMap<>(); this.idMap = new HashMap<>(); } @@ -40,9 +40,26 @@ public class StreamDataMapping implements StreamDataMappingGetter, StreamDataMap return; } - id++; - classMap.put(streamDataClass, id); - idMap.put(id, streamDataClass); + streamClassList.add(streamDataClass); + } + + public void init() { + /** + * The stream protocol use this list order to assign the ID, + * which is used in across node communication. This order must be certain. + */ + Collections.sort(streamClassList, new Comparator() { + @Override public int compare(Class streamClass1, Class streamClass2) { + return streamClass1.getName().compareTo(streamClass2.getName()); + } + }); + + for (int i = 0; i < streamClassList.size(); i++) { + Class streamClass = streamClassList.get(i); + int streamId = i + 1; + classMap.put(streamClass, streamId); + idMap.put(streamId, streamClass); + } } @Override public int findIdByClass(Class streamDataClass) {