From e11e4e5aca8d10cdd61ea048b1fa4a461e9e070f Mon Sep 17 00:00:00 2001 From: peng-yongsheng <8082209@qq.com> Date: Fri, 3 Nov 2017 18:36:48 +0800 Subject: [PATCH] Provide stream module. --- apm-collector/apm-collector-boot/pom.xml | 2 +- .../apm/collector/core/data/Data.java | 4 +- .../apm/collector/core/data/DataDefine.java | 4 +- .../collector/remote/RemoteDataMapping.java | 2 +- .../remote/service/RemoteClient.java | 3 +- .../remote/grpc/service/GRPCRemoteClient.java | 6 +- .../table/global/GlobalTraceDataDefine.java | 5 + .../instance/InstPerformanceDataDefine.java | 5 + .../table/jvm/CpuMetricDataDefine.java | 5 + .../storage/table/jvm/GCMetricDataDefine.java | 5 + .../table/jvm/MemoryMetricDataDefine.java | 5 + .../table/jvm/MemoryPoolMetricDataDefine.java | 5 + .../table/node/NodeComponentDataDefine.java | 5 + .../table/node/NodeMappingDataDefine.java | 5 + .../noderef/NodeReferenceDataDefine.java | 5 + .../table/register/ApplicationDataDefine.java | 7 +- .../table/register/InstanceDataDefine.java | 7 +- .../table/register/ServiceNameDataDefine.java | 7 +- .../table/segment/SegmentCostDataDefine.java | 5 + .../table/segment/SegmentDataDefine.java | 5 + .../table/service/ServiceEntryDataDefine.java | 5 + .../ServiceReferenceDataDefine.java | 5 + .../apm/collector/stream/StreamModule.java | 37 +++++ ...kywalking.apm.collector.core.module.Module | 19 +++ .../collector-stream-provider/pom.xml | 7 + .../stream/StreamModuleProvider.java | 64 ++++++++ .../stream/timer/PersistenceTimer.java | 74 +++++++++ .../base}/AbstractLocalAsyncWorker.java | 4 +- .../AbstractLocalAsyncWorkerProvider.java | 24 +-- .../worker/base}/AbstractRemoteWorker.java | 2 +- .../base}/AbstractRemoteWorkerProvider.java | 21 ++- .../stream/worker/base}/AbstractWorker.java | 2 +- .../worker/base}/AbstractWorkerProvider.java | 2 +- .../worker/base}/ClusterWorkerContext.java | 2 +- .../stream/worker/base}/Context.java | 2 +- .../LocalAsyncWorkerProviderDefineLoader.java | 2 +- .../worker/base}/LocalAsyncWorkerRef.java | 4 +- .../LocalWorkerProviderDefinitionFile.java | 2 +- .../collector/stream/worker/base}/LookUp.java | 2 +- .../stream/worker/base}/Provider.java | 4 +- .../base}/ProviderNotFoundException.java | 2 +- .../RemoteWorkerProviderDefineLoader.java | 2 +- .../RemoteWorkerProviderDefinitionFile.java | 2 +- .../stream/worker/base}/RemoteWorkerRef.java | 35 +--- .../worker/base/RemoteWorkerRefCounter.java} | 4 +- .../collector/stream/worker/base}/Role.java | 4 +- .../worker/base}/UsedRoleNameException.java | 2 +- .../stream/worker/base}/WorkerContext.java | 2 +- .../worker/base}/WorkerCreateListener.java | 9 +- .../stream/worker/base}/WorkerException.java | 2 +- .../worker/base}/WorkerInvokeException.java | 2 +- .../worker/base}/WorkerNotFoundException.java | 2 +- .../stream/worker/base}/WorkerRef.java | 2 +- .../stream/worker/base}/WorkerRefs.java | 4 +- .../base}/selector/ForeverFirstSelector.java | 4 +- .../base}/selector/HashCodeSelector.java | 6 +- .../base}/selector/RollingSelector.java | 6 +- .../worker/base}/selector/WorkerSelector.java | 6 +- .../stream/worker/impl/AggregationWorker.java | 100 ++++++++++++ .../stream/worker/impl/FlushAndSwitch.java | 25 +++ .../stream/worker/impl/PersistenceWorker.java | 154 ++++++++++++++++++ .../impl/PersistenceWorkerContainer.java | 39 +++++ .../stream/worker/impl/data/DataCache.java | 54 ++++++ .../worker/impl/data/DataCollection.java | 86 ++++++++++ .../stream/worker/impl/data/Window.java | 84 ++++++++++ ...g.apm.collector.core.module.ModuleProvider | 19 +++ apm-collector/apm-collector-stream/pom.xml | 15 ++ 67 files changed, 955 insertions(+), 98 deletions(-) create mode 100644 apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/StreamModule.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-define/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.Module create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/StreamModuleProvider.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/timer/PersistenceTimer.java rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/AbstractLocalAsyncWorker.java (95%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/AbstractLocalAsyncWorkerProvider.java (62%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/AbstractRemoteWorker.java (97%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/AbstractRemoteWorkerProvider.java (72%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/AbstractWorker.java (96%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/AbstractWorkerProvider.java (95%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/ClusterWorkerContext.java (95%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/Context.java (94%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/LocalAsyncWorkerProviderDefineLoader.java (97%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/LocalAsyncWorkerRef.java (90%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/LocalWorkerProviderDefinitionFile.java (94%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/LookUp.java (93%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/Provider.java (83%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/ProviderNotFoundException.java (93%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/RemoteWorkerProviderDefineLoader.java (97%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/RemoteWorkerProviderDefinitionFile.java (94%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/RemoteWorkerRef.java (55%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/ClusterWorkerRefCounter.java => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/RemoteWorkerRefCounter.java} (92%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/Role.java (87%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/UsedRoleNameException.java (93%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/WorkerContext.java (98%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/WorkerCreateListener.java (82%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/WorkerException.java (94%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/WorkerInvokeException.java (95%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/WorkerNotFoundException.java (93%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/WorkerRef.java (94%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/WorkerRefs.java (93%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/selector/ForeverFirstSelector.java (89%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/selector/HashCodeSelector.java (91%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/selector/RollingSelector.java (88%) rename apm-collector/apm-collector-stream/{collector-stream-define/src/main/java/org/skywalking/apm/collector/stream => collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base}/selector/WorkerSelector.java (87%) create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java create mode 100644 apm-collector/apm-collector-stream/collector-stream-provider/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.ModuleProvider diff --git a/apm-collector/apm-collector-boot/pom.xml b/apm-collector/apm-collector-boot/pom.xml index 69f7de2dc..c396f3d11 100644 --- a/apm-collector/apm-collector-boot/pom.xml +++ b/apm-collector/apm-collector-boot/pom.xml @@ -113,7 +113,7 @@ org.skywalking - collector-remote-grpc-define + collector-remote-grpc-provider ${project.version} diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/Data.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/Data.java index 280411890..dec0f37e7 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/Data.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/Data.java @@ -29,8 +29,8 @@ public class Data extends AbstractHashMessage { private Boolean[] dataBooleans; private byte[][] dataBytes; - public Data(String id, int stringCapacity, int longCapacity, int doubleCapacity, int integerCapacity, - int booleanCapacity, int byteCapacity) { + public Data(String id, int stringCapacity, int longCapacity, int doubleCapacity, + int integerCapacity, int booleanCapacity, int byteCapacity) { super(id); this.dataStrings = new String[stringCapacity]; this.dataStrings[0] = id; diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/DataDefine.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/DataDefine.java index bbfcdd7bc..35057419d 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/DataDefine.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/data/DataDefine.java @@ -58,6 +58,8 @@ public abstract class DataDefine { attributes[position] = attribute; } + public abstract int remoteDataMappingId(); + protected abstract int initialCapacity(); protected abstract void attributeDefine(); @@ -66,7 +68,7 @@ public abstract class DataDefine { return new Data(id, stringCapacity, longCapacity, doubleCapacity, integerCapacity, booleanCapacity, byteCapacity); } - public void mergeData(Data newData, Data oldData) { + public final void mergeData(Data newData, Data oldData) { int stringPosition = 0; int longPosition = 0; int doublePosition = 0; diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/RemoteDataMapping.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/RemoteDataMapping.java index e18811331..30897c93c 100644 --- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/RemoteDataMapping.java +++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/RemoteDataMapping.java @@ -22,5 +22,5 @@ package org.skywalking.apm.collector.remote; * @author peng-yongsheng */ public enum RemoteDataMapping { - InstPerformance, NodeComponent, NodeMapping, NodeReference, Application, Instance, ServiceName, ServiceEntry, ServiceReference + GlobalTrace, Segment, SegmentCost, InstPerformance, NodeComponent, NodeMapping, NodeReference, Application, Instance, ServiceName, ServiceEntry, ServiceReference, CpuMetric, MemoryMetric, MemoryPoolMetric, GCMetric } diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/service/RemoteClient.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/service/RemoteClient.java index 2bb01d130..23d31d009 100644 --- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/service/RemoteClient.java +++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/skywalking/apm/collector/remote/service/RemoteClient.java @@ -19,11 +19,10 @@ package org.skywalking.apm.collector.remote.service; import org.skywalking.apm.collector.core.data.Data; -import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public interface RemoteClient { - void send(String roleName, Data data, RemoteDataMapping mapping); + void send(String roleName, Data data, int remoteDataMappingId); } diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java index ab196ccde..4a414ad7f 100644 --- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java +++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java @@ -19,11 +19,11 @@ package org.skywalking.apm.collector.remote.grpc.service; import io.grpc.stub.StreamObserver; +import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.remote.RemoteDataMapping; import org.skywalking.apm.collector.remote.RemoteDataMappingContainer; import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; import org.skywalking.apm.collector.remote.grpc.proto.RemoteMessage; -import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.remote.service.RemoteClient; /** @@ -39,8 +39,8 @@ public class GRPCRemoteClient implements RemoteClient { this.streamObserver = streamObserver; } - @Override public void send(String roleName, Data data, RemoteDataMapping mapping) { - RemoteData remoteData = (RemoteData)container.get(mapping.ordinal()).serialize(data); + @Override public void send(String roleName, Data data, int remoteDataMappingId) { + RemoteData remoteData = (RemoteData)container.get(remoteDataMappingId).serialize(data); RemoteMessage.Builder builder = RemoteMessage.newBuilder(); builder.setWorkerRole(roleName); builder.setRemoteData(remoteData); diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/global/GlobalTraceDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/global/GlobalTraceDataDefine.java index be5dcb8d0..dea99cd32 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/global/GlobalTraceDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/global/GlobalTraceDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class GlobalTraceDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.GlobalTrace.ordinal(); + } + @Override protected int initialCapacity() { return 4; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/instance/InstPerformanceDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/instance/InstPerformanceDataDefine.java index eca3fa454..0c92a7d5f 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/instance/InstPerformanceDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/instance/InstPerformanceDataDefine.java @@ -24,12 +24,17 @@ import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.AddOperation; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class InstPerformanceDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.InstPerformance.ordinal(); + } + @Override protected int initialCapacity() { return 6; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/CpuMetricDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/CpuMetricDataDefine.java index 8cdea74ca..4aad02485 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/CpuMetricDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/CpuMetricDataDefine.java @@ -24,12 +24,17 @@ import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.AddOperation; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class CpuMetricDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.CpuMetric.ordinal(); + } + @Override protected int initialCapacity() { return 4; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/GCMetricDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/GCMetricDataDefine.java index 1a284ee52..6c77d9052 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/GCMetricDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/GCMetricDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class GCMetricDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.GCMetric.ordinal(); + } + @Override protected int initialCapacity() { return 6; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryMetricDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryMetricDataDefine.java index 5c31dd132..96dfdd9fc 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryMetricDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryMetricDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class MemoryMetricDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.MemoryMetric.ordinal(); + } + @Override protected int initialCapacity() { return 8; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetricDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetricDataDefine.java index e5aa9b8cf..b1ffce0e9 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetricDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetricDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class MemoryPoolMetricDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.MemoryPoolMetric.ordinal(); + } + @Override protected int initialCapacity() { return 8; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeComponentDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeComponentDataDefine.java index faa3682ac..ea87cf3ad 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeComponentDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeComponentDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class NodeComponentDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.NodeComponent.ordinal(); + } + @Override protected int initialCapacity() { return 6; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeMappingDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeMappingDataDefine.java index 4046112a6..0898f0bb3 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeMappingDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/node/NodeMappingDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class NodeMappingDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.NodeMapping.ordinal(); + } + @Override protected int initialCapacity() { return 5; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/noderef/NodeReferenceDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/noderef/NodeReferenceDataDefine.java index 71d30fb3b..4908e1d10 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/noderef/NodeReferenceDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/noderef/NodeReferenceDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.AddOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class NodeReferenceDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.NodeReference.ordinal(); + } + @Override protected int initialCapacity() { return 11; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ApplicationDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ApplicationDataDefine.java index d0dec2c9b..e8cbb6b2e 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ApplicationDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ApplicationDataDefine.java @@ -18,18 +18,23 @@ package org.skywalking.apm.collector.storage.table.register; -import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.core.data.Attribute; import org.skywalking.apm.collector.core.data.AttributeType; +import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class ApplicationDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.Application.ordinal(); + } + @Override protected int initialCapacity() { return 3; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/InstanceDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/InstanceDataDefine.java index 15d0da209..fdf6bfcb8 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/InstanceDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/InstanceDataDefine.java @@ -18,18 +18,23 @@ package org.skywalking.apm.collector.storage.table.register; -import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.core.data.Attribute; import org.skywalking.apm.collector.core.data.AttributeType; +import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class InstanceDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.Instance.ordinal(); + } + @Override protected int initialCapacity() { return 7; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ServiceNameDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ServiceNameDataDefine.java index 6b56ca657..d37d5cfe7 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ServiceNameDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/register/ServiceNameDataDefine.java @@ -18,18 +18,23 @@ package org.skywalking.apm.collector.storage.table.register; -import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.core.data.Attribute; import org.skywalking.apm.collector.core.data.AttributeType; +import org.skywalking.apm.collector.core.data.Data; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class ServiceNameDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.ServiceName.ordinal(); + } + @Override protected int initialCapacity() { return 4; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentCostDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentCostDataDefine.java index fffb6b85d..4db565c72 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentCostDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentCostDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class SegmentCostDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.SegmentCost.ordinal(); + } + @Override protected int initialCapacity() { return 9; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentDataDefine.java index fb040b4fd..6a602d2df 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/segment/SegmentDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class SegmentDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.Segment.ordinal(); + } + @Override protected int initialCapacity() { return 2; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/service/ServiceEntryDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/service/ServiceEntryDataDefine.java index 84b023b80..3f00bfce0 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/service/ServiceEntryDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/service/ServiceEntryDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.CoverOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class ServiceEntryDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.ServiceEntry.ordinal(); + } + @Override protected int initialCapacity() { return 6; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/serviceref/ServiceReferenceDataDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/serviceref/ServiceReferenceDataDefine.java index b38129f2f..9c79f5f3a 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/serviceref/ServiceReferenceDataDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/skywalking/apm/collector/storage/table/serviceref/ServiceReferenceDataDefine.java @@ -23,12 +23,17 @@ import org.skywalking.apm.collector.core.data.AttributeType; import org.skywalking.apm.collector.core.data.DataDefine; import org.skywalking.apm.collector.core.data.operator.AddOperation; import org.skywalking.apm.collector.core.data.operator.NonOperation; +import org.skywalking.apm.collector.remote.RemoteDataMapping; /** * @author peng-yongsheng */ public class ServiceReferenceDataDefine extends DataDefine { + @Override public int remoteDataMappingId() { + return RemoteDataMapping.ServiceReference.ordinal(); + } + @Override protected int initialCapacity() { return 15; } diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/StreamModule.java b/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/StreamModule.java new file mode 100644 index 000000000..29ae5467c --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/StreamModule.java @@ -0,0 +1,37 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream; + +import org.skywalking.apm.collector.core.module.Module; + +/** + * @author peng-yongsheng + */ +public class StreamModule extends Module { + + public static final String NAME = "stream"; + + @Override public String name() { + return NAME; + } + + @Override public Class[] services() { + return new Class[0]; + } +} diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.Module b/apm-collector/apm-collector-stream/collector-stream-define/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.Module new file mode 100644 index 000000000..468c08d65 --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-define/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.Module @@ -0,0 +1,19 @@ +# +# Copyright 2017, OpenSkywalking Organization All rights reserved. +# +# Licensed 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. +# +# Project repository: https://github.com/OpenSkywalking/skywalking +# + +org.skywalking.apm.collector.stream.StreamModule \ No newline at end of file diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/pom.xml b/apm-collector/apm-collector-stream/collector-stream-provider/pom.xml index 199971bf2..dab054e56 100644 --- a/apm-collector/apm-collector-stream/collector-stream-provider/pom.xml +++ b/apm-collector/apm-collector-stream/collector-stream-provider/pom.xml @@ -30,4 +30,11 @@ collector-stream-provider jar + + + org.skywalking + collector-stream-define + ${project.version} + + \ No newline at end of file diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/StreamModuleProvider.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/StreamModuleProvider.java new file mode 100644 index 000000000..3d67fc526 --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/StreamModuleProvider.java @@ -0,0 +1,64 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream; + +import java.util.Properties; +import org.skywalking.apm.collector.core.module.Module; +import org.skywalking.apm.collector.core.module.ModuleNotFoundException; +import org.skywalking.apm.collector.core.module.ModuleProvider; +import org.skywalking.apm.collector.core.module.ServiceNotProvidedException; +import org.skywalking.apm.collector.queue.QueueModule; +import org.skywalking.apm.collector.queue.service.QueueCreatorService; +import org.skywalking.apm.collector.remote.RemoteModule; +import org.skywalking.apm.collector.remote.service.RemoteClientService; +import org.skywalking.apm.collector.storage.StorageModule; + +/** + * @author peng-yongsheng + */ +public class StreamModuleProvider extends ModuleProvider { + + @Override public String name() { + return "worker"; + } + + @Override public Class module() { + return StreamModule.class; + } + + @Override public void prepare(Properties config) throws ServiceNotProvidedException { + } + + @Override public void start(Properties config) throws ServiceNotProvidedException { + try { + QueueCreatorService queueCreatorService = getManager().find(QueueModule.NAME).getService(QueueCreatorService.class); + RemoteClientService remoteClientService = getManager().find(RemoteModule.NAME).getService(RemoteClientService.class); + } catch (ModuleNotFoundException e) { + throw new ServiceNotProvidedException(e.getMessage()); + } + } + + @Override public void notifyAfterCompleted() throws ServiceNotProvidedException { + + } + + @Override public String[] requiredModules() { + return new String[] {RemoteModule.NAME, QueueModule.NAME, StorageModule.NAME}; + } +} diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/timer/PersistenceTimer.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/timer/PersistenceTimer.java new file mode 100644 index 000000000..83d00431a --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/timer/PersistenceTimer.java @@ -0,0 +1,74 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.timer; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import org.skywalking.apm.collector.core.framework.Starter; +import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.storage.dao.IBatchDAO; +import org.skywalking.apm.collector.stream.worker.WorkerException; +import org.skywalking.apm.collector.stream.worker.impl.FlushAndSwitch; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorker; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorkerContainer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class PersistenceTimer implements Starter { + + private final Logger logger = LoggerFactory.getLogger(PersistenceTimer.class); + + public void start() { + logger.info("persistence timer start"); + //TODO timer value config +// final long timeInterval = EsConfig.Es.Persistence.Timer.VALUE * 1000; + final long timeInterval = 3; + Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(() -> extractDataAndSave(), 1, timeInterval, TimeUnit.SECONDS); + } + + private void extractDataAndSave() { + try { + List workers = PersistenceWorkerContainer.INSTANCE.getPersistenceWorkers(); + List batchAllCollection = new ArrayList<>(); + workers.forEach((PersistenceWorker worker) -> { + logger.debug("extract {} worker data and save", worker.getRole().roleName()); + try { + worker.allocateJob(new FlushAndSwitch()); + List batchCollection = worker.buildBatchCollection(); + logger.debug("extract {} worker data size: {}", worker.getRole().roleName(), batchCollection.size()); + batchAllCollection.addAll(batchCollection); + } catch (WorkerException e) { + logger.error(e.getMessage(), e); + } + }); + + IBatchDAO dao = (IBatchDAO)DAOContainer.INSTANCE.get(IBatchDAO.class.getName()); + dao.batchPersistence(batchAllCollection); + } catch (Throwable e) { + logger.error(e.getMessage(), e); + } finally { + logger.debug("persistence data save finish"); + } + } +} diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractLocalAsyncWorker.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractLocalAsyncWorker.java similarity index 95% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractLocalAsyncWorker.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractLocalAsyncWorker.java index 288d107eb..f08747498 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractLocalAsyncWorker.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractLocalAsyncWorker.java @@ -16,9 +16,9 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; -import org.skywalking.apm.collector.queue.QueueExecutor; +import org.skywalking.apm.collector.queue.base.QueueExecutor; /** * The AbstractLocalAsyncWorker implementations represent workers, diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractLocalAsyncWorkerProvider.java similarity index 62% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractLocalAsyncWorkerProvider.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractLocalAsyncWorkerProvider.java index 484604c78..7192cce36 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractLocalAsyncWorkerProvider.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractLocalAsyncWorkerProvider.java @@ -16,11 +16,11 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; -import org.skywalking.apm.collector.queue.QueueCreator; -import org.skywalking.apm.collector.queue.QueueEventHandler; -import org.skywalking.apm.collector.queue.QueueExecutor; +import org.skywalking.apm.collector.queue.base.QueueEventHandler; +import org.skywalking.apm.collector.queue.base.QueueExecutor; +import org.skywalking.apm.collector.queue.service.QueueCreatorService; /** * @author peng-yongsheng @@ -29,16 +29,20 @@ public abstract class AbstractLocalAsyncWorkerProviderAbstractRemoteWorker implementations represent workers, diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractRemoteWorkerProvider.java similarity index 72% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractRemoteWorkerProvider.java index 921759885..691bccb17 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/AbstractRemoteWorkerProvider.java @@ -16,7 +16,9 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; + +import org.skywalking.apm.collector.remote.service.RemoteClientService; /** * The AbstractRemoteWorkerProvider implementations represent providers, @@ -28,6 +30,12 @@ package org.skywalking.apm.collector.stream; */ public abstract class AbstractRemoteWorkerProvider extends AbstractWorkerProvider { + private final RemoteClientService remoteClientService; + + public AbstractRemoteWorkerProvider(RemoteClientService remoteClientService) { + this.remoteClientService = remoteClientService; + } + /** * Create the worker instance into akka system, the akka system will control the cluster worker life cycle. * @@ -35,15 +43,16 @@ public abstract class AbstractRemoteWorkerProvider streamObserver; private final AbstractRemoteWorker remoteWorker; - private final String address; + private final RemoteClient remoteClient; public RemoteWorkerRef(Role role, AbstractRemoteWorker remoteWorker) { super(role); this.remoteWorker = remoteWorker; this.acrossJVM = false; - this.stub = null; - this.address = Const.EMPTY_STRING; + this.remoteClient = null; } - public RemoteWorkerRef(Role role, GRPCClient client) { + public RemoteWorkerRef(Role role, RemoteClient remoteClient) { super(role); this.remoteWorker = null; this.acrossJVM = true; - this.stub = RemoteCommonServiceGrpc.newStub(client.getChannel()); - this.address = client.toString(); - createStreamObserver(); + this.remoteClient = remoteClient; } @Override public void tell(Object message) throws WorkerInvokeException { if (acrossJVM) { try { - RemoteData remoteData = getRole().dataDefine().serialize(message); - RemoteMessage.Builder builder = RemoteMessage.newBuilder(); - builder.setWorkerRole(getRole().roleName()); - builder.setRemoteData(remoteData); - - streamObserver.onNext(builder.build()); + remoteClient.send(getRole().roleName(), (Data)message, getRole().dataDefine().remoteDataMappingId()); } catch (Throwable e) { logger.error(e.getMessage(), e); } @@ -80,9 +66,6 @@ public class RemoteWorkerRef extends WorkerRef { } @Override public String toString() { - StringBuilder toString = new StringBuilder(); - toString.append("acrossJVM: ").append(acrossJVM); - toString.append(", address: ").append(address); - return toString.toString(); + return "acrossJVM: " + isAcrossJVM(); } } diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/ClusterWorkerRefCounter.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/RemoteWorkerRefCounter.java similarity index 92% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/ClusterWorkerRefCounter.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/RemoteWorkerRefCounter.java index 754ad0365..6d3e0f174 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/ClusterWorkerRefCounter.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/RemoteWorkerRefCounter.java @@ -16,7 +16,7 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -25,7 +25,7 @@ import java.util.concurrent.atomic.AtomicInteger; /** * @author peng-yongsheng */ -public enum ClusterWorkerRefCounter { +public enum RemoteWorkerRefCounter { INSTANCE; private Map counter = new ConcurrentHashMap<>(); diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/Role.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/Role.java similarity index 87% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/Role.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/Role.java index 9cf22125c..0e7f42c78 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/Role.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/Role.java @@ -16,10 +16,10 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; import org.skywalking.apm.collector.core.data.DataDefine; -import org.skywalking.apm.collector.stream.selector.WorkerSelector; +import org.skywalking.apm.collector.stream.worker.base.selector.WorkerSelector; /** * @author peng-yongsheng diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/UsedRoleNameException.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/UsedRoleNameException.java similarity index 93% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/UsedRoleNameException.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/UsedRoleNameException.java index d385b0d07..ce7842f20 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/UsedRoleNameException.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/UsedRoleNameException.java @@ -16,7 +16,7 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; public class UsedRoleNameException extends Exception { public UsedRoleNameException(String message) { diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerContext.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerContext.java similarity index 98% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerContext.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerContext.java index 90e84c772..e2bfe911b 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerContext.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerContext.java @@ -16,7 +16,7 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; import java.util.ArrayList; import java.util.HashMap; diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerCreateListener.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerCreateListener.java similarity index 82% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerCreateListener.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerCreateListener.java index 7b3c08eba..bc8903065 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerCreateListener.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerCreateListener.java @@ -16,11 +16,14 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; /** * @author peng-yongsheng */ -public interface WorkerCreateListener { - void onCreate(W workerx); +public class WorkerCreateListener { + + public void addWorker(AbstractWorker worker) { + + } } diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerException.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerException.java similarity index 94% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerException.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerException.java index 70360ecc7..2b3304702 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerException.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerException.java @@ -16,7 +16,7 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; /** * Defines a general exception a worker can throw when it diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerInvokeException.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerInvokeException.java similarity index 95% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerInvokeException.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerInvokeException.java index ce8fdeed6..5338828c4 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerInvokeException.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerInvokeException.java @@ -16,7 +16,7 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; /** * This exception is raised when worker fails to process job during "call" or "ask" diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerNotFoundException.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerNotFoundException.java similarity index 93% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerNotFoundException.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerNotFoundException.java index fd49a3384..b9e656aff 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerNotFoundException.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerNotFoundException.java @@ -16,7 +16,7 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; public class WorkerNotFoundException extends WorkerException { public WorkerNotFoundException(String message) { diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerRef.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerRef.java similarity index 94% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerRef.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerRef.java index 9237948bf..22440358d 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerRef.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerRef.java @@ -16,7 +16,7 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; /** * @author peng-yongsheng diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerRefs.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerRefs.java similarity index 93% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerRefs.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerRefs.java index ac8ad8480..6c180212a 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/WorkerRefs.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/WorkerRefs.java @@ -16,10 +16,10 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream; +package org.skywalking.apm.collector.stream.worker.base; import java.util.List; -import org.skywalking.apm.collector.stream.selector.WorkerSelector; +import org.skywalking.apm.collector.stream.worker.base.selector.WorkerSelector; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/ForeverFirstSelector.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/ForeverFirstSelector.java similarity index 89% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/ForeverFirstSelector.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/ForeverFirstSelector.java index c404b3b70..2e3405f59 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/ForeverFirstSelector.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/ForeverFirstSelector.java @@ -16,10 +16,10 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream.selector; +package org.skywalking.apm.collector.stream.worker.base.selector; import java.util.List; -import org.skywalking.apm.collector.stream.WorkerRef; +import org.skywalking.apm.collector.stream.worker.base.WorkerRef; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/HashCodeSelector.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/HashCodeSelector.java similarity index 91% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/HashCodeSelector.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/HashCodeSelector.java index 38d1bb1fe..af9bcdbc7 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/HashCodeSelector.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/HashCodeSelector.java @@ -16,12 +16,12 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream.selector; +package org.skywalking.apm.collector.stream.worker.base.selector; import java.util.List; import org.skywalking.apm.collector.core.data.AbstractHashMessage; -import org.skywalking.apm.collector.stream.WorkerRef; -import org.skywalking.apm.collector.stream.AbstractWorker; +import org.skywalking.apm.collector.stream.worker.base.WorkerRef; +import org.skywalking.apm.collector.stream.worker.base.AbstractWorker; /** * The HashCodeSelector is a simple implementation of {@link WorkerSelector}. It choose {@link WorkerRef} diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/RollingSelector.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/RollingSelector.java similarity index 88% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/RollingSelector.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/RollingSelector.java index 1a238ece8..2985a3a04 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/RollingSelector.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/RollingSelector.java @@ -16,11 +16,11 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream.selector; +package org.skywalking.apm.collector.stream.worker.base.selector; import java.util.List; -import org.skywalking.apm.collector.stream.WorkerRef; -import org.skywalking.apm.collector.stream.AbstractWorker; +import org.skywalking.apm.collector.stream.worker.base.WorkerRef; +import org.skywalking.apm.collector.stream.worker.base.AbstractWorker; /** * The RollingSelector is a simple implementation of {@link WorkerSelector}. diff --git a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/WorkerSelector.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/WorkerSelector.java similarity index 87% rename from apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/WorkerSelector.java rename to apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/WorkerSelector.java index 0d25ed8db..9c3cbc928 100644 --- a/apm-collector/apm-collector-stream/collector-stream-define/src/main/java/org/skywalking/apm/collector/stream/selector/WorkerSelector.java +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/base/selector/WorkerSelector.java @@ -16,11 +16,11 @@ * Project repository: https://github.com/OpenSkywalking/skywalking */ -package org.skywalking.apm.collector.stream.selector; +package org.skywalking.apm.collector.stream.worker.base.selector; import java.util.List; -import org.skywalking.apm.collector.stream.WorkerRef; -import org.skywalking.apm.collector.stream.AbstractWorker; +import org.skywalking.apm.collector.stream.worker.base.WorkerRef; +import org.skywalking.apm.collector.stream.worker.base.AbstractWorker; /** * The WorkerSelector should be implemented by any class whose instances diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java new file mode 100644 index 000000000..c5d3f4146 --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java @@ -0,0 +1,100 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.worker.impl; + +import org.skywalking.apm.collector.core.data.Data; +import org.skywalking.apm.collector.queue.base.EndOfBatchCommand; +import org.skywalking.apm.collector.stream.worker.base.AbstractLocalAsyncWorker; +import org.skywalking.apm.collector.stream.worker.base.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.worker.base.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.base.Role; +import org.skywalking.apm.collector.stream.worker.base.WorkerException; +import org.skywalking.apm.collector.stream.worker.base.WorkerInvokeException; +import org.skywalking.apm.collector.stream.worker.base.WorkerNotFoundException; +import org.skywalking.apm.collector.stream.worker.base.WorkerRefs; +import org.skywalking.apm.collector.stream.worker.impl.data.DataCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public abstract class AggregationWorker extends AbstractLocalAsyncWorker { + + private final Logger logger = LoggerFactory.getLogger(AggregationWorker.class); + + private DataCache dataCache; + private int messageNum; + + public AggregationWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + dataCache = new DataCache(); + } + + @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 WorkerRefs nextWorkRef(String id) throws WorkerNotFoundException; + + private void sendToNext() throws WorkerException { + dataCache.switchPointer(); + while (dataCache.getLast().isWriting()) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + throw new WorkerException(e.getMessage(), e); + } + } + dataCache.getLast().asMap().forEach((id, data) -> { + try { + logger.debug(data.toString()); + nextWorkRef(id).tell(data); + } catch (WorkerNotFoundException | WorkerInvokeException e) { + logger.error(e.getMessage(), e); + } + }); + dataCache.finishReadingLast(); + } + + protected final void aggregate(Object message) { + Data data = (Data)message; + dataCache.writing(); + if (dataCache.containsKey(data.id())) { + getRole().dataDefine().mergeData(dataCache.get(data.id()), data); + } else { + dataCache.put(data.id(), data); + } + dataCache.finishWriting(); + } +} diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java new file mode 100644 index 000000000..b6148a732 --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java @@ -0,0 +1,25 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.worker.impl; + +/** + * @author peng-yongsheng + */ +public class FlushAndSwitch { +} diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java new file mode 100644 index 000000000..d99d1ef48 --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java @@ -0,0 +1,154 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.worker.impl; + +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.core.data.Data; +import org.skywalking.apm.collector.core.util.ObjectUtils; +import org.skywalking.apm.collector.queue.base.EndOfBatchCommand; +import org.skywalking.apm.collector.storage.base.dao.DAOContainer; +import org.skywalking.apm.collector.storage.base.dao.IBatchDAO; +import org.skywalking.apm.collector.storage.base.dao.IPersistenceDAO; +import org.skywalking.apm.collector.stream.worker.base.AbstractLocalAsyncWorker; +import org.skywalking.apm.collector.stream.worker.base.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.worker.base.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.base.Role; +import org.skywalking.apm.collector.stream.worker.base.WorkerException; +import org.skywalking.apm.collector.stream.worker.impl.data.DataCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { + + private final Logger logger = LoggerFactory.getLogger(PersistenceWorker.class); + + private DataCache dataCache; + + public PersistenceWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + dataCache = new DataCache(); + } + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected final void onWork(Object message) throws WorkerException { + if (message instanceof FlushAndSwitch) { + try { + if (dataCache.trySwitchPointer()) { + dataCache.switchPointer(); + } + } finally { + dataCache.trySwitchPointerFinally(); + } + } else if (message instanceof EndOfBatchCommand) { + } else { + if (dataCache.currentCollectionSize() >= 5000) { + try { + if (dataCache.trySwitchPointer()) { + dataCache.switchPointer(); + + List collection = buildBatchCollection(); + IBatchDAO dao = (IBatchDAO)DAOContainer.INSTANCE.get(IBatchDAO.class.getName()); + dao.batchPersistence(collection); + } + } finally { + dataCache.trySwitchPointerFinally(); + } + } + aggregate(message); + } + } + + public final List buildBatchCollection() throws WorkerException { + List batchCollection = new LinkedList<>(); + try { + while (dataCache.getLast().isWriting()) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + logger.warn("thread wake up"); + } + } + + if (dataCache.getLast().asMap() != null) { + batchCollection = prepareBatch(dataCache.getLast().asMap()); + } + } finally { + dataCache.finishReadingLast(); + } + return batchCollection; + } + + protected final List prepareBatch(Map dataMap) { + List insertBatchCollection = new LinkedList<>(); + List updateBatchCollection = new LinkedList<>(); + dataMap.forEach((id, data) -> { + if (needMergeDBData()) { + Data dbData = persistenceDAO().get(id, getRole().dataDefine()); + if (ObjectUtils.isNotEmpty(dbData)) { + getRole().dataDefine().mergeData(data, dbData); + try { + updateBatchCollection.add(persistenceDAO().prepareBatchUpdate(data)); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } + } else { + try { + insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } + } + } else { + try { + insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } + } + }); + + insertBatchCollection.addAll(updateBatchCollection); + return insertBatchCollection; + } + + private void aggregate(Object message) { + dataCache.writing(); + Data data = (Data)message; + + if (dataCache.containsKey(data.id())) { + getRole().dataDefine().mergeData(dataCache.get(data.id()), data); + } else { + dataCache.put(data.id(), data); + } + + dataCache.finishWriting(); + } + + protected abstract IPersistenceDAO persistenceDAO(); + + protected abstract boolean needMergeDBData(); +} diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java new file mode 100644 index 000000000..9dd5e4330 --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java @@ -0,0 +1,39 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.worker.impl; + +import java.util.ArrayList; +import java.util.List; + +/** + * @author peng-yongsheng + */ +public enum PersistenceWorkerContainer { + INSTANCE; + + private List persistenceWorkers = new ArrayList<>(); + + public void addWorker(PersistenceWorker worker) { + persistenceWorkers.add(worker); + } + + public List getPersistenceWorkers() { + return persistenceWorkers; + } +} diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java new file mode 100644 index 000000000..b8877b084 --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java @@ -0,0 +1,54 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.worker.impl.data; + +import org.skywalking.apm.collector.core.data.Data; + +/** + * @author peng-yongsheng + */ +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 writing() { + lockedDataCollection = getCurrentAndWriting(); + } + + public int currentCollectionSize() { + return getCurrent().size(); + } + + public void finishWriting() { + lockedDataCollection.finishWriting(); + lockedDataCollection = null; + } +} diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java new file mode 100644 index 000000000..eff0e468c --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java @@ -0,0 +1,86 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.worker.impl.data; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import org.skywalking.apm.collector.core.data.Data; + +/** + * @author peng-yongsheng + */ +public class DataCollection { + private Map data; + private volatile boolean writing; + private volatile boolean reading; + + public DataCollection() { + this.data = new ConcurrentHashMap<>(); + this.writing = false; + this.reading = false; + } + + public void finishWriting() { + writing = false; + } + + public void writing() { + writing = true; + } + + public boolean isWriting() { + return writing; + } + + public void finishReading() { + reading = false; + } + + public void reading() { + reading = true; + } + + public boolean isReading() { + return reading; + } + + 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/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java new file mode 100644 index 000000000..a1b53450a --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java @@ -0,0 +1,84 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.stream.worker.impl.data; + +import java.util.concurrent.atomic.AtomicInteger; + +/** + * @author peng-yongsheng + */ +public abstract class Window { + + private AtomicInteger windowSwitch = new AtomicInteger(0); + + private DataCollection pointer; + + private DataCollection windowDataA; + private DataCollection windowDataB; + + public Window() { + windowDataA = new DataCollection(); + windowDataB = new DataCollection(); + pointer = windowDataA; + } + + public boolean trySwitchPointer() { + return windowSwitch.incrementAndGet() == 1 && !getLast().isReading(); + } + + public void trySwitchPointerFinally() { + windowSwitch.addAndGet(-1); + } + + public void switchPointer() { + if (pointer == windowDataA) { + pointer = windowDataB; + } else { + pointer = windowDataA; + } + getLast().reading(); + } + + protected DataCollection getCurrentAndWriting() { + if (pointer == windowDataA) { + windowDataA.writing(); + return windowDataA; + } else { + windowDataB.writing(); + return windowDataB; + } + } + + protected DataCollection getCurrent() { + return pointer; + } + + public DataCollection getLast() { + if (pointer == windowDataA) { + return windowDataB; + } else { + return windowDataA; + } + } + + public void finishReadingLast() { + getLast().clear(); + getLast().finishReading(); + } +} diff --git a/apm-collector/apm-collector-stream/collector-stream-provider/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.ModuleProvider b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.ModuleProvider new file mode 100644 index 000000000..ee323473e --- /dev/null +++ b/apm-collector/apm-collector-stream/collector-stream-provider/src/main/resources/META-INF/services/org.skywalking.apm.collector.core.module.ModuleProvider @@ -0,0 +1,19 @@ +# +# Copyright 2017, OpenSkywalking Organization All rights reserved. +# +# Licensed 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. +# +# Project repository: https://github.com/OpenSkywalking/skywalking +# + +org.skywalking.apm.collector.stream.StreamModuleProvider \ No newline at end of file diff --git a/apm-collector/apm-collector-stream/pom.xml b/apm-collector/apm-collector-stream/pom.xml index e2debf0f5..f9b2086da 100644 --- a/apm-collector/apm-collector-stream/pom.xml +++ b/apm-collector/apm-collector-stream/pom.xml @@ -40,5 +40,20 @@ apm-collector-core ${project.version} + + org.skywalking + collector-queue-define + ${project.version} + + + org.skywalking + collector-storage-define + ${project.version} + + + org.skywalking + collector-remote-define + ${project.version} + \ No newline at end of file