data;
private volatile boolean writing;
private volatile boolean reading;
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractData.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractData.java
new file mode 100644
index 000000000..b8636a6e6
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractData.java
@@ -0,0 +1,196 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.skywalking.apm.collector.core.data;
+
+/**
+ * @author peng-yongsheng
+ */
+public abstract class AbstractData {
+ private String[] dataStrings;
+ private Long[] dataLongs;
+ private Double[] dataDoubles;
+ private Integer[] dataIntegers;
+ private Boolean[] dataBooleans;
+ private byte[][] dataBytes;
+ private final Column[] stringColumns;
+ private final Column[] longColumns;
+ private final Column[] doubleColumns;
+ private final Column[] integerColumns;
+ private final Column[] booleanColumns;
+ private final Column[] byteColumns;
+
+ public AbstractData(Column[] stringColumns, Column[] longColumns, Column[] doubleColumns,
+ Column[] integerColumns, Column[] booleanColumns, Column[] byteColumns) {
+ this.dataStrings = new String[stringColumns.length];
+ this.dataLongs = new Long[longColumns.length];
+ this.dataDoubles = new Double[doubleColumns.length];
+ this.dataIntegers = new Integer[integerColumns.length];
+ this.dataBooleans = new Boolean[booleanColumns.length];
+ this.dataBytes = new byte[byteColumns.length][];
+ this.stringColumns = stringColumns;
+ this.longColumns = longColumns;
+ this.doubleColumns = doubleColumns;
+ this.integerColumns = integerColumns;
+ this.booleanColumns = booleanColumns;
+ this.byteColumns = byteColumns;
+ }
+
+ public int getDataStringsCount() {
+ return dataStrings.length;
+ }
+
+ public int getDataLongsCount() {
+ return dataLongs.length;
+ }
+
+ public int getDataDoublesCount() {
+ return dataDoubles.length;
+ }
+
+ public int getDataIntegersCount() {
+ return dataIntegers.length;
+ }
+
+ public int getDataBooleansCount() {
+ return dataBooleans.length;
+ }
+
+ public int getDataBytesCount() {
+ return dataBytes.length;
+ }
+
+ public void setDataString(int position, String value) {
+ dataStrings[position] = value;
+ }
+
+ public void setDataLong(int position, Long value) {
+ dataLongs[position] = value;
+ }
+
+ public void setDataDouble(int position, Double value) {
+ dataDoubles[position] = value;
+ }
+
+ public void setDataInteger(int position, Integer value) {
+ dataIntegers[position] = value;
+ }
+
+ public void setDataBoolean(int position, Boolean value) {
+ dataBooleans[position] = value;
+ }
+
+ public void setDataBytes(int position, byte[] dataBytes) {
+ this.dataBytes[position] = dataBytes;
+ }
+
+ public String getDataString(int position) {
+ return dataStrings[position];
+ }
+
+ public Long getDataLong(int position) {
+ if (position + 1 > dataLongs.length) {
+ throw new IndexOutOfBoundsException();
+ } else if (dataLongs[position] == null) {
+ return 0L;
+ } else {
+ return dataLongs[position];
+ }
+ }
+
+ public Double getDataDouble(int position) {
+ if (position + 1 > dataDoubles.length) {
+ throw new IndexOutOfBoundsException();
+ } else if (dataDoubles[position] == null) {
+ return 0D;
+ } else {
+ return dataDoubles[position];
+ }
+ }
+
+ public Integer getDataInteger(int position) {
+ if (position + 1 > dataIntegers.length) {
+ throw new IndexOutOfBoundsException();
+ } else if (dataIntegers[position] == null) {
+ return 0;
+ } else {
+ return dataIntegers[position];
+ }
+ }
+
+ public Boolean getDataBoolean(int position) {
+ return dataBooleans[position];
+ }
+
+ public byte[] getDataBytes(int position) {
+ return dataBytes[position];
+ }
+
+ public void mergeData(AbstractData newData) {
+ for (int i = 0; i < stringColumns.length; i++) {
+ String stringData = stringColumns[i].getOperation().operate(newData.getDataString(i), this.getDataString(i));
+ this.dataStrings[i] = stringData;
+ }
+ for (int i = 0; i < longColumns.length; i++) {
+ Long longData = longColumns[i].getOperation().operate(newData.getDataLong(i), this.getDataLong(i));
+ this.dataLongs[i] = longData;
+ }
+ for (int i = 0; i < doubleColumns.length; i++) {
+ Double doubleData = doubleColumns[i].getOperation().operate(newData.getDataDouble(i), this.getDataDouble(i));
+ this.dataDoubles[i] = doubleData;
+ }
+ for (int i = 0; i < integerColumns.length; i++) {
+ Integer integerData = integerColumns[i].getOperation().operate(newData.getDataInteger(i), this.getDataInteger(i));
+ this.dataIntegers[i] = integerData;
+ }
+ for (int i = 0; i < booleanColumns.length; i++) {
+ Boolean booleanData = booleanColumns[i].getOperation().operate(newData.getDataBoolean(i), this.getDataBoolean(i));
+ this.dataBooleans[i] = booleanData;
+ }
+ for (int i = 0; i < byteColumns.length; i++) {
+ byte[] byteData = byteColumns[i].getOperation().operate(newData.getDataBytes(i), this.getDataBytes(i));
+ this.dataBytes[i] = byteData;
+ }
+ }
+
+ @Override public String toString() {
+ StringBuilder dataStr = new StringBuilder();
+ dataStr.append("string: [");
+ for (String dataString : dataStrings) {
+ dataStr.append(dataString).append(",");
+ }
+ dataStr.append("], longs: [");
+ for (Long dataLong : dataLongs) {
+ dataStr.append(dataLong).append(",");
+ }
+ dataStr.append("], double: [");
+ for (Double dataDouble : dataDoubles) {
+ dataStr.append(dataDouble).append(",");
+ }
+ dataStr.append("], integer: [");
+ for (Integer dataInteger : dataIntegers) {
+ dataStr.append(dataInteger).append(",");
+ }
+ dataStr.append("], boolean: [");
+ for (Boolean dataBoolean : dataBooleans) {
+ dataStr.append(dataBoolean).append(",");
+ }
+ dataStr.append("]");
+ return dataStr.toString();
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/Data.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/Data.java
index 9e190797c..f93b68f07 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/Data.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/Data.java
@@ -16,193 +16,45 @@
*
*/
-
package org.apache.skywalking.apm.collector.core.data;
/**
* @author peng-yongsheng
*/
-public abstract class Data extends EndOfBatchQueueMessage {
- private String[] dataStrings;
- private Long[] dataLongs;
- private Double[] dataDoubles;
- private Integer[] dataIntegers;
- private Boolean[] dataBooleans;
- private byte[][] dataBytes;
- private final Column[] stringColumns;
- private final Column[] longColumns;
- private final Column[] doubleColumns;
- private final Column[] integerColumns;
- private final Column[] booleanColumns;
- private final Column[] byteColumns;
+public interface Data {
+ int getDataStringsCount();
- public Data(String id, Column[] stringColumns, Column[] longColumns, Column[] doubleColumns,
- Column[] integerColumns, Column[] booleanColumns, Column[] byteColumns) {
- super(id);
- this.dataStrings = new String[stringColumns.length];
- this.dataStrings[0] = id;
- this.dataLongs = new Long[longColumns.length];
- this.dataDoubles = new Double[doubleColumns.length];
- this.dataIntegers = new Integer[integerColumns.length];
- this.dataBooleans = new Boolean[booleanColumns.length];
- this.dataBytes = new byte[byteColumns.length][];
- this.stringColumns = stringColumns;
- this.longColumns = longColumns;
- this.doubleColumns = doubleColumns;
- this.integerColumns = integerColumns;
- this.booleanColumns = booleanColumns;
- this.byteColumns = byteColumns;
- }
+ int getDataLongsCount();
- public int getDataStringsCount() {
- return dataStrings.length;
- }
+ int getDataDoublesCount();
- public int getDataLongsCount() {
- return dataLongs.length;
- }
+ int getDataIntegersCount();
- public int getDataDoublesCount() {
- return dataDoubles.length;
- }
+ int getDataBooleansCount();
- public int getDataIntegersCount() {
- return dataIntegers.length;
- }
+ int getDataBytesCount();
- public int getDataBooleansCount() {
- return dataBooleans.length;
- }
+ void setDataString(int position, String value);
- public int getDataBytesCount() {
- return dataBytes.length;
- }
+ void setDataLong(int position, Long value);
- public void setDataString(int position, String value) {
- dataStrings[position] = value;
- }
+ void setDataDouble(int position, Double value);
- public void setDataLong(int position, Long value) {
- dataLongs[position] = value;
- }
+ void setDataInteger(int position, Integer value);
- public void setDataDouble(int position, Double value) {
- dataDoubles[position] = value;
- }
+ void setDataBoolean(int position, Boolean value);
- public void setDataInteger(int position, Integer value) {
- dataIntegers[position] = value;
- }
+ void setDataBytes(int position, byte[] dataBytes);
- public void setDataBoolean(int position, Boolean value) {
- dataBooleans[position] = value;
- }
+ String getDataString(int position);
- public void setDataBytes(int position, byte[] dataBytes) {
- this.dataBytes[position] = dataBytes;
- }
+ Long getDataLong(int position);
- public String getDataString(int position) {
- return dataStrings[position];
- }
+ Double getDataDouble(int position);
- public Long getDataLong(int position) {
- if (position + 1 > dataLongs.length) {
- throw new IndexOutOfBoundsException();
- } else if (dataLongs[position] == null) {
- return 0L;
- } else {
- return dataLongs[position];
- }
- }
+ Integer getDataInteger(int position);
- public Double getDataDouble(int position) {
- if (position + 1 > dataDoubles.length) {
- throw new IndexOutOfBoundsException();
- } else if (dataDoubles[position] == null) {
- return 0D;
- } else {
- return dataDoubles[position];
- }
- }
+ Boolean getDataBoolean(int position);
- public Integer getDataInteger(int position) {
- if (position + 1 > dataIntegers.length) {
- throw new IndexOutOfBoundsException();
- } else if (dataIntegers[position] == null) {
- return 0;
- } else {
- return dataIntegers[position];
- }
- }
-
- public Boolean getDataBoolean(int position) {
- return dataBooleans[position];
- }
-
- public byte[] getDataBytes(int position) {
- return dataBytes[position];
- }
-
- public String getId() {
- return dataStrings[0];
- }
-
- public void setId(String id) {
- setKey(id);
- this.dataStrings[0] = id;
- }
-
- public void mergeData(Data newData) {
- for (int i = 0; i < stringColumns.length; i++) {
- String stringData = stringColumns[i].getOperation().operate(newData.getDataString(i), this.getDataString(i));
- this.dataStrings[i] = stringData;
- }
- for (int i = 0; i < longColumns.length; i++) {
- Long longData = longColumns[i].getOperation().operate(newData.getDataLong(i), this.getDataLong(i));
- this.dataLongs[i] = longData;
- }
- for (int i = 0; i < doubleColumns.length; i++) {
- Double doubleData = doubleColumns[i].getOperation().operate(newData.getDataDouble(i), this.getDataDouble(i));
- this.dataDoubles[i] = doubleData;
- }
- for (int i = 0; i < integerColumns.length; i++) {
- Integer integerData = integerColumns[i].getOperation().operate(newData.getDataInteger(i), this.getDataInteger(i));
- this.dataIntegers[i] = integerData;
- }
- for (int i = 0; i < booleanColumns.length; i++) {
- Boolean booleanData = booleanColumns[i].getOperation().operate(newData.getDataBoolean(i), this.getDataBoolean(i));
- this.dataBooleans[i] = booleanData;
- }
- for (int i = 0; i < byteColumns.length; i++) {
- byte[] byteData = byteColumns[i].getOperation().operate(newData.getDataBytes(i), this.getDataBytes(i));
- this.dataBytes[i] = byteData;
- }
- }
-
- @Override public String toString() {
- StringBuilder dataStr = new StringBuilder();
- dataStr.append("string: [");
- for (String dataString : dataStrings) {
- dataStr.append(dataString).append(",");
- }
- dataStr.append("], longs: [");
- for (Long dataLong : dataLongs) {
- dataStr.append(dataLong).append(",");
- }
- dataStr.append("], double: [");
- for (Double dataDouble : dataDoubles) {
- dataStr.append(dataDouble).append(",");
- }
- dataStr.append("], integer: [");
- for (Integer dataInteger : dataIntegers) {
- dataStr.append(dataInteger).append(",");
- }
- dataStr.append("], boolean: [");
- for (Boolean dataBoolean : dataBooleans) {
- dataStr.append(dataBoolean).append(",");
- }
- dataStr.append("]");
- return dataStr.toString();
- }
+ byte[] getDataBytes(int position);
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/EndOfBatchQueueMessage.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/QueueData.java
similarity index 69%
rename from apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/EndOfBatchQueueMessage.java
rename to apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/QueueData.java
index 54319d9b1..9036dfd4a 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/EndOfBatchQueueMessage.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/QueueData.java
@@ -16,26 +16,16 @@
*
*/
-
package org.apache.skywalking.apm.collector.core.data;
+import org.apache.skywalking.apm.collector.core.queue.EndOfBatchContext;
+
/**
* @author peng-yongsheng
*/
-public abstract class EndOfBatchQueueMessage extends AbstractHashMessage {
+public interface QueueData {
- private boolean endOfBatch;
+ EndOfBatchContext getEndOfBatchContext();
- public EndOfBatchQueueMessage(String key) {
- super(key);
- endOfBatch = false;
- }
-
- public final boolean isEndOfBatch() {
- return endOfBatch;
- }
-
- public final void setEndOfBatch(boolean endOfBatch) {
- this.endOfBatch = endOfBatch;
- }
+ void setEndOfBatchContext(EndOfBatchContext context);
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/RemoteData.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/RemoteData.java
new file mode 100644
index 000000000..001397e80
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/RemoteData.java
@@ -0,0 +1,26 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.skywalking.apm.collector.core.data;
+
+/**
+ * @author peng-yongsheng
+ */
+public interface RemoteData extends Data {
+ String selectKey();
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractHashMessage.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/StreamData.java
similarity index 67%
rename from apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractHashMessage.java
rename to apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/StreamData.java
index 51b9714da..87d8dae0b 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractHashMessage.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/StreamData.java
@@ -16,29 +16,22 @@
*
*/
-
package org.apache.skywalking.apm.collector.core.data;
+import org.apache.skywalking.apm.collector.core.queue.EndOfBatchContext;
+
/**
- * The AbstractHashMessage implementations represent aggregate message,
- * which use to aggregate metric.
- *
- *
* @author peng-yongsheng
- * @since v3.0-2017
*/
-public abstract class AbstractHashMessage {
- private int hashCode;
+public abstract class StreamData implements RemoteData, QueueData {
- public AbstractHashMessage(String key) {
- this.hashCode = key.hashCode();
+ private EndOfBatchContext endOfBatchContext;
+
+ @Override public final EndOfBatchContext getEndOfBatchContext() {
+ return this.endOfBatchContext;
}
- public int getHashCode() {
- return hashCode;
- }
-
- public void setKey(String key) {
- this.hashCode = key.hashCode();
+ @Override public final void setEndOfBatchContext(EndOfBatchContext context) {
+ this.endOfBatchContext = endOfBatchContext;
}
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/queue/EndOfBatchContext.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/queue/EndOfBatchContext.java
new file mode 100644
index 000000000..e2507d162
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/queue/EndOfBatchContext.java
@@ -0,0 +1,39 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.skywalking.apm.collector.core.queue;
+
+/**
+ * @author peng-yongsheng
+ */
+public class EndOfBatchContext {
+
+ private boolean isEndOfBatch;
+
+ public EndOfBatchContext(boolean isEndOfBatch) {
+ this.isEndOfBatch = isEndOfBatch;
+ }
+
+ public boolean isEndOfBatch() {
+ return isEndOfBatch;
+ }
+
+ public void setEndOfBatch(boolean endOfBatch) {
+ isEndOfBatch = endOfBatch;
+ }
+}
diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/CommonRemoteDataRegisterService.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/CommonRemoteDataRegisterService.java
index 81c0ee5be..e1df66448 100644
--- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/CommonRemoteDataRegisterService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/CommonRemoteDataRegisterService.java
@@ -16,12 +16,11 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.service;
import java.util.HashMap;
import java.util.Map;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -33,7 +32,7 @@ public class CommonRemoteDataRegisterService implements RemoteDataRegisterServic
private final Logger logger = LoggerFactory.getLogger(CommonRemoteDataRegisterService.class);
private Integer id;
- private final Map, Integer> dataClassMapping;
+ private final Map, Integer> dataClassMapping;
private final Map dataInstanceCreatorMapping;
public CommonRemoteDataRegisterService() {
@@ -42,7 +41,7 @@ public class CommonRemoteDataRegisterService implements RemoteDataRegisterServic
this.dataInstanceCreatorMapping = new HashMap<>();
}
- @Override public void register(Class extends Data> dataClass, RemoteDataInstanceCreator instanceCreator) {
+ @Override public void register(Class extends RemoteData> dataClass, RemoteDataInstanceCreator instanceCreator) {
if (!dataClassMapping.containsKey(dataClass)) {
dataClassMapping.put(dataClass, this.id);
dataInstanceCreatorMapping.put(this.id, instanceCreator);
@@ -53,7 +52,7 @@ public class CommonRemoteDataRegisterService implements RemoteDataRegisterServic
}
@Override
- public Integer getRemoteDataId(Class extends Data> dataClass) throws RemoteDataMappingIdNotFoundException {
+ public Integer getRemoteDataId(Class extends RemoteData> dataClass) throws RemoteDataMappingIdNotFoundException {
if (dataClassMapping.containsKey(dataClass)) {
return dataClassMapping.get(dataClass);
} else {
diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteClient.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteClient.java
index e4d7d5782..153eaba16 100644
--- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteClient.java
+++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteClient.java
@@ -16,10 +16,9 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.service;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
/**
* @author peng-yongsheng
@@ -27,7 +26,7 @@ import org.apache.skywalking.apm.collector.core.data.Data;
public interface RemoteClient extends Comparable {
String getAddress();
- void push(int graphId, int nodeId, Data data);
+ void push(int graphId, int nodeId, RemoteData data);
boolean equals(String address);
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataIDGetter.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataIDGetter.java
index 807d4ec3f..8e993dcf1 100644
--- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataIDGetter.java
+++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataIDGetter.java
@@ -16,14 +16,13 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.service;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
/**
* @author peng-yongsheng
*/
public interface RemoteDataIDGetter {
- Integer getRemoteDataId(Class extends Data> dataClass) throws RemoteDataMappingIdNotFoundException;
+ Integer getRemoteDataId(Class extends RemoteData> dataClass) throws RemoteDataMappingIdNotFoundException;
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataRegisterService.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataRegisterService.java
index 410d5b81d..943b4d31b 100644
--- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataRegisterService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDataRegisterService.java
@@ -16,19 +16,18 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.service;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.core.module.Service;
-import org.apache.skywalking.apm.collector.core.data.Data;
/**
* @author peng-yongsheng
*/
public interface RemoteDataRegisterService extends Service {
- void register(Class extends Data> dataClass, RemoteDataInstanceCreator instanceCreator);
+ void register(Class extends RemoteData> dataClass, RemoteDataInstanceCreator instanceCreator);
- interface RemoteDataInstanceCreator {
- RemoteData createInstance(String id);
+ interface RemoteDataInstanceCreator {
+ REMOTE_DATA createInstance(String id);
}
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDeserializeService.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDeserializeService.java
index 348e85b2e..ef6ad0cf8 100644
--- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDeserializeService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteDeserializeService.java
@@ -16,14 +16,11 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.service;
-import org.apache.skywalking.apm.collector.core.data.Data;
-
/**
* @author peng-yongsheng
*/
public interface RemoteDeserializeService {
- void deserialize(RemoteData remoteData, Data data);
+ void deserialize(RemoteData remoteData, org.apache.skywalking.apm.collector.core.data.RemoteData data);
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSenderService.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSenderService.java
index d14dc61c0..e20fbad16 100644
--- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSenderService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSenderService.java
@@ -16,17 +16,16 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.service;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.core.module.Service;
-import org.apache.skywalking.apm.collector.core.data.Data;
/**
* @author peng-yongsheng
*/
public interface RemoteSenderService extends Service {
- Mode send(int graphId, int nodeId, Data data, Selector selector);
+ Mode send(int graphId, int nodeId, RemoteData remoteData, Selector selector);
enum Mode {
Remote, Local
diff --git a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSerializeService.java b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSerializeService.java
index 908d99a8d..d7b51064e 100644
--- a/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSerializeService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-define/src/main/java/org/apache/skywalking/apm/collector/remote/service/RemoteSerializeService.java
@@ -16,14 +16,13 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.service;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
/**
* @author peng-yongsheng
*/
public interface RemoteSerializeService {
- Builder serialize(Data data);
+ Builder serialize(RemoteData data);
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/handler/RemoteCommonServiceHandler.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/handler/RemoteCommonServiceHandler.java
index 36a037889..c9cc8f53e 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/handler/RemoteCommonServiceHandler.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/handler/RemoteCommonServiceHandler.java
@@ -16,19 +16,17 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.handler;
import io.grpc.stub.StreamObserver;
-import org.apache.skywalking.apm.collector.core.data.Data;
import org.apache.skywalking.apm.collector.core.graph.GraphManager;
import org.apache.skywalking.apm.collector.core.graph.Next;
import org.apache.skywalking.apm.collector.core.util.Const;
-import org.apache.skywalking.apm.collector.remote.grpc.service.GRPCRemoteDeserializeService;
import org.apache.skywalking.apm.collector.remote.grpc.proto.Empty;
import org.apache.skywalking.apm.collector.remote.grpc.proto.RemoteCommonServiceGrpc;
import org.apache.skywalking.apm.collector.remote.grpc.proto.RemoteData;
import org.apache.skywalking.apm.collector.remote.grpc.proto.RemoteMessage;
+import org.apache.skywalking.apm.collector.remote.grpc.service.GRPCRemoteDeserializeService;
import org.apache.skywalking.apm.collector.remote.service.RemoteDataInstanceCreatorGetter;
import org.apache.skywalking.apm.collector.remote.service.RemoteDataInstanceCreatorNotFoundException;
import org.apache.skywalking.apm.collector.server.grpc.GRPCHandler;
@@ -60,7 +58,7 @@ public class RemoteCommonServiceHandler extends RemoteCommonServiceGrpc.RemoteCo
RemoteData remoteData = message.getRemoteData();
try {
- Data output = instanceCreatorGetter.getInstanceCreator(remoteDataId).createInstance(Const.EMPTY_STRING);
+ org.apache.skywalking.apm.collector.core.data.RemoteData output = instanceCreatorGetter.getInstanceCreator(remoteDataId).createInstance(Const.EMPTY_STRING);
service.deserialize(remoteData, output);
Next next = GraphManager.INSTANCE.findGraph(graphId).toFinder().findNext(nodeId);
next.execute(output);
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java
index a120f619e..eda70781d 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClient.java
@@ -16,13 +16,11 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service;
import io.grpc.stub.StreamObserver;
import java.util.List;
import org.apache.skywalking.apm.collector.client.grpc.GRPCClient;
-import org.apache.skywalking.apm.collector.core.data.Data;
import org.apache.skywalking.apm.collector.remote.grpc.proto.Empty;
import org.apache.skywalking.apm.collector.remote.grpc.proto.RemoteCommonServiceGrpc;
import org.apache.skywalking.apm.collector.remote.grpc.proto.RemoteMessage;
@@ -62,7 +60,7 @@ public class GRPCRemoteClient implements RemoteClient {
return this.address;
}
- @Override public void push(int graphId, int nodeId, Data data) {
+ @Override public void push(int graphId, int nodeId, org.apache.skywalking.apm.collector.core.data.RemoteData data) {
try {
Integer remoteDataId = remoteDataIDGetter.getRemoteDataId(data.getClass());
RemoteMessage.Builder builder = RemoteMessage.newBuilder();
@@ -72,7 +70,6 @@ public class GRPCRemoteClient implements RemoteClient {
builder.setRemoteData(service.serialize(data));
this.carrier.produce(builder.build());
- logger.debug("put remote message into queue, id: {}", data.getId());
} catch (RemoteDataMappingIdNotFoundException e) {
logger.error(e.getMessage(), e);
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClientService.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClientService.java
index 46fa0d831..3e36a0e2f 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClientService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteClientService.java
@@ -16,7 +16,6 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service;
import org.apache.skywalking.apm.collector.client.ClientException;
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteDeserializeService.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteDeserializeService.java
index 709c65108..4ef7e5eed 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteDeserializeService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteDeserializeService.java
@@ -16,19 +16,18 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service;
-import org.apache.skywalking.apm.collector.remote.service.RemoteDeserializeService;
-import org.apache.skywalking.apm.collector.core.data.Data;
import org.apache.skywalking.apm.collector.remote.grpc.proto.RemoteData;
+import org.apache.skywalking.apm.collector.remote.service.RemoteDeserializeService;
/**
* @author peng-yongsheng
*/
public class GRPCRemoteDeserializeService implements RemoteDeserializeService {
- @Override public void deserialize(RemoteData remoteData, Data data) {
+ @Override
+ public void deserialize(RemoteData remoteData, org.apache.skywalking.apm.collector.core.data.RemoteData data) {
for (int i = 0; i < remoteData.getDataStringsCount(); i++) {
data.setDataString(i, remoteData.getDataStrings(i));
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSenderService.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSenderService.java
index fe2827765..00bb61e1c 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSenderService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSenderService.java
@@ -16,7 +16,6 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service;
import java.util.ArrayList;
@@ -24,7 +23,7 @@ import java.util.Collections;
import java.util.List;
import org.apache.skywalking.apm.collector.cluster.ClusterModuleListener;
import org.apache.skywalking.apm.collector.core.UnexpectedException;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.remote.RemoteModule;
import org.apache.skywalking.apm.collector.remote.grpc.RemoteModuleGRPCProvider;
import org.apache.skywalking.apm.collector.remote.grpc.service.selector.ForeverFirstSelector;
@@ -50,27 +49,27 @@ public class GRPCRemoteSenderService extends ClusterModuleListener implements Re
private final int channelSize;
private final int bufferSize;
- @Override public Mode send(int graphId, int nodeId, Data data, Selector selector) {
+ @Override public Mode send(int graphId, int nodeId, RemoteData remoteData, Selector selector) {
RemoteClient remoteClient;
switch (selector) {
case HashCode:
- remoteClient = hashCodeSelector.select(remoteClients, data);
- return sendToRemoteWhenNotSelf(remoteClient, graphId, nodeId, data);
+ remoteClient = hashCodeSelector.select(remoteClients, remoteData);
+ return sendToRemoteWhenNotSelf(remoteClient, graphId, nodeId, remoteData);
case Rolling:
- remoteClient = rollingSelector.select(remoteClients, data);
- return sendToRemoteWhenNotSelf(remoteClient, graphId, nodeId, data);
+ remoteClient = rollingSelector.select(remoteClients, remoteData);
+ return sendToRemoteWhenNotSelf(remoteClient, graphId, nodeId, remoteData);
case ForeverFirst:
- remoteClient = foreverFirstSelector.select(remoteClients, data);
- return sendToRemoteWhenNotSelf(remoteClient, graphId, nodeId, data);
+ remoteClient = foreverFirstSelector.select(remoteClients, remoteData);
+ return sendToRemoteWhenNotSelf(remoteClient, graphId, nodeId, remoteData);
}
throw new UnexpectedException("Selector not match, Just support hash, rolling, forever first selector.");
}
- private Mode sendToRemoteWhenNotSelf(RemoteClient remoteClient, int graphId, int nodeId, Data data) {
+ private Mode sendToRemoteWhenNotSelf(RemoteClient remoteClient, int graphId, int nodeId, RemoteData remoteData) {
if (remoteClient.equals(selfAddress)) {
return Mode.Local;
} else {
- remoteClient.push(graphId, nodeId, data);
+ remoteClient.push(graphId, nodeId, remoteData);
return Mode.Remote;
}
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSerializeService.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSerializeService.java
index 762fb9f6b..b4d6a0d49 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSerializeService.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/GRPCRemoteSerializeService.java
@@ -16,10 +16,8 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service;
-import org.apache.skywalking.apm.collector.core.data.Data;
import org.apache.skywalking.apm.collector.remote.grpc.proto.RemoteData;
import org.apache.skywalking.apm.collector.remote.service.RemoteSerializeService;
@@ -28,7 +26,7 @@ import org.apache.skywalking.apm.collector.remote.service.RemoteSerializeService
*/
public class GRPCRemoteSerializeService implements RemoteSerializeService {
- @Override public RemoteData.Builder serialize(Data data) {
+ @Override public RemoteData.Builder serialize(org.apache.skywalking.apm.collector.core.data.RemoteData data) {
RemoteData.Builder builder = RemoteData.newBuilder();
for (int i = 0; i < data.getDataStringsCount(); i++) {
builder.addDataStrings(data.getDataString(i));
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/ForeverFirstSelector.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/ForeverFirstSelector.java
index 0ad2efc61..0820d5ba0 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/ForeverFirstSelector.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/ForeverFirstSelector.java
@@ -16,11 +16,10 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service.selector;
import java.util.List;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.remote.service.RemoteClient;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -32,7 +31,7 @@ public class ForeverFirstSelector implements RemoteClientSelector {
private final Logger logger = LoggerFactory.getLogger(ForeverFirstSelector.class);
- @Override public RemoteClient select(List clients, Data message) {
+ @Override public RemoteClient select(List clients, RemoteData remoteData) {
logger.debug("clients size: {}", clients.size());
return clients.get(0);
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/HashCodeSelector.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/HashCodeSelector.java
index 613f40e71..ffdc3ae21 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/HashCodeSelector.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/HashCodeSelector.java
@@ -16,11 +16,10 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service.selector;
import java.util.List;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.remote.service.RemoteClient;
/**
@@ -28,9 +27,9 @@ import org.apache.skywalking.apm.collector.remote.service.RemoteClient;
*/
public class HashCodeSelector implements RemoteClientSelector {
- @Override public RemoteClient select(List clients, Data message) {
+ @Override public RemoteClient select(List clients, RemoteData remoteData) {
int size = clients.size();
- int selectIndex = Math.abs(message.getHashCode()) % size;
+ int selectIndex = Math.abs(remoteData.selectKey().hashCode()) % size;
return clients.get(selectIndex);
}
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RemoteClientSelector.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RemoteClientSelector.java
index d5d96cec0..51137d1b6 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RemoteClientSelector.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RemoteClientSelector.java
@@ -16,16 +16,15 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service.selector;
import java.util.List;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.remote.service.RemoteClient;
/**
* @author peng-yongsheng
*/
public interface RemoteClientSelector {
- RemoteClient select(List clients, Data message);
+ RemoteClient select(List clients, RemoteData remoteData);
}
diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RollingSelector.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RollingSelector.java
index faa8d7975..3a242c596 100644
--- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RollingSelector.java
+++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/remote/grpc/service/selector/RollingSelector.java
@@ -16,11 +16,10 @@
*
*/
-
package org.apache.skywalking.apm.collector.remote.grpc.service.selector;
import java.util.List;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.remote.service.RemoteClient;
/**
@@ -30,7 +29,7 @@ public class RollingSelector implements RemoteClientSelector {
private int index = 0;
- @Override public RemoteClient select(List clients, Data message) {
+ @Override public RemoteClient select(List clients, RemoteData remoteData) {
int size = clients.size();
index++;
int selectIndex = Math.abs(index) % size;
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java
index 7f6152114..0a476d3de 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/base/dao/IPersistenceDAO.java
@@ -19,12 +19,12 @@
package org.apache.skywalking.apm.collector.storage.base.dao;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
/**
* @author peng-yongsheng
*/
-public interface IPersistenceDAO extends DAO {
+public interface IPersistenceDAO extends DAO {
DataImpl get(String id);
Insert prepareBatchInsert(DataImpl data);
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java
index 525c30dc9..0d82707bc 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ApplicationAlarm extends Data implements Alarm {
+public class ApplicationAlarm extends AbstractData implements Alarm {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationAlarmTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarmList.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarmList.java
index 83f23f52b..2e3fd977e 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarmList.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarmList.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ApplicationAlarmList extends Data {
+public class ApplicationAlarmList extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationAlarmListTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarm.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarm.java
index 566f1aa7d..e23ece51b 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarm.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarm.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ApplicationReferenceAlarm extends Data implements Alarm {
+public class ApplicationReferenceAlarm extends AbstractData implements Alarm {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationReferenceAlarmTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarmList.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarmList.java
index db47b8c90..79038c8fd 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarmList.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationReferenceAlarmList.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ApplicationReferenceAlarmList extends Data {
+public class ApplicationReferenceAlarmList extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationReferenceAlarmListTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarm.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarm.java
index 40eb96b97..72310052f 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarm.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarm.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class InstanceAlarm extends Data implements Alarm {
+public class InstanceAlarm extends AbstractData implements Alarm {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceAlarmTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarmList.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarmList.java
index 3516842e5..ec2aa6021 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarmList.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceAlarmList.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class InstanceAlarmList extends Data {
+public class InstanceAlarmList extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceAlarmListTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarm.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarm.java
index 18c0fb8d0..ff291ff7e 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarm.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarm.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class InstanceReferenceAlarm extends Data implements Alarm {
+public class InstanceReferenceAlarm extends AbstractData implements Alarm {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceReferenceAlarmTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarmList.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarmList.java
index 38b85bbca..bd71e92cc 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarmList.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/InstanceReferenceAlarmList.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class InstanceReferenceAlarmList extends Data {
+public class InstanceReferenceAlarmList extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceReferenceAlarmListTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarm.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarm.java
index 5ccdebb5f..373a6e05e 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarm.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarm.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ServiceAlarm extends Data implements Alarm {
+public class ServiceAlarm extends AbstractData implements Alarm {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceAlarmTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarmList.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarmList.java
index 0976aae0c..1b2c47259 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarmList.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceAlarmList.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ServiceAlarmList extends Data {
+public class ServiceAlarmList extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceAlarmListTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarm.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarm.java
index 67ff09c1a..c1d76b4e7 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarm.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarm.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ServiceReferenceAlarm extends Data implements Alarm {
+public class ServiceReferenceAlarm extends AbstractData implements Alarm {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceReferenceAlarmTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarmList.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarmList.java
index cb173acd3..140a4367b 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarmList.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ServiceReferenceAlarmList.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ServiceReferenceAlarmList extends Data {
+public class ServiceReferenceAlarmList extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceReferenceAlarmListTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationComponent.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationComponent.java
index 34b48445e..23bf9e1e1 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationComponent.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationComponent.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.application;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ApplicationComponent extends Data {
+public class ApplicationComponent extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationComponentTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMapping.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMapping.java
index ec091a131..b466767a2 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMapping.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMapping.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.application;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ApplicationMapping extends Data {
+public class ApplicationMapping extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationMappingTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMetric.java
index 1eed6b07c..80188bf2b 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationMetric.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.application;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.AddOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
import org.apache.skywalking.apm.collector.storage.table.Metric;
@@ -27,7 +27,7 @@ import org.apache.skywalking.apm.collector.storage.table.Metric;
/**
* @author peng-yongsheng
*/
-public class ApplicationMetric extends Data implements Metric {
+public class ApplicationMetric extends AbstractData implements Metric {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationReferenceMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationReferenceMetric.java
index 08222a911..d092ca7a5 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationReferenceMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/application/ApplicationReferenceMetric.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.application;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.AddOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
import org.apache.skywalking.apm.collector.storage.table.Metric;
@@ -27,7 +27,7 @@ import org.apache.skywalking.apm.collector.storage.table.Metric;
/**
* @author peng-yongsheng
*/
-public class ApplicationReferenceMetric extends Data implements Metric {
+public class ApplicationReferenceMetric extends AbstractData implements Metric {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationReferenceMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/global/GlobalTrace.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/global/GlobalTrace.java
index c8ed30ec9..f7f442251 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/global/GlobalTrace.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/global/GlobalTrace.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.global;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class GlobalTrace extends Data {
+public class GlobalTrace extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(GlobalTraceTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMapping.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMapping.java
index 47f9af9aa..c260c5711 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMapping.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMapping.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.instance;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class InstanceMapping extends Data {
+public class InstanceMapping extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceMappingTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMetric.java
index e9e051ed9..9dbf3a355 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceMetric.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.instance;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.AddOperation;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
@@ -28,7 +28,7 @@ import org.apache.skywalking.apm.collector.storage.table.Metric;
/**
* @author peng-yongsheng
*/
-public class InstanceMetric extends Data implements Metric {
+public class InstanceMetric extends AbstractData implements Metric {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceReferenceMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceReferenceMetric.java
index da9a36400..2b6018cf4 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceReferenceMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/instance/InstanceReferenceMetric.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.instance;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.AddOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
import org.apache.skywalking.apm.collector.storage.table.Metric;
@@ -27,7 +27,7 @@ import org.apache.skywalking.apm.collector.storage.table.Metric;
/**
* @author peng-yongsheng
*/
-public class InstanceReferenceMetric extends Data implements Metric {
+public class InstanceReferenceMetric extends AbstractData implements Metric {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceReferenceMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java
index 3cfb4b6d7..b0da0d8a9 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/CpuMetric.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.jvm;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.Column;
import org.apache.skywalking.apm.collector.core.data.operator.AddOperation;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
@@ -28,7 +28,7 @@ import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class CpuMetric extends Data {
+public class CpuMetric extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(CpuMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/GCMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/GCMetric.java
index bad0fe3c7..752b6a5eb 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/GCMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/GCMetric.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.jvm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class GCMetric extends Data {
+public class GCMetric extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(GCMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryMetric.java
index cf342ec8c..1cc3b46f2 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryMetric.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.jvm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class MemoryMetric extends Data {
+public class MemoryMetric extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(MemoryMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetric.java
index a75274432..b0b741a03 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/jvm/MemoryPoolMetric.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.jvm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class MemoryPoolMetric extends Data {
+public class MemoryPoolMetric extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(MemoryPoolMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Application.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Application.java
index 10a4cf12b..47c51b41e 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Application.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Application.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.register;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class Application extends Data {
+public class Application extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Instance.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Instance.java
index c9fee82c1..12ee65d77 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Instance.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/Instance.java
@@ -19,14 +19,14 @@
package org.apache.skywalking.apm.collector.storage.table.register;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class Instance extends Data {
+public class Instance extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(InstanceTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java
index 99aafbf95..ff230e19b 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java
@@ -19,13 +19,13 @@
package org.apache.skywalking.apm.collector.storage.table.register;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class NetworkAddress extends Data {
+public class NetworkAddress extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(NetworkAddressTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServiceName.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServiceName.java
index e95028742..6694fe2c7 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServiceName.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServiceName.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.register;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ServiceName extends Data {
+public class ServiceName extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceNameTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/Segment.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/Segment.java
index a59c5ab9c..c5162a2fe 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/Segment.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/Segment.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.segment;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.Column;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
@@ -27,7 +27,7 @@ import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class Segment extends Data {
+public class Segment extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(SegmentTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/SegmentCost.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/SegmentCost.java
index 51011dfcc..d5c13dbe5 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/SegmentCost.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/segment/SegmentCost.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.segment;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class SegmentCost extends Data {
+public class SegmentCost extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(SegmentCostTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceEntry.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceEntry.java
index b97df5a8b..ad86a4ff6 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceEntry.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceEntry.java
@@ -20,14 +20,14 @@
package org.apache.skywalking.apm.collector.storage.table.service;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ServiceEntry extends Data {
+public class ServiceEntry extends AbstractData {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceEntryTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceMetric.java
index 73c5e6893..bda10fdf3 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceMetric.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.service;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.AddOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
import org.apache.skywalking.apm.collector.storage.table.Metric;
@@ -27,7 +27,7 @@ import org.apache.skywalking.apm.collector.storage.table.Metric;
/**
* @author peng-yongsheng
*/
-public class ServiceMetric extends Data implements Metric {
+public class ServiceMetric extends AbstractData implements Metric {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceMetricTable.COLUMN_ID, new NonOperation()),
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceReferenceMetric.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceReferenceMetric.java
index 574d470a4..abecd56bf 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceReferenceMetric.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/service/ServiceReferenceMetric.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.apm.collector.storage.table.service;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.Data;
+import org.apache.skywalking.apm.collector.core.data.AbstractData;
import org.apache.skywalking.apm.collector.core.data.operator.AddOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
import org.apache.skywalking.apm.collector.storage.table.Metric;
@@ -27,7 +27,7 @@ import org.apache.skywalking.apm.collector.storage.table.Metric;
/**
* @author peng-yongsheng
*/
-public class ServiceReferenceMetric extends Data implements Metric {
+public class ServiceReferenceMetric extends AbstractData implements Metric {
private static final Column[] STRING_COLUMNS = {
new Column(ServiceReferenceMetricTable.COLUMN_ID, new NonOperation()),
From 3f30e6dca251eaab7ba45760f62f9fd8267cd450 Mon Sep 17 00:00:00 2001
From: peng-yongsheng <8082209@qq.com>
Date: Sun, 7 Jan 2018 11:25:33 +0800
Subject: [PATCH 04/43] =?UTF-8?q?Modify=20base=20worker=20modal=E2=80=99s?=
=?UTF-8?q?=20generic=20type=20definition.?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
.../worker/model/base/AbstractLocalAsyncWorker.java | 4 ++--
.../model/base/AbstractLocalAsyncWorkerProvider.java | 4 ++--
.../worker/model/base/AbstractRemoteWorker.java | 4 ++--
.../model/base/AbstractRemoteWorkerProvider.java | 4 ++--
.../analysis/worker/model/base/AbstractWorker.java | 5 ++---
.../worker/model/base/AbstractWorkerProvider.java | 5 ++---
.../worker/model/base/LocalAsyncWorkerRef.java | 10 ++++++----
.../analysis/worker/model/base/RemoteWorkerRef.java | 4 ++--
.../worker/model/base/WorkerCreateListener.java | 2 +-
.../analysis/worker/model/base/WorkerRef.java | 2 +-
10 files changed, 22 insertions(+), 22 deletions(-)
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorker.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorker.java
index 73a48df6c..d85ac9f5d 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorker.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorker.java
@@ -18,7 +18,7 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.base;
-import org.apache.skywalking.apm.collector.core.data.EndOfBatchQueueMessage;
+import org.apache.skywalking.apm.collector.core.data.QueueData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
/**
@@ -28,7 +28,7 @@ import org.apache.skywalking.apm.collector.core.module.ModuleManager;
* @author peng-yongsheng
* @since v3.0-2017
*/
-public abstract class AbstractLocalAsyncWorker extends AbstractWorker {
+public abstract class AbstractLocalAsyncWorker extends AbstractWorker {
public AbstractLocalAsyncWorker(ModuleManager moduleManager) {
super(moduleManager);
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java
index 54d92ef92..472f4d173 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java
@@ -18,14 +18,14 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.base;
-import org.apache.skywalking.apm.collector.core.data.EndOfBatchQueueMessage;
+import org.apache.skywalking.apm.collector.core.data.QueueData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.commons.datacarrier.DataCarrier;
/**
* @author peng-yongsheng
*/
-public abstract class AbstractLocalAsyncWorkerProvider> extends AbstractWorkerProvider {
+public abstract class AbstractLocalAsyncWorkerProvider> extends AbstractWorkerProvider {
public abstract int queueSize();
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorker.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorker.java
index 6f13e83b4..42769c464 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorker.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorker.java
@@ -18,7 +18,7 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.base;
-import org.apache.skywalking.apm.collector.core.data.AbstractData;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.remote.service.Selector;
@@ -31,7 +31,7 @@ import org.apache.skywalking.apm.collector.remote.service.Selector;
* @author peng-yongsheng
* @since v3.0-2017
*/
-public abstract class AbstractRemoteWorker extends AbstractWorker {
+public abstract class AbstractRemoteWorker extends AbstractWorker {
public AbstractRemoteWorker(ModuleManager moduleManager) {
super(moduleManager);
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorkerProvider.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorkerProvider.java
index 339b489f9..e007524d9 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorkerProvider.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractRemoteWorkerProvider.java
@@ -18,7 +18,7 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.base;
-import org.apache.skywalking.apm.collector.core.data.AbstractData;
+import org.apache.skywalking.apm.collector.core.data.RemoteData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.remote.service.RemoteSenderService;
@@ -30,7 +30,7 @@ import org.apache.skywalking.apm.collector.remote.service.RemoteSenderService;
* @author peng-yongsheng
* @since v3.0-2017
*/
-public abstract class AbstractRemoteWorkerProvider> extends AbstractWorkerProvider {
+public abstract class AbstractRemoteWorkerProvider> extends AbstractWorkerProvider {
private final RemoteSenderService remoteSenderService;
private final int graphId;
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorker.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorker.java
index da4efdbbc..e3b697ba2 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorker.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorker.java
@@ -18,7 +18,6 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.base;
-import org.apache.skywalking.apm.collector.core.data.EndOfBatchQueueMessage;
import org.apache.skywalking.apm.collector.core.graph.Next;
import org.apache.skywalking.apm.collector.core.graph.NodeProcessor;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
@@ -28,13 +27,13 @@ import org.slf4j.LoggerFactory;
/**
* @author peng-yongsheng
*/
-public abstract class AbstractWorker implements NodeProcessor {
+public abstract class AbstractWorker implements NodeProcessor {
private final Logger logger = LoggerFactory.getLogger(AbstractWorker.class);
private final ModuleManager moduleManager;
- public AbstractWorker(ModuleManager moduleManager) {
+ AbstractWorker(ModuleManager moduleManager) {
this.moduleManager = moduleManager;
}
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorkerProvider.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorkerProvider.java
index 8ebc62d70..b312af532 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorkerProvider.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractWorkerProvider.java
@@ -18,17 +18,16 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.base;
-import org.apache.skywalking.apm.collector.core.data.EndOfBatchQueueMessage;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
/**
* @author peng-yongsheng
*/
-public abstract class AbstractWorkerProvider> implements Provider {
+public abstract class AbstractWorkerProvider> implements Provider {
private final ModuleManager moduleManager;
- public AbstractWorkerProvider(ModuleManager moduleManager) {
+ AbstractWorkerProvider(ModuleManager moduleManager) {
this.moduleManager = moduleManager;
}
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/LocalAsyncWorkerRef.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/LocalAsyncWorkerRef.java
index 58d831a20..da2e8107b 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/LocalAsyncWorkerRef.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/LocalAsyncWorkerRef.java
@@ -20,8 +20,9 @@ package org.apache.skywalking.apm.collector.analysis.worker.model.base;
import java.util.Iterator;
import java.util.List;
-import org.apache.skywalking.apm.collector.core.data.EndOfBatchQueueMessage;
+import org.apache.skywalking.apm.collector.core.data.QueueData;
import org.apache.skywalking.apm.collector.core.graph.NodeProcessor;
+import org.apache.skywalking.apm.collector.core.queue.EndOfBatchContext;
import org.apache.skywalking.apm.commons.datacarrier.DataCarrier;
import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer;
import org.slf4j.Logger;
@@ -30,7 +31,7 @@ import org.slf4j.LoggerFactory;
/**
* @author peng-yongsheng
*/
-public class LocalAsyncWorkerRef extends WorkerRef implements IConsumer {
+public class LocalAsyncWorkerRef extends WorkerRef implements IConsumer {
private final Logger logger = LoggerFactory.getLogger(LocalAsyncWorkerRef.class);
@@ -40,7 +41,7 @@ public class LocalAsyncWorkerRef dataCarrier) {
+ void setQueueEventHandler(DataCarrier dataCarrier) {
this.dataCarrier = dataCarrier;
}
@@ -52,7 +53,7 @@ public class LocalAsyncWorkerRef extends WorkerRef {
+public class RemoteWorkerRef extends WorkerRef {
private final Logger logger = LoggerFactory.getLogger(RemoteWorkerRef.class);
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerCreateListener.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerCreateListener.java
index 23e602d28..fe1a48f72 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerCreateListener.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerCreateListener.java
@@ -33,7 +33,7 @@ public class WorkerCreateListener {
this.persistenceWorkers = new ArrayList<>();
}
- public void addWorker(AbstractWorker worker) {
+ void addWorker(AbstractWorker worker) {
if (worker instanceof PersistenceWorker) {
persistenceWorkers.add((PersistenceWorker)worker);
}
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerRef.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerRef.java
index b1fa00b5d..3dbb71d27 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerRef.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/WorkerRef.java
@@ -24,7 +24,7 @@ import org.apache.skywalking.apm.collector.core.graph.WayToNode;
/**
* @author peng-yongsheng
*/
-public abstract class WorkerRef extends WayToNode {
+abstract class WorkerRef extends WayToNode {
WorkerRef(NodeProcessor destinationHandler) {
super(destinationHandler);
}
From 0a4171b97fc75f892ba59f8ac1e996149d149315 Mon Sep 17 00:00:00 2001
From: peng-yongsheng <8082209@qq.com>
Date: Sun, 7 Jan 2018 11:26:09 +0800
Subject: [PATCH 05/43] Define stream data.
---
.../apm/collector/core/data/AbstractData.java | 40 +++++++++----------
.../apm/collector/core/data/CommonTable.java | 2 +-
.../apm/collector/core/data/StreamData.java | 21 +++++++++-
.../core/data/AbstractHashMessageTest.java | 40 -------------------
.../storage/table/alarm/ApplicationAlarm.java | 25 ++++++++++--
5 files changed, 61 insertions(+), 67 deletions(-)
delete mode 100644 apm-collector/apm-collector-core/src/test/java/org/apache/skywalking/apm/collector/core/data/AbstractHashMessageTest.java
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractData.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractData.java
index b8636a6e6..51f2a3f08 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractData.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/AbstractData.java
@@ -51,59 +51,59 @@ public abstract class AbstractData {
this.byteColumns = byteColumns;
}
- public int getDataStringsCount() {
+ public final int getDataStringsCount() {
return dataStrings.length;
}
- public int getDataLongsCount() {
+ public final int getDataLongsCount() {
return dataLongs.length;
}
- public int getDataDoublesCount() {
+ public final int getDataDoublesCount() {
return dataDoubles.length;
}
- public int getDataIntegersCount() {
+ public final int getDataIntegersCount() {
return dataIntegers.length;
}
- public int getDataBooleansCount() {
+ public final int getDataBooleansCount() {
return dataBooleans.length;
}
- public int getDataBytesCount() {
+ public final int getDataBytesCount() {
return dataBytes.length;
}
- public void setDataString(int position, String value) {
+ public final void setDataString(int position, String value) {
dataStrings[position] = value;
}
- public void setDataLong(int position, Long value) {
+ public final void setDataLong(int position, Long value) {
dataLongs[position] = value;
}
- public void setDataDouble(int position, Double value) {
+ public final void setDataDouble(int position, Double value) {
dataDoubles[position] = value;
}
- public void setDataInteger(int position, Integer value) {
+ public final void setDataInteger(int position, Integer value) {
dataIntegers[position] = value;
}
- public void setDataBoolean(int position, Boolean value) {
+ public final void setDataBoolean(int position, Boolean value) {
dataBooleans[position] = value;
}
- public void setDataBytes(int position, byte[] dataBytes) {
+ public final void setDataBytes(int position, byte[] dataBytes) {
this.dataBytes[position] = dataBytes;
}
- public String getDataString(int position) {
+ public final String getDataString(int position) {
return dataStrings[position];
}
- public Long getDataLong(int position) {
+ public final Long getDataLong(int position) {
if (position + 1 > dataLongs.length) {
throw new IndexOutOfBoundsException();
} else if (dataLongs[position] == null) {
@@ -113,7 +113,7 @@ public abstract class AbstractData {
}
}
- public Double getDataDouble(int position) {
+ public final Double getDataDouble(int position) {
if (position + 1 > dataDoubles.length) {
throw new IndexOutOfBoundsException();
} else if (dataDoubles[position] == null) {
@@ -123,7 +123,7 @@ public abstract class AbstractData {
}
}
- public Integer getDataInteger(int position) {
+ public final Integer getDataInteger(int position) {
if (position + 1 > dataIntegers.length) {
throw new IndexOutOfBoundsException();
} else if (dataIntegers[position] == null) {
@@ -133,15 +133,15 @@ public abstract class AbstractData {
}
}
- public Boolean getDataBoolean(int position) {
+ public final Boolean getDataBoolean(int position) {
return dataBooleans[position];
}
- public byte[] getDataBytes(int position) {
+ public final byte[] getDataBytes(int position) {
return dataBytes[position];
}
- public void mergeData(AbstractData newData) {
+ public final void mergeData(AbstractData newData) {
for (int i = 0; i < stringColumns.length; i++) {
String stringData = stringColumns[i].getOperation().operate(newData.getDataString(i), this.getDataString(i));
this.dataStrings[i] = stringData;
@@ -168,7 +168,7 @@ public abstract class AbstractData {
}
}
- @Override public String toString() {
+ @Override public final String toString() {
StringBuilder dataStr = new StringBuilder();
dataStr.append("string: [");
for (String dataString : dataStrings) {
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/CommonTable.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/CommonTable.java
index 5110571af..8d8e7e420 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/CommonTable.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/CommonTable.java
@@ -16,7 +16,6 @@
*
*/
-
package org.apache.skywalking.apm.collector.core.data;
/**
@@ -25,5 +24,6 @@ package org.apache.skywalking.apm.collector.core.data;
public abstract class CommonTable {
public static final String TABLE_TYPE = "type";
public static final String COLUMN_ID = "id";
+ public static final String COLUMN_METRIC_ID = "metric_id";
public static final String COLUMN_TIME_BUCKET = "time_bucket";
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/StreamData.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/StreamData.java
index 87d8dae0b..f5fb8eade 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/StreamData.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/data/StreamData.java
@@ -23,7 +23,7 @@ import org.apache.skywalking.apm.collector.core.queue.EndOfBatchContext;
/**
* @author peng-yongsheng
*/
-public abstract class StreamData implements RemoteData, QueueData {
+public abstract class StreamData extends AbstractData implements RemoteData, QueueData {
private EndOfBatchContext endOfBatchContext;
@@ -32,6 +32,23 @@ public abstract class StreamData implements RemoteData, QueueData {
}
@Override public final void setEndOfBatchContext(EndOfBatchContext context) {
- this.endOfBatchContext = endOfBatchContext;
+ this.endOfBatchContext = context;
}
+
+ public StreamData(Column[] stringColumns, Column[] longColumns, Column[] doubleColumns,
+ Column[] integerColumns, Column[] booleanColumns, Column[] byteColumns) {
+ super(stringColumns, longColumns, doubleColumns, integerColumns, booleanColumns, byteColumns);
+ }
+
+ @Override public final String selectKey() {
+ return getMetricId();
+ }
+
+ public abstract String getId();
+
+ public abstract void setId(String id);
+
+ public abstract String getMetricId();
+
+ public abstract void setMetricId(String metricId);
}
diff --git a/apm-collector/apm-collector-core/src/test/java/org/apache/skywalking/apm/collector/core/data/AbstractHashMessageTest.java b/apm-collector/apm-collector-core/src/test/java/org/apache/skywalking/apm/collector/core/data/AbstractHashMessageTest.java
deleted file mode 100644
index 841e12716..000000000
--- a/apm-collector/apm-collector-core/src/test/java/org/apache/skywalking/apm/collector/core/data/AbstractHashMessageTest.java
+++ /dev/null
@@ -1,40 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-
-
-package org.apache.skywalking.apm.collector.core.data;
-
-import org.junit.Assert;
-import org.junit.Test;
-
-/**
- * @author wu-sheng
- */
-public class AbstractHashMessageTest {
- public class NewMessage extends AbstractHashMessage {
- public NewMessage() {
- super("key");
- }
- }
-
- @Test
- public void testHash() {
- NewMessage message = new NewMessage();
- Assert.assertEquals("key".hashCode(), message.getHashCode());
- }
-}
diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java
index 0d82707bc..5a17b8768 100644
--- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java
+++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/alarm/ApplicationAlarm.java
@@ -19,17 +19,18 @@
package org.apache.skywalking.apm.collector.storage.table.alarm;
import org.apache.skywalking.apm.collector.core.data.Column;
-import org.apache.skywalking.apm.collector.core.data.AbstractData;
+import org.apache.skywalking.apm.collector.core.data.StreamData;
import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation;
import org.apache.skywalking.apm.collector.core.data.operator.NonOperation;
/**
* @author peng-yongsheng
*/
-public class ApplicationAlarm extends AbstractData implements Alarm {
+public class ApplicationAlarm extends StreamData implements Alarm {
private static final Column[] STRING_COLUMNS = {
new Column(ApplicationAlarmTable.COLUMN_ID, new NonOperation()),
+ new Column(ApplicationAlarmTable.COLUMN_METRIC_ID, new NonOperation()),
new Column(ApplicationAlarmTable.COLUMN_ALARM_CONTENT, new CoverOperation()),
};
@@ -49,8 +50,24 @@ public class ApplicationAlarm extends AbstractData implements Alarm {
private static final Column[] BYTE_COLUMNS = {};
- public ApplicationAlarm(String id) {
- super(id, STRING_COLUMNS, LONG_COLUMNS, DOUBLE_COLUMNS, INTEGER_COLUMNS, BOOLEAN_COLUMNS, BYTE_COLUMNS);
+ public ApplicationAlarm() {
+ super(STRING_COLUMNS, LONG_COLUMNS, DOUBLE_COLUMNS, INTEGER_COLUMNS, BOOLEAN_COLUMNS, BYTE_COLUMNS);
+ }
+
+ @Override public String getId() {
+ return getDataString(0);
+ }
+
+ @Override public void setId(String id) {
+ setDataString(0, id);
+ }
+
+ @Override public String getMetricId() {
+ return getDataString(1);
+ }
+
+ @Override public void setMetricId(String metricId) {
+ setDataString(1, metricId);
}
@Override
From 506c00b4254e741e27d05fbf82f6574d026cb9fc Mon Sep 17 00:00:00 2001
From: peng-yongsheng <8082209@qq.com>
Date: Sun, 7 Jan 2018 11:35:44 +0800
Subject: [PATCH 06/43] =?UTF-8?q?Change=20the=20aggregation=20and=20persis?=
=?UTF-8?q?tence=20worker=E2=80=99s=20generic=20type=20definition.?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
.../worker/model/impl/AggregationWorker.java | 9 ++--
.../worker/model/impl/FlushAndSwitch.java | 25 -----------
.../worker/model/impl/MessageHolder.java | 41 -------------------
.../worker/model/impl/PersistenceWorker.java | 4 +-
.../model/impl/PersistenceWorkerProvider.java | 4 +-
.../worker/model/impl/data/DataCache.java | 13 +++---
.../model/impl/data/DataCollection.java | 13 +++---
.../apm/collector/core/cache/Window.java | 1 -
8 files changed, 19 insertions(+), 91 deletions(-)
delete mode 100644 apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/FlushAndSwitch.java
delete mode 100644 apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/MessageHolder.java
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/AggregationWorker.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/AggregationWorker.java
index ba8593673..0b6f3240d 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/AggregationWorker.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/AggregationWorker.java
@@ -21,7 +21,7 @@ package org.apache.skywalking.apm.collector.analysis.worker.model.impl;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.impl.data.DataCache;
-import org.apache.skywalking.apm.collector.core.data.AbstractData;
+import org.apache.skywalking.apm.collector.core.data.StreamData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -29,7 +29,7 @@ import org.slf4j.LoggerFactory;
/**
* @author peng-yongsheng
*/
-public abstract class AggregationWorker extends AbstractLocalAsyncWorker {
+public abstract class AggregationWorker extends AbstractLocalAsyncWorker {
private final Logger logger = LoggerFactory.getLogger(AggregationWorker.class);
@@ -47,9 +47,7 @@ public abstract class AggregationWorker {
- private MESSAGE message;
-
- public MESSAGE getMessage() {
- return message;
- }
-
- public void setMessage(MESSAGE message) {
- this.message = message;
- }
-
- public void reset() {
- message = null;
- }
-}
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorker.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorker.java
index 9f104645c..6dfb39b49 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorker.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorker.java
@@ -24,7 +24,7 @@ import java.util.Map;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.impl.data.DataCache;
-import org.apache.skywalking.apm.collector.core.data.AbstractData;
+import org.apache.skywalking.apm.collector.core.data.StreamData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.core.util.ObjectUtils;
import org.apache.skywalking.apm.collector.storage.StorageModule;
@@ -36,7 +36,7 @@ import org.slf4j.LoggerFactory;
/**
* @author peng-yongsheng
*/
-public abstract class PersistenceWorker extends AbstractLocalAsyncWorker {
+public abstract class PersistenceWorker extends AbstractLocalAsyncWorker {
private final Logger logger = LoggerFactory.getLogger(PersistenceWorker.class);
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorkerProvider.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorkerProvider.java
index 6fcdce93a..b9769a871 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorkerProvider.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/PersistenceWorkerProvider.java
@@ -19,13 +19,13 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.impl;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorkerProvider;
-import org.apache.skywalking.apm.collector.core.data.AbstractData;
+import org.apache.skywalking.apm.collector.core.data.StreamData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
/**
* @author peng-yongsheng
*/
-public abstract class PersistenceWorkerProvider> extends AbstractLocalAsyncWorkerProvider {
+public abstract class PersistenceWorkerProvider> extends AbstractLocalAsyncWorkerProvider {
public PersistenceWorkerProvider(ModuleManager moduleManager) {
super(moduleManager);
diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/data/DataCache.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/data/DataCache.java
index 2e7273ec4..7cfa06637 100644
--- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/data/DataCache.java
+++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/impl/data/DataCache.java
@@ -16,20 +16,19 @@
*
*/
-
package org.apache.skywalking.apm.collector.analysis.worker.model.impl.data;
import org.apache.skywalking.apm.collector.core.cache.Window;
-import org.apache.skywalking.apm.collector.core.data.AbstractData;
+import org.apache.skywalking.apm.collector.core.data.StreamData;
/**
* @author peng-yongsheng
*/
-public class DataCache extends Window> {
+public class DataCache extends Window> {
- private DataCollection lockedDataCollection;
+ private DataCollection lockedDataCollection;
- @Override public DataCollection collectionInstance() {
+ @Override public DataCollection collectionInstance() {
return new DataCollection<>();
}
@@ -37,11 +36,11 @@ public class DataCache extends Window implements Collection