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 extends Module> 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