From dd6ae268cf1042fe4f0e62d969f81c02f4a3669f Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 23 Jul 2017 23:32:32 +0800 Subject: [PATCH] 1. remote worker use common service, send data object with bytes and data define id, use define id to find the grpc deserialize object. 2. merge the metric data and record data to data object. --- .../agentstream/worker/CommonTable.java | 9 ++ .../agentstream/worker/TimeSlice.java | 28 ++++ .../worker/config/CacheSizeConfig.java | 17 +++ .../worker/config/WorkerConfig.java | 125 ++++++++++++++++++ .../NodeComponentAggDayWorker.java | 25 ++++ .../node/define/NodeComponentDataDefine.java | 28 ++++ .../define/NodeComponentEsTableDefine.java | 2 +- .../define/NodeComponentH2TableDefine.java | 2 +- .../node/{ => define}/NodeComponentTable.java | 6 +- .../node/define/NodeMappingEsTableDefine.java | 33 +++++ .../node/define/NodeMappingH2TableDefine.java | 21 +++ .../worker/node/define/NodeMappingTable.java | 12 ++ .../src/main/proto/NodeComponent.proto | 11 ++ .../resources/META-INF/defines/data.define | 1 + .../resources/META-INF/defines/storage.define | 4 +- apm-collector/apm-collector-client/pom.xml | 8 ++ .../elasticsearch/ElasticSearchClient.java | 5 + .../core/storage/StorageBatchBuilder.java | 11 ++ apm-collector/apm-collector-remote/pom.xml | 5 + .../collector/remote/RemoteModuleDefine.java | 34 ++++- .../remote/RemoteModuleException.java | 17 +++ .../remote/grpc/RemoteGRPCConfig.java | 9 ++ .../remote/grpc/RemoteGRPCConfigParser.java | 27 ++++ .../remote/grpc/RemoteGRPCDataListener.java | 17 +++ .../remote/grpc/RemoteGRPCModuleDefine.java | 40 ++++-- .../grpc/RemoteGRPCModuleRegistration.java | 13 ++ .../handler/RemoteHandlerDefineException.java | 17 +++ .../handler/RemoteHandlerDefineLoader.java | 30 +++++ .../handler/RemoteHandlerDefinitionFile.java | 13 ++ .../src/main/proto/RemoteCommonService.proto | 4 +- apm-collector/apm-collector-stream/pom.xml | 5 + ...rWorker.java => AbstractRemoteWorker.java} | 8 +- ...java => AbstractRemoteWorkerProvider.java} | 12 +- .../apm/collector/stream/AbstractWorker.java | 12 +- ...terWorkerRef.java => RemoteWorkerRef.java} | 6 +- .../apm/collector/stream/WorkerContext.java | 16 ++- .../stream/WorkerModuleInstaller.java | 25 ---- .../stream/impl/AggregationWorker.java | 55 ++++++++ .../apm/collector/stream/impl/Const.java | 13 ++ .../stream/impl/GRPCRemoteWorker.java | 26 ++++ .../impl/RemoteCommonServiceHandler.java | 24 ++++ .../collector/stream/impl/data/Attribute.java | 28 ++++ .../stream/impl/data/AttributeType.java | 8 ++ .../apm/collector/stream/impl/data/Data.java | 50 +++++++ .../collector/stream/impl/data/DataCache.java | 30 +++++ .../stream/impl/data/DataCollection.java | 53 ++++++++ .../stream/impl/data/DataDefine.java | 78 +++++++++++ .../stream/impl/data/DataDefineLoader.java | 29 ++++ .../stream/impl/data/DataDefinitionFile.java | 13 ++ .../collector/stream/impl/data/Operation.java | 12 ++ .../collector/stream/impl/data/Window.java | 44 ++++++ .../impl/data/operate/CoverOperation.java | 20 +++ .../impl/data/operate/NonOperation.java | 20 +++ .../META-INF/defines/remote_handler.define | 1 + 54 files changed, 1093 insertions(+), 69 deletions(-) create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/CacheSizeConfig.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/WorkerConfig.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/aggregation/NodeComponentAggDayWorker.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentDataDefine.java rename apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/{ => define}/NodeComponentTable.java (50%) create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingEsTableDefine.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingH2TableDefine.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingTable.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/proto/NodeComponent.proto create mode 100644 apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define create mode 100644 apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageBatchBuilder.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleException.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfig.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfigParser.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCDataListener.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleRegistration.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineException.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineLoader.java create mode 100644 apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefinitionFile.java rename apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/{AbstractClusterWorker.java => AbstractRemoteWorker.java} (79%) rename apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/{AbstractClusterWorkerProvider.java => AbstractRemoteWorkerProvider.java} (66%) rename apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/{ClusterWorkerRef.java => RemoteWorkerRef.java} (61%) delete mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/WorkerModuleInstaller.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/AggregationWorker.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/Const.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/GRPCRemoteWorker.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/RemoteCommonServiceHandler.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Attribute.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/AttributeType.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Data.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCache.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCollection.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefine.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefineLoader.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefinitionFile.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Operation.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Window.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/CoverOperation.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/NonOperation.java create mode 100644 apm-collector/apm-collector-stream/src/main/resources/META-INF/defines/remote_handler.define diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java new file mode 100644 index 000000000..4f8b27ac4 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java @@ -0,0 +1,9 @@ +package org.skywalking.apm.collector.agentstream.worker; + +/** + * @author pengys5 + */ +public class CommonTable { + public static final String COLUMN_AGG = "agg"; + public static final String COLUMN_TIME_BUCKET = "time_bucket"; +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java new file mode 100644 index 000000000..c75132a0e --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java @@ -0,0 +1,28 @@ +package org.skywalking.apm.collector.agentstream.worker; + +/** + * @author pengys5 + */ +public abstract class TimeSlice { + private String sliceType; + private long startTime; + private long endTime; + + public TimeSlice(String sliceType, long startTime, long endTime) { + this.startTime = startTime; + this.endTime = endTime; + this.sliceType = sliceType; + } + + public String getSliceType() { + return sliceType; + } + + public long getStartTime() { + return startTime; + } + + public long getEndTime() { + return endTime; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/CacheSizeConfig.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/CacheSizeConfig.java new file mode 100644 index 000000000..fe792ff4f --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/CacheSizeConfig.java @@ -0,0 +1,17 @@ +package org.skywalking.apm.collector.agentstream.worker.config; + +/** + * @author pengys5 + */ +public class CacheSizeConfig { + + public static class Cache { + public static class Analysis { + public static int SIZE = 1024; + } + + public static class Persistence { + public static int SIZE = 5000; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/WorkerConfig.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/WorkerConfig.java new file mode 100644 index 000000000..9c94a6ae8 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/WorkerConfig.java @@ -0,0 +1,125 @@ +package org.skywalking.apm.collector.agentstream.worker.config; + +/** + * @author pengys5 + */ +public class WorkerConfig { + + public static class WorkerNum { + public static class Node { + public static class NodeCompAgg { + public static int VALUE = 2; + } + + public static class NodeMappingDayAgg { + public static int VALUE = 2; + } + + public static class NodeMappingHourAgg { + public static int VALUE = 2; + } + + public static class NodeMappingMinuteAgg { + public static int VALUE = 2; + } + } + + public static class NodeRef { + public static class NodeRefDayAgg { + public static int VALUE = 2; + } + + public static class NodeRefHourAgg { + public static int VALUE = 2; + } + + public static class NodeRefMinuteAgg { + public static int VALUE = 2; + } + + public static class NodeRefResSumDayAgg { + public static int VALUE = 2; + } + + public static class NodeRefResSumHourAgg { + public static int VALUE = 2; + } + + public static class NodeRefResSumMinuteAgg { + public static int VALUE = 2; + } + } + + public static class GlobalTrace { + public static class GlobalTraceAgg { + public static int VALUE = 2; + } + } + } + + public static class Queue { + public static class GlobalTrace { + public static class GlobalTraceAnalysis { + public static int SIZE = 1024; + } + } + + public static class Segment { + public static class SegmentAnalysis { + public static int SIZE = 1024; + } + + public static class SegmentCostAnalysis { + public static int SIZE = 4096; + } + + public static class SegmentExceptionAnalysis { + public static int SIZE = 4096; + } + } + + public static class Node { + public static class NodeCompAnalysis { + public static int SIZE = 1024; + } + + public static class NodeMappingDayAnalysis { + public static int SIZE = 1024; + } + + public static class NodeMappingHourAnalysis { + public static int SIZE = 1024; + } + + public static class NodeMappingMinuteAnalysis { + public static int SIZE = 1024; + } + } + + public static class NodeRef { + public static class NodeRefDayAnalysis { + public static int SIZE = 1024; + } + + public static class NodeRefHourAnalysis { + public static int SIZE = 1024; + } + + public static class NodeRefMinuteAnalysis { + public static int SIZE = 1024; + } + + public static class NodeRefResSumDayAnalysis { + public static int SIZE = 1024; + } + + public static class NodeRefResSumHourAnalysis { + public static int SIZE = 1024; + } + + public static class NodeRefResSumMinuteAnalysis { + public static int SIZE = 1024; + } + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/aggregation/NodeComponentAggDayWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/aggregation/NodeComponentAggDayWorker.java new file mode 100644 index 000000000..f73e122c1 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/aggregation/NodeComponentAggDayWorker.java @@ -0,0 +1,25 @@ +package org.skywalking.apm.collector.agentstream.worker.node.aggregation; + +import org.skywalking.apm.collector.stream.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.LocalWorkerContext; +import org.skywalking.apm.collector.stream.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.Role; +import org.skywalking.apm.collector.stream.impl.AggregationWorker; + +/** + * @author pengys5 + */ +public class NodeComponentAggDayWorker extends AggregationWorker { + + public NodeComponentAggDayWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected void sendToNext() { + + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentDataDefine.java new file mode 100644 index 000000000..4bc3e36ce --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentDataDefine.java @@ -0,0 +1,28 @@ +package org.skywalking.apm.collector.agentstream.worker.node.define; + +import org.skywalking.apm.collector.stream.impl.data.Attribute; +import org.skywalking.apm.collector.stream.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.impl.data.operate.CoverOperation; +import org.skywalking.apm.collector.stream.impl.data.operate.NonOperation; + +/** + * @author pengys5 + */ +public class NodeComponentDataDefine extends DataDefine { + + @Override protected int defineId() { + return 0; + } + + @Override protected int initialCapacity() { + return 4; + } + + @Override protected void attributeDefine() { + addAttribute(0, new Attribute("id", AttributeType.STRING, new NonOperation())); + addAttribute(1, new Attribute("name", AttributeType.STRING, new CoverOperation())); + addAttribute(2, new Attribute("peers", AttributeType.STRING, new CoverOperation())); + addAttribute(3, new Attribute("aggregation", AttributeType.STRING, new CoverOperation())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentEsTableDefine.java index c40c0856a..a2b4fe89e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentEsTableDefine.java @@ -1,6 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.node.define; -import org.skywalking.apm.collector.agentstream.worker.node.NodeComponentTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; @@ -28,5 +27,6 @@ public class NodeComponentEsTableDefine extends ElasticSearchTableDefine { @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_PEERS, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentH2TableDefine.java index 586c75bf3..8b6ab6c67 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentH2TableDefine.java @@ -1,6 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.node.define; -import org.skywalking.apm.collector.agentstream.worker.node.NodeComponentTable; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; @@ -16,5 +15,6 @@ public class NodeComponentH2TableDefine extends H2TableDefine { @Override public void initialize() { addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_NAME, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_PEERS, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/NodeComponentTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentTable.java similarity index 50% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/NodeComponentTable.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentTable.java index 1fccd3597..9d479cd10 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/NodeComponentTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentTable.java @@ -1,9 +1,11 @@ -package org.skywalking.apm.collector.agentstream.worker.node; +package org.skywalking.apm.collector.agentstream.worker.node.define; + +import org.skywalking.apm.collector.agentstream.worker.CommonTable; /** * @author pengys5 */ -public class NodeComponentTable { +public class NodeComponentTable extends CommonTable { public static final String TABLE = "node_component"; public static final String COLUMN_NAME = "name"; public static final String COLUMN_PEERS = "peers"; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingEsTableDefine.java new file mode 100644 index 000000000..ef8a1a2ef --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingEsTableDefine.java @@ -0,0 +1,33 @@ +package org.skywalking.apm.collector.agentstream.worker.node.define; + +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; + +/** + * @author pengys5 + */ +public class NodeMappingEsTableDefine extends ElasticSearchTableDefine { + + public NodeMappingEsTableDefine() { + super(NodeMappingTable.TABLE); + } + + @Override public int refreshInterval() { + return 0; + } + + @Override public int numberOfShards() { + return 2; + } + + @Override public int numberOfReplicas() { + return 0; + } + + @Override public void initialize() { + addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_PEERS, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingH2TableDefine.java new file mode 100644 index 000000000..f11217ccc --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingH2TableDefine.java @@ -0,0 +1,21 @@ +package org.skywalking.apm.collector.agentstream.worker.node.define; + +import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; +import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; + +/** + * @author pengys5 + */ +public class NodeMappingH2TableDefine extends H2TableDefine { + + public NodeMappingH2TableDefine() { + super(NodeMappingTable.TABLE); + } + + @Override public void initialize() { + addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_NAME, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_PEERS, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingTable.java new file mode 100644 index 000000000..26f728fa9 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingTable.java @@ -0,0 +1,12 @@ +package org.skywalking.apm.collector.agentstream.worker.node.define; + +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + +/** + * @author pengys5 + */ +public class NodeMappingTable extends CommonTable { + public static final String TABLE = "node_mapping"; + public static final String COLUMN_NAME = "name"; + public static final String COLUMN_PEERS = "peers"; +} diff --git a/apm-collector/apm-collector-agentstream/src/main/proto/NodeComponent.proto b/apm-collector/apm-collector-agentstream/src/main/proto/NodeComponent.proto new file mode 100644 index 000000000..4ad963754 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/proto/NodeComponent.proto @@ -0,0 +1,11 @@ +syntax = "proto3"; + +option java_multiple_files = true; +option java_package = "org.skywalking.apm.collector.agentstream.worker.node.define.proto"; + +message Message { + string id = 1; + string name = 2; + string peers = 3; + string aggregation = 4; +} \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define new file mode 100644 index 000000000..a25812b85 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define @@ -0,0 +1 @@ +org.skywalking.apm.collector.agentstream.worker.node.define.NodeComponentDataDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define index c0a307879..82a6af78a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define @@ -1,2 +1,4 @@ org.skywalking.apm.collector.agentstream.worker.node.define.NodeComponentEsTableDefine -org.skywalking.apm.collector.agentstream.worker.node.define.NodeComponentH2TableDefine \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.node.define.NodeComponentH2TableDefine +org.skywalking.apm.collector.agentstream.worker.node.define.NodeMappingEsTableDefine +org.skywalking.apm.collector.agentstream.worker.node.define.NodeMappingH2TableDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-client/pom.xml b/apm-collector/apm-collector-client/pom.xml index 2e2c47fff..01bd8c710 100644 --- a/apm-collector/apm-collector-client/pom.xml +++ b/apm-collector/apm-collector-client/pom.xml @@ -37,6 +37,10 @@ snakeyaml org.yaml + + log4j-api + org.apache.logging.log4j + @@ -52,6 +56,10 @@ slf4j-log4j12 org.slf4j + + log4j + log4j + diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java index b975a12c1..73865df2c 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java @@ -7,6 +7,7 @@ import java.util.List; import org.elasticsearch.action.admin.indices.create.CreateIndexResponse; import org.elasticsearch.action.admin.indices.delete.DeleteIndexResponse; import org.elasticsearch.action.admin.indices.exists.indices.IndicesExistsResponse; +import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.client.IndicesAdminClient; import org.elasticsearch.common.settings.Settings; import org.elasticsearch.common.transport.InetSocketTransportAddress; @@ -99,4 +100,8 @@ public class ElasticSearchClient implements Client { IndicesExistsResponse response = adminClient.prepareExists(indexName).get(); return response.isExists(); } + + public IndexRequestBuilder prepareIndex(String indexName) { + return null; + } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageBatchBuilder.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageBatchBuilder.java new file mode 100644 index 000000000..5f1665a76 --- /dev/null +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageBatchBuilder.java @@ -0,0 +1,11 @@ +package org.skywalking.apm.collector.core.storage; + +import java.util.List; +import java.util.Map; + +/** + * @author pengys5 + */ +public abstract class StorageBatchBuilder { + public abstract List build(C client, Map lastData); +} diff --git a/apm-collector/apm-collector-remote/pom.xml b/apm-collector/apm-collector-remote/pom.xml index 81c877d81..f20ebd870 100644 --- a/apm-collector/apm-collector-remote/pom.xml +++ b/apm-collector/apm-collector-remote/pom.xml @@ -29,6 +29,11 @@ apm-collector-server ${project.version} + + org.skywalking + apm-collector-cluster + ${project.version} + io.grpc grpc-netty diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java index ad6875615..87796d38c 100644 --- a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java @@ -1,9 +1,41 @@ package org.skywalking.apm.collector.remote; +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine; +import org.skywalking.apm.collector.core.client.ClientException; +import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; +import org.skywalking.apm.collector.core.cluster.ClusterModuleContext; +import org.skywalking.apm.collector.core.config.ConfigParseException; +import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.core.framework.DefineException; +import org.skywalking.apm.collector.core.framework.Handler; import org.skywalking.apm.collector.core.module.ModuleDefine; +import org.skywalking.apm.collector.core.server.Server; +import org.skywalking.apm.collector.core.server.ServerException; +import org.skywalking.apm.collector.core.server.ServerHolder; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public abstract class RemoteModuleDefine extends ModuleDefine { +public abstract class RemoteModuleDefine extends ModuleDefine implements ClusterDataListenerDefine { + + private final Logger logger = LoggerFactory.getLogger(RemoteModuleDefine.class); + + @Override + public final void initialize(Map config, ServerHolder serverHolder) throws DefineException, ClientException { + try { + configParser().parse(config); + Server server = server(); + serverHolder.holdServer(server, handlerList()); + + ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getDataMonitor().addListener(listener(), registration()); + } catch (ConfigParseException | ServerException e) { + throw new RemoteModuleException(e.getMessage(), e); + } + } + + public abstract List handlerList() throws DefineException; } diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleException.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleException.java new file mode 100644 index 000000000..87af6e986 --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleException.java @@ -0,0 +1,17 @@ +package org.skywalking.apm.collector.remote; + +import org.skywalking.apm.collector.core.module.ModuleException; + +/** + * @author pengys5 + */ +public class RemoteModuleException extends ModuleException { + + public RemoteModuleException(String message) { + super(message); + } + + public RemoteModuleException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfig.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfig.java new file mode 100644 index 000000000..6b465917a --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfig.java @@ -0,0 +1,9 @@ +package org.skywalking.apm.collector.remote.grpc; + +/** + * @author pengys5 + */ +public class RemoteGRPCConfig { + public static String HOST; + public static int PORT; +} diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfigParser.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfigParser.java new file mode 100644 index 000000000..73beb8a49 --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfigParser.java @@ -0,0 +1,27 @@ +package org.skywalking.apm.collector.remote.grpc; + +import java.util.Map; +import org.skywalking.apm.collector.core.config.ConfigParseException; +import org.skywalking.apm.collector.core.module.ModuleConfigParser; +import org.skywalking.apm.collector.core.util.ObjectUtils; +import org.skywalking.apm.collector.core.util.StringUtils; + +/** + * @author pengys5 + */ +public class RemoteGRPCConfigParser implements ModuleConfigParser { + + private static final String HOST = "host"; + private static final String PORT = "port"; + + @Override public void parse(Map config) throws ConfigParseException { + if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(HOST))) { + RemoteGRPCConfig.HOST = "localhost"; + } + if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(PORT))) { + RemoteGRPCConfig.PORT = 11800; + } else { + RemoteGRPCConfig.PORT = (Integer)config.get(PORT); + } + } +} diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCDataListener.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCDataListener.java new file mode 100644 index 000000000..7d2e96c84 --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCDataListener.java @@ -0,0 +1,17 @@ +package org.skywalking.apm.collector.remote.grpc; + +import org.skywalking.apm.collector.cluster.ClusterModuleDefine; +import org.skywalking.apm.collector.core.cluster.ClusterDataListener; +import org.skywalking.apm.collector.remote.RemoteModuleGroupDefine; + +/** + * @author pengys5 + */ +public class RemoteGRPCDataListener extends ClusterDataListener { + + public static final String PATH = ClusterModuleDefine.BASE_CATALOG + "." + RemoteModuleGroupDefine.GROUP_NAME + "." + RemoteGRPCModuleDefine.MODULE_NAME; + + @Override public String path() { + return PATH; + } +} diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java index 877952cff..22b884e43 100644 --- a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java @@ -1,31 +1,42 @@ package org.skywalking.apm.collector.remote.grpc; -import java.util.Map; +import java.util.List; import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.client.ClientException; import org.skywalking.apm.collector.core.client.DataMonitor; +import org.skywalking.apm.collector.core.cluster.ClusterDataListener; +import org.skywalking.apm.collector.core.config.ConfigException; import org.skywalking.apm.collector.core.framework.DefineException; +import org.skywalking.apm.collector.core.framework.Handler; import org.skywalking.apm.collector.core.module.ModuleConfigParser; import org.skywalking.apm.collector.core.module.ModuleRegistration; import org.skywalking.apm.collector.core.server.Server; -import org.skywalking.apm.collector.core.server.ServerHolder; import org.skywalking.apm.collector.remote.RemoteModuleDefine; import org.skywalking.apm.collector.remote.RemoteModuleGroupDefine; +import org.skywalking.apm.collector.remote.grpc.handler.RemoteHandlerDefineException; +import org.skywalking.apm.collector.remote.grpc.handler.RemoteHandlerDefineLoader; +import org.skywalking.apm.collector.server.grpc.GRPCServer; /** * @author pengys5 */ public class RemoteGRPCModuleDefine extends RemoteModuleDefine { + + public static final String MODULE_NAME = "remote"; + + @Override public String name() { + return MODULE_NAME; + } + @Override protected String group() { return RemoteModuleGroupDefine.GROUP_NAME; } @Override public boolean defaultModule() { - return false; + return true; } @Override protected ModuleConfigParser configParser() { - return null; + return new RemoteGRPCConfigParser(); } @Override protected Client createClient(DataMonitor dataMonitor) { @@ -33,18 +44,25 @@ public class RemoteGRPCModuleDefine extends RemoteModuleDefine { } @Override protected Server server() { - return null; + return new GRPCServer(RemoteGRPCConfig.HOST, RemoteGRPCConfig.PORT); } @Override protected ModuleRegistration registration() { - return null; + return new RemoteGRPCModuleRegistration(); } - @Override public void initialize(Map config, ServerHolder serverHolder) throws DefineException, ClientException { - + @Override public ClusterDataListener listener() { + return new RemoteGRPCDataListener(); } - @Override public String name() { - return null; + @Override public List handlerList() throws DefineException { + RemoteHandlerDefineLoader loader = new RemoteHandlerDefineLoader(); + List handlers = null; + try { + handlers = loader.load(); + } catch (ConfigException e) { + throw new RemoteHandlerDefineException(e.getMessage(), e); + } + return handlers; } } diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleRegistration.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleRegistration.java new file mode 100644 index 000000000..4f5a371d2 --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleRegistration.java @@ -0,0 +1,13 @@ +package org.skywalking.apm.collector.remote.grpc; + +import org.skywalking.apm.collector.core.module.ModuleRegistration; + +/** + * @author pengys5 + */ +public class RemoteGRPCModuleRegistration extends ModuleRegistration { + + @Override public Value buildValue() { + return new Value(RemoteGRPCConfig.HOST, RemoteGRPCConfig.PORT, null); + } +} diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineException.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineException.java new file mode 100644 index 000000000..7a5c6063b --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineException.java @@ -0,0 +1,17 @@ +package org.skywalking.apm.collector.remote.grpc.handler; + +import org.skywalking.apm.collector.core.framework.DefineException; + +/** + * @author pengys5 + */ +public class RemoteHandlerDefineException extends DefineException { + + public RemoteHandlerDefineException(String message) { + super(message); + } + + public RemoteHandlerDefineException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineLoader.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineLoader.java new file mode 100644 index 000000000..b125bc090 --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineLoader.java @@ -0,0 +1,30 @@ +package org.skywalking.apm.collector.remote.grpc.handler; + +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.collector.core.config.ConfigException; +import org.skywalking.apm.collector.core.framework.Handler; +import org.skywalking.apm.collector.core.framework.Loader; +import org.skywalking.apm.collector.core.util.DefinitionLoader; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class RemoteHandlerDefineLoader implements Loader> { + + private final Logger logger = LoggerFactory.getLogger(RemoteHandlerDefineLoader.class); + + @Override public List load() throws ConfigException { + List handlers = new ArrayList<>(); + + RemoteHandlerDefinitionFile definitionFile = new RemoteHandlerDefinitionFile(); + DefinitionLoader definitionLoader = DefinitionLoader.load(Handler.class, definitionFile); + for (Handler handler : definitionLoader) { + logger.info("loaded remote handler definition class: {}", handler.getClass().getName()); + handlers.add(handler); + } + return handlers; + } +} diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefinitionFile.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefinitionFile.java new file mode 100644 index 000000000..4231229c4 --- /dev/null +++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefinitionFile.java @@ -0,0 +1,13 @@ +package org.skywalking.apm.collector.remote.grpc.handler; + +import org.skywalking.apm.collector.core.framework.DefinitionFile; + +/** + * @author pengys5 + */ +public class RemoteHandlerDefinitionFile extends DefinitionFile { + + @Override protected String fileName() { + return "remote_handler.define"; + } +} diff --git a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto index ae1da50f0..bf0f2d8bc 100644 --- a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto +++ b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto @@ -10,8 +10,8 @@ service RemoteCommonService { message Message { string workerRole = 1; - int32 objectId = 2; - bytes objectBytes = 3; // the byte array of data object + int32 dataDefineId = 2; + bytes dataBytes = 3; } message Empty { diff --git a/apm-collector/apm-collector-stream/pom.xml b/apm-collector/apm-collector-stream/pom.xml index e9be3a610..9ace78e53 100644 --- a/apm-collector/apm-collector-stream/pom.xml +++ b/apm-collector/apm-collector-stream/pom.xml @@ -23,5 +23,10 @@ apm-collector-queue ${project.version} + + org.skywalking + apm-collector-remote + ${project.version} + \ No newline at end of file diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorker.java similarity index 79% rename from apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorker.java rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorker.java index cc2991b4d..20c552bee 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorker.java @@ -1,7 +1,7 @@ package org.skywalking.apm.collector.stream; /** - * The AbstractClusterWorker implementations represent workers, + * The AbstractRemoteWorker implementations represent workers, * which receive remote messages. *

* Usually, the implementations are doing persistent, or aggregate works. @@ -9,17 +9,17 @@ package org.skywalking.apm.collector.stream; * @author pengys5 * @since v3.0-2017 */ -public abstract class AbstractClusterWorker extends AbstractWorker { +public abstract class AbstractRemoteWorker extends AbstractWorker { /** - * Construct an AbstractClusterWorker with the worker role and context. + * Construct an AbstractRemoteWorker with the worker role and context. * * @param role If multi-workers are for load balance, they should be more likely called worker instance. Meaning, * each worker have multi instances. * @param clusterContext See {@link ClusterWorkerContext} * @param selfContext See {@link LocalWorkerContext} */ - protected AbstractClusterWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + protected AbstractRemoteWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { super(role, clusterContext, selfContext); } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java similarity index 66% rename from apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorkerProvider.java rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java index 3e1fd1337..3774be6cb 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorkerProvider.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java @@ -1,17 +1,17 @@ package org.skywalking.apm.collector.stream; /** - * The AbstractClusterWorkerProvider implementations represent providers, - * which create instance of cluster workers whose implemented {@link AbstractClusterWorker}. + * The AbstractRemoteWorkerProvider implementations represent providers, + * which create instance of cluster workers whose implemented {@link AbstractRemoteWorker}. *

* * @author pengys5 * @since v3.0-2017 */ -public abstract class AbstractClusterWorkerProvider extends AbstractWorkerProvider { +public abstract class AbstractRemoteWorkerProvider extends AbstractWorkerProvider { /** - * Create how many worker instance of {@link AbstractClusterWorker} in one jvm. + * Create how many worker instance of {@link AbstractRemoteWorker} in one jvm. * * @return The worker instance number. */ @@ -21,7 +21,7 @@ public abstract class AbstractClusterWorkerProvider dataDefineMap; + private Map> roleWorkers; public WorkerContext() { @@ -20,8 +23,11 @@ public abstract class WorkerContext implements Context { return this.roleWorkers; } - @Override - final public WorkerRefs lookup(Role role) throws WorkerNotFoundException { + public final DataDefine getDataDefine(int defineId) { + return dataDefineMap.get(defineId); + } + + @Override final public WorkerRefs lookup(Role role) throws WorkerNotFoundException { if (getRoleWorkers().containsKey(role.roleName())) { WorkerRefs refs = new WorkerRefs(getRoleWorkers().get(role.roleName()), role.workerSelector()); return refs; @@ -30,16 +36,14 @@ public abstract class WorkerContext implements Context { } } - @Override - final public void put(WorkerRef workerRef) { + @Override final public void put(WorkerRef workerRef) { if (!getRoleWorkers().containsKey(workerRef.getRole().roleName())) { getRoleWorkers().putIfAbsent(workerRef.getRole().roleName(), new ArrayList()); } getRoleWorkers().get(workerRef.getRole().roleName()).add(workerRef); } - @Override - final public void remove(WorkerRef workerRef) { + @Override final public void remove(WorkerRef workerRef) { getRoleWorkers().remove(workerRef.getRole().roleName()); } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/WorkerModuleInstaller.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/WorkerModuleInstaller.java deleted file mode 100644 index 5dd6095ca..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/WorkerModuleInstaller.java +++ /dev/null @@ -1,25 +0,0 @@ -package org.skywalking.apm.collector.stream; - -import java.util.Map; -import org.skywalking.apm.collector.core.client.ClientException; -import org.skywalking.apm.collector.core.framework.DefineException; -import org.skywalking.apm.collector.core.module.ModuleDefine; -import org.skywalking.apm.collector.core.module.ModuleInstaller; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * @author pengys5 - */ -public class WorkerModuleInstaller implements ModuleInstaller { - - private final Logger logger = LoggerFactory.getLogger(WorkerModuleInstaller.class); - - @Override public void install(Map moduleConfig, - Map moduleDefineMap) throws DefineException, ClientException { - logger.info("beginning worker module install"); - Map.Entry workerConfigEntry = moduleConfig.entrySet().iterator().next(); - ModuleDefine moduleDefine = moduleDefineMap.get(workerConfigEntry.getKey()); - moduleDefine.initialize(workerConfigEntry.getValue()); - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/AggregationWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/AggregationWorker.java new file mode 100644 index 000000000..bd892035c --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/AggregationWorker.java @@ -0,0 +1,55 @@ +package org.skywalking.apm.collector.stream.impl; + +import org.skywalking.apm.collector.core.queue.EndOfBatchCommand; +import org.skywalking.apm.collector.stream.AbstractLocalAsyncWorker; +import org.skywalking.apm.collector.stream.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.LocalWorkerContext; +import org.skywalking.apm.collector.stream.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.Role; +import org.skywalking.apm.collector.stream.WorkerException; +import org.skywalking.apm.collector.stream.impl.data.Data; +import org.skywalking.apm.collector.stream.impl.data.DataCache; + +/** + * @author pengys5 + */ +public abstract class AggregationWorker extends AbstractLocalAsyncWorker { + + private DataCache dataCache; + + public AggregationWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); + dataCache = new DataCache(); + } + + private int messageNum; + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected final void onWork(Object message) throws WorkerException { + if (message instanceof EndOfBatchCommand) { + sendToNext(); + } else { + messageNum++; + aggregate(message); + + if (messageNum >= 100) { + sendToNext(); + messageNum = 0; + } + } + } + + protected abstract void sendToNext(); + + protected final void aggregate(Object message) { + Data data = (Data)message; + if (dataCache.containsKey(data.id())) { + getClusterContext().getDataDefine(data.getDefineId()).mergeData(data, dataCache.get(data.id())); + } else { + dataCache.put(data.id(), data); + } + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/Const.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/Const.java new file mode 100644 index 000000000..1850a2ee2 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/Const.java @@ -0,0 +1,13 @@ +package org.skywalking.apm.collector.stream.impl; + +/** + * @author pengys5 + */ +public class Const { + public static final String ID_SPLIT = "..-.."; + public static final String IDS_SPLIT = "\\.\\.-\\.\\."; + public static final String PEERS_FRONT_SPLIT = "["; + public static final String PEERS_BEHIND_SPLIT = "]"; + public static final String USER_CODE = "User"; + public static final String RESULT = "result"; +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/GRPCRemoteWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/GRPCRemoteWorker.java new file mode 100644 index 000000000..3bd9a33fa --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/GRPCRemoteWorker.java @@ -0,0 +1,26 @@ +package org.skywalking.apm.collector.stream.impl; + +import org.skywalking.apm.collector.stream.AbstractRemoteWorker; +import org.skywalking.apm.collector.stream.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.LocalWorkerContext; +import org.skywalking.apm.collector.stream.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.Role; +import org.skywalking.apm.collector.stream.WorkerException; + +/** + * @author pengys5 + */ +public class GRPCRemoteWorker extends AbstractRemoteWorker { + + protected GRPCRemoteWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + + } + + @Override protected final void onWork(Object message) throws WorkerException { + + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/RemoteCommonServiceHandler.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/RemoteCommonServiceHandler.java new file mode 100644 index 000000000..ea560ef12 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/RemoteCommonServiceHandler.java @@ -0,0 +1,24 @@ +package org.skywalking.apm.collector.stream.impl; + +import com.google.protobuf.ByteString; +import io.grpc.stub.StreamObserver; +import org.skywalking.apm.collector.remote.grpc.proto.Empty; +import org.skywalking.apm.collector.remote.grpc.proto.Message; +import org.skywalking.apm.collector.remote.grpc.proto.RemoteCommonServiceGrpc; +import org.skywalking.apm.collector.server.grpc.GRPCHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class RemoteCommonServiceHandler extends RemoteCommonServiceGrpc.RemoteCommonServiceImplBase implements GRPCHandler { + + private final Logger logger = LoggerFactory.getLogger(RemoteCommonServiceHandler.class); + + @Override public void call(Message request, StreamObserver responseObserver) { + String workerRole = request.getWorkerRole(); + int dataDefineId = request.getDataDefineId(); + ByteString bytesData = request.getDataBytes(); + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Attribute.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Attribute.java new file mode 100644 index 000000000..3075517de --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Attribute.java @@ -0,0 +1,28 @@ +package org.skywalking.apm.collector.stream.impl.data; + +/** + * @author pengys5 + */ +public class Attribute { + private final String name; + private final AttributeType type; + private final Operation operation; + + public Attribute(String name, AttributeType type, Operation operation) { + this.name = name; + this.type = type; + this.operation = operation; + } + + public String getName() { + return name; + } + + public AttributeType getType() { + return type; + } + + public Operation getOperation() { + return operation; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/AttributeType.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/AttributeType.java new file mode 100644 index 000000000..84e317ec0 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/AttributeType.java @@ -0,0 +1,8 @@ +package org.skywalking.apm.collector.stream.impl.data; + +/** + * @author pengys5 + */ +public enum AttributeType { + STRING, LONG, FLOAT +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Data.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Data.java new file mode 100644 index 000000000..be8d4eb21 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Data.java @@ -0,0 +1,50 @@ +package org.skywalking.apm.collector.stream.impl.data; + +/** + * @author pengys5 + */ +public class Data { + private int defineId; + private String[] dataStrings; + private Long[] dataLongs; + private Float[] dataFloats; + + public Data(int defineId, int stringCapacity, int longCapacity, int floatCapacity) { + this.defineId = defineId; + this.dataStrings = new String[stringCapacity]; + this.dataLongs = new Long[longCapacity]; + this.dataFloats = new Float[floatCapacity]; + } + + public void setDataString(int position, String value) { + dataStrings[position] = value; + } + + public void setDataLong(int position, Long value) { + dataLongs[position] = value; + } + + public void setDataFloat(int position, Float value) { + dataFloats[position] = value; + } + + public String getDataString(int position) { + return dataStrings[position]; + } + + public Long getDataLong(int position) { + return dataLongs[position]; + } + + public Float getDataFloat(int position) { + return dataFloats[position]; + } + + public String id() { + return dataStrings[0]; + } + + public int getDefineId() { + return defineId; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCache.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCache.java new file mode 100644 index 000000000..db51cf91d --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCache.java @@ -0,0 +1,30 @@ +package org.skywalking.apm.collector.stream.impl.data; + +/** + * @author pengys5 + */ +public class DataCache extends Window { + + private DataCollection lockedDataCollection; + + public boolean containsKey(String id) { + return lockedDataCollection.containsKey(id); + } + + public Data get(String id) { + return lockedDataCollection.get(id); + } + + public void put(String id, Data data) { + lockedDataCollection.put(id, data); + } + + public void hold() { + lockedDataCollection = getCurrentAndHold(); + } + + public void release() { + lockedDataCollection.release(); + lockedDataCollection = null; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCollection.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCollection.java new file mode 100644 index 000000000..725dcca71 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCollection.java @@ -0,0 +1,53 @@ +package org.skywalking.apm.collector.stream.impl.data; + +import java.util.HashMap; +import java.util.Map; + +/** + * @author pengys5 + */ +public class DataCollection { + private Map data; + private volatile boolean isHold; + + public DataCollection() { + this.data = new HashMap<>(); + this.isHold = false; + } + + public void release() { + isHold = false; + } + + public void hold() { + isHold = true; + } + + public boolean isHolding() { + return isHold; + } + + public boolean containsKey(String key) { + return data.containsKey(key); + } + + public void put(String key, Data value) { + data.put(key, value); + } + + public Data get(String key) { + return data.get(key); + } + + public int size() { + return data.size(); + } + + public void clear() { + data.clear(); + } + + public Map asMap() { + return data; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefine.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefine.java new file mode 100644 index 000000000..5571ac264 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefine.java @@ -0,0 +1,78 @@ +package org.skywalking.apm.collector.stream.impl.data; + +/** + * @author pengys5 + */ +public abstract class DataDefine { + private Attribute[] attributes; + private int stringCapacity; + private int longCapacity; + private int floatCapacity; + + public DataDefine() { + stringCapacity = 0; + longCapacity = 0; + floatCapacity = 0; + } + + public final void initial() { + for (Attribute attribute : attributes) { + if (AttributeType.STRING.equals(attribute.getType())) { + stringCapacity++; + } else if (AttributeType.LONG.equals(attribute.getType())) { + longCapacity++; + } else if (AttributeType.FLOAT.equals(attribute.getType())) { + floatCapacity++; + } + } + } + + public final void addAttribute(int position, Attribute attribute) { + attributes[position] = attribute; + } + + public final void define() { + attributes = new Attribute[initialCapacity()]; + } + + protected abstract int defineId(); + + protected abstract int initialCapacity(); + + protected abstract void attributeDefine(); + + public int getStringCapacity() { + return stringCapacity; + } + + public int getLongCapacity() { + return longCapacity; + } + + public int getFloatCapacity() { + return floatCapacity; + } + + public Data build() { + return new Data(defineId(), getStringCapacity(), getLongCapacity(), getFloatCapacity()); + } + + public void mergeData(Data newData, Data oldData) { + int stringPosition = 0; + int longPosition = 0; + int floatPosition = 0; + for (int i = 0; i < initialCapacity(); i++) { + Attribute attribute = attributes[i]; + if (AttributeType.STRING.equals(attribute.getType())) { + attribute.getOperation().operate(newData.getDataString(stringPosition), oldData.getDataString(stringPosition)); + stringPosition++; + } else if (AttributeType.LONG.equals(attribute.getType())) { + attribute.getOperation().operate(newData.getDataLong(longPosition), oldData.getDataLong(longPosition)); + longPosition++; + } else if (AttributeType.FLOAT.equals(attribute.getType())) { + attribute.getOperation().operate(newData.getDataFloat(floatPosition), oldData.getDataFloat(floatPosition)); + floatPosition++; + } + } + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefineLoader.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefineLoader.java new file mode 100644 index 000000000..8ce454440 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefineLoader.java @@ -0,0 +1,29 @@ +package org.skywalking.apm.collector.stream.impl.data; + +import java.util.HashMap; +import java.util.Map; +import org.skywalking.apm.collector.core.config.ConfigException; +import org.skywalking.apm.collector.core.framework.Loader; +import org.skywalking.apm.collector.core.util.DefinitionLoader; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class DataDefineLoader implements Loader> { + + private final Logger logger = LoggerFactory.getLogger(DataDefineLoader.class); + + @Override public Map load() throws ConfigException { + Map dataDefineMap = new HashMap<>(); + + DataDefinitionFile definitionFile = new DataDefinitionFile(); + DefinitionLoader definitionLoader = DefinitionLoader.load(DataDefine.class, definitionFile); + for (DataDefine dataDefine : definitionLoader) { + logger.info("loaded data definition class: {}", dataDefine.getClass().getName()); + dataDefineMap.put(dataDefine.defineId(), dataDefine); + } + return dataDefineMap; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefinitionFile.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefinitionFile.java new file mode 100644 index 000000000..93df98019 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefinitionFile.java @@ -0,0 +1,13 @@ +package org.skywalking.apm.collector.stream.impl.data; + +import org.skywalking.apm.collector.core.framework.DefinitionFile; + +/** + * @author pengys5 + */ +public class DataDefinitionFile extends DefinitionFile { + + @Override protected String fileName() { + return "data.define"; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Operation.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Operation.java new file mode 100644 index 000000000..b94e16b11 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Operation.java @@ -0,0 +1,12 @@ +package org.skywalking.apm.collector.stream.impl.data; + +/** + * @author pengys5 + */ +public interface Operation { + String operate(String newValue, String oldValue); + + Long operate(Long newValue, Long oldValue); + + Float operate(Float newValue, Float oldValue); +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Window.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Window.java new file mode 100644 index 000000000..05b48c197 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Window.java @@ -0,0 +1,44 @@ +package org.skywalking.apm.collector.stream.impl.data; + +/** + * @author pengys5 + */ +public abstract class Window { + + private DataCollection pointer; + + private DataCollection windowDataA; + private DataCollection windowDataB; + + public Window() { + windowDataA = new DataCollection(); + windowDataB = new DataCollection(); + pointer = windowDataA; + } + + public void switchPointer() { + if (pointer == windowDataA) { + pointer = windowDataB; + } else { + pointer = windowDataA; + } + } + + protected DataCollection getCurrentAndHold() { + if (pointer == windowDataA) { + windowDataA.hold(); + return windowDataA; + } else { + windowDataB.hold(); + return windowDataB; + } + } + + public DataCollection getLast() { + if (pointer == windowDataA) { + return windowDataB; + } else { + return windowDataA; + } + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/CoverOperation.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/CoverOperation.java new file mode 100644 index 000000000..bae2f6308 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/CoverOperation.java @@ -0,0 +1,20 @@ +package org.skywalking.apm.collector.stream.impl.data.operate; + +import org.skywalking.apm.collector.stream.impl.data.Operation; + +/** + * @author pengys5 + */ +public class CoverOperation implements Operation { + @Override public String operate(String newValue, String oldValue) { + return newValue; + } + + @Override public Long operate(Long newValue, Long oldValue) { + return newValue; + } + + @Override public Float operate(Float newValue, Float oldValue) { + return newValue; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/NonOperation.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/NonOperation.java new file mode 100644 index 000000000..b33950073 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/operate/NonOperation.java @@ -0,0 +1,20 @@ +package org.skywalking.apm.collector.stream.impl.data.operate; + +import org.skywalking.apm.collector.stream.impl.data.Operation; + +/** + * @author pengys5 + */ +public class NonOperation implements Operation { + @Override public String operate(String newValue, String oldValue) { + return oldValue; + } + + @Override public Long operate(Long newValue, Long oldValue) { + return oldValue; + } + + @Override public Float operate(Float newValue, Float oldValue) { + return oldValue; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/resources/META-INF/defines/remote_handler.define b/apm-collector/apm-collector-stream/src/main/resources/META-INF/defines/remote_handler.define new file mode 100644 index 000000000..389dfb80e --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/resources/META-INF/defines/remote_handler.define @@ -0,0 +1 @@ +org.skywalking.apm.collector.stream.impl.RemoteCommonServiceHandler \ No newline at end of file