From d80e2ff43dc3639d79e5af48d3b914ab1170ec6f Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sat, 29 Jul 2017 12:07:42 +0800 Subject: [PATCH] node component save to es success. --- .../AgentStreamModuleInstaller.java | 3 + .../handler/TraceSegmentServiceHandler.java | 5 +- .../agentstream/worker/CommonTable.java | 1 + .../agentstream/worker/TimeSlice.java | 28 ------ ...va => NodeComponentAggregationWorker.java} | 27 +++--- .../NodeComponentPersistenceWorker.java | 70 ++++++++++++++ .../component/NodeComponentRemoteWorker.java | 35 +++++++ .../component/NodeComponentSpanListener.java | 16 +++- .../node/component/dao/INodeComponentDAO.java | 12 +++ .../component/dao/NodeComponentEsDAO.java | 29 ++++++ .../component/dao/NodeComponentH2DAO.java | 15 +++ .../define/NodeComponentDataDefine.java | 13 ++- .../define/NodeComponentEsTableDefine.java | 2 +- .../register/instance/InstanceDataDefine.java | 4 +- .../register/instance/dao/InstanceH2DAO.java | 4 + .../worker/segment/SegmentParse.java | 7 ++ .../worker/storage/PersistenceTimer.java | 55 +++++++++++ .../resources/META-INF/defines/es_dao.define | 3 +- .../resources/META-INF/defines/h2_dao.define | 3 +- .../local_async_worker_provider.define | 3 +- .../defines/remote_worker_provider.define | 4 +- .../TraceSegmentServiceHandlerTestCase.java | 92 +++++++++++++++++++ apm-collector/apm-collector-client/pom.xml | 4 - .../elasticsearch/ElasticSearchClient.java | 5 + .../src/main/resources/logback.xml | 2 +- .../apm/collector/storage/dao/IBatchDAO.java | 10 ++ .../storage/elasticsearch/dao/BatchEsDAO.java | 35 +++++++ .../collector/storage/h2/dao/BatchH2DAO.java | 14 +++ .../resources/META-INF/defines/es_dao.define | 1 + .../resources/META-INF/defines/h2_dao.define | 1 + .../stream/StreamModuleInstaller.java | 1 - .../AbstractLocalAsyncWorkerProvider.java | 6 ++ .../stream/worker/WorkerContext.java | 5 + .../stream/worker/impl/AggregationWorker.java | 32 ++++++- .../stream/worker/impl/FlushAndSwitch.java | 7 ++ .../stream/worker/impl/PersistenceWorker.java | 82 +++++++++++++++++ .../impl/PersistenceWorkerContainer.java | 21 +++++ .../stream/worker/impl/data/Data.java | 6 +- .../stream/worker/impl/data/DataCache.java | 4 + .../stream/worker/impl/data/DataDefine.java | 11 +-- .../worker/impl/data/DataDefineLoader.java | 1 - .../worker/impl/data/TransformToData.java | 8 ++ .../stream/worker/impl/data/Window.java | 18 ++++ apm-collector/pom.xml | 10 ++ 44 files changed, 644 insertions(+), 71 deletions(-) delete mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java rename apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/{NodeComponentAggWorker.java => NodeComponentAggregationWorker.java} (57%) create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentPersistenceWorker.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/INodeComponentDAO.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandlerTestCase.java create mode 100644 apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/IBatchDAO.java create mode 100644 apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java create mode 100644 apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java create mode 100644 apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/es_dao.define create mode 100644 apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/h2_dao.define create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java create mode 100644 apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/TransformToData.java diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java index ace5f61c0..bf5227d48 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java @@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream; import java.util.Iterator; import java.util.Map; +import org.skywalking.apm.collector.agentstream.worker.storage.PersistenceTimer; import org.skywalking.apm.collector.core.client.ClientException; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.framework.DefineException; @@ -32,5 +33,7 @@ public class AgentStreamModuleInstaller implements ModuleInstaller { logger.info("module {} initialize", moduleDefine.getClass().getName()); moduleDefine.initialize((ObjectUtils.isNotEmpty(moduleConfig) && moduleConfig.containsKey(moduleDefine.name())) ? moduleConfig.get(moduleDefine.name()) : null, serverHolder); } + + new PersistenceTimer().start(); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java index 17ab6e578..70887d0c1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java @@ -20,11 +20,11 @@ public class TraceSegmentServiceHandler extends TraceSegmentServiceGrpc.TraceSeg private final Logger logger = LoggerFactory.getLogger(TraceSegmentServiceHandler.class); - private SegmentParse segmentParse = new SegmentParse(); - @Override public StreamObserver collect(StreamObserver responseObserver) { return new StreamObserver() { @Override public void onNext(UpstreamSegment segment) { + logger.debug("receive segment"); + SegmentParse segmentParse = new SegmentParse(); try { List traceIds = segment.getGlobalTraceIdsList(); TraceSegmentObject segmentObject = TraceSegmentObject.parseFrom(segment.getSegment()); @@ -39,6 +39,7 @@ public class TraceSegmentServiceHandler extends TraceSegmentServiceGrpc.TraceSeg } @Override public void onCompleted() { + responseObserver.onNext(Downstream.newBuilder().build()); responseObserver.onCompleted(); } }; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java index 4d5669c16..c66f627ad 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java @@ -4,6 +4,7 @@ package org.skywalking.apm.collector.agentstream.worker; * @author pengys5 */ public class CommonTable { + public static final String TABLE_TYPE = "type"; public static final String COLUMN_ID = "id"; public static final String COLUMN_AGG = "agg"; public static final String COLUMN_TIME_BUCKET = "time_bucket"; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java deleted file mode 100644 index c75132a0e..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/TimeSlice.java +++ /dev/null @@ -1,28 +0,0 @@ -package org.skywalking.apm.collector.agentstream.worker; - -/** - * @author pengys5 - */ -public abstract class TimeSlice { - private String sliceType; - private long startTime; - private long endTime; - - public TimeSlice(String sliceType, long startTime, long endTime) { - this.startTime = startTime; - this.endTime = endTime; - this.sliceType = sliceType; - } - - public String getSliceType() { - return sliceType; - } - - public long getStartTime() { - return startTime; - } - - public long getEndTime() { - return endTime; - } -} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggregationWorker.java similarity index 57% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggWorker.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggregationWorker.java index a3dec152e..85349b10c 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggregationWorker.java @@ -4,17 +4,20 @@ import org.skywalking.apm.collector.agentstream.worker.node.component.define.Nod import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.Role; +import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; +import org.skywalking.apm.collector.stream.worker.WorkerRefs; import org.skywalking.apm.collector.stream.worker.impl.AggregationWorker; import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; -import org.skywalking.apm.collector.stream.worker.selector.RollingSelector; +import org.skywalking.apm.collector.stream.worker.selector.HashCodeSelector; import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; /** * @author pengys5 */ -public class NodeComponentAggWorker extends AggregationWorker { +public class NodeComponentAggregationWorker extends AggregationWorker { - public NodeComponentAggWorker(Role role, ClusterWorkerContext clusterContext) { + public NodeComponentAggregationWorker(Role role, ClusterWorkerContext clusterContext) { super(role, clusterContext); } @@ -22,19 +25,19 @@ public class NodeComponentAggWorker extends AggregationWorker { super.preStart(); } - @Override protected void sendToNext() { - + @Override protected WorkerRefs nextWorkRef(String id) throws WorkerNotFoundException { + return getClusterContext().lookup(NodeComponentRemoteWorker.WorkerRole.INSTANCE); } - public static class Factory extends AbstractLocalAsyncWorkerProvider { + public static class Factory extends AbstractLocalAsyncWorkerProvider { @Override public Role role() { - return Role.INSTANCE; + return WorkerRole.INSTANCE; } @Override - public NodeComponentAggWorker workerInstance(ClusterWorkerContext clusterContext) { - return new NodeComponentAggWorker(role(), clusterContext); + public NodeComponentAggregationWorker workerInstance(ClusterWorkerContext clusterContext) { + return new NodeComponentAggregationWorker(role(), clusterContext); } @Override @@ -43,17 +46,17 @@ public class NodeComponentAggWorker extends AggregationWorker { } } - public enum Role implements org.skywalking.apm.collector.stream.worker.Role { + public enum WorkerRole implements Role { INSTANCE; @Override public String roleName() { - return NodeComponentAggWorker.class.getSimpleName(); + return NodeComponentAggregationWorker.class.getSimpleName(); } @Override public WorkerSelector workerSelector() { - return new RollingSelector(); + return new HashCodeSelector(); } @Override public DataDefine dataDefine() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentPersistenceWorker.java new file mode 100644 index 000000000..2e7ccf0b9 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentPersistenceWorker.java @@ -0,0 +1,70 @@ +package org.skywalking.apm.collector.agentstream.worker.node.component; + +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.agentstream.worker.node.component.dao.INodeComponentDAO; +import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine; +import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; +import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.Role; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorker; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.selector.HashCodeSelector; +import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; + +/** + * @author pengys5 + */ +public class NodeComponentPersistenceWorker extends PersistenceWorker { + + public NodeComponentPersistenceWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected List prepareBatch(Map dataMap) { + INodeComponentDAO dao = (INodeComponentDAO)DAOContainer.INSTANCE.get(INodeComponentDAO.class.getName()); + return dao.prepareBatch(dataMap); + } + + public static class Factory extends AbstractLocalAsyncWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public NodeComponentPersistenceWorker workerInstance(ClusterWorkerContext clusterContext) { + return new NodeComponentPersistenceWorker(role(), clusterContext); + } + + @Override + public int queueSize() { + return 1024; + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return NodeComponentPersistenceWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new NodeComponentDataDefine(); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentRemoteWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentRemoteWorker.java index 41b9096ce..8ed1e4114 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentRemoteWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentRemoteWorker.java @@ -1,10 +1,15 @@ package org.skywalking.apm.collector.agentstream.worker.node.component; +import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorker; +import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorkerProvider; import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; import org.skywalking.apm.collector.stream.worker.Role; import org.skywalking.apm.collector.stream.worker.WorkerException; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.selector.HashCodeSelector; +import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; /** * @author pengys5 @@ -20,6 +25,36 @@ public class NodeComponentRemoteWorker extends AbstractRemoteWorker { } @Override protected void onWork(Object message) throws WorkerException { + getClusterContext().lookup(NodeComponentPersistenceWorker.WorkerRole.INSTANCE).tell(message); + } + public static class Factory extends AbstractRemoteWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public NodeComponentRemoteWorker workerInstance(ClusterWorkerContext clusterContext) { + return new NodeComponentRemoteWorker(role(), clusterContext); + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return NodeComponentRemoteWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new NodeComponentDataDefine(); + } } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java index ee9fca213..615f1bb10 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java @@ -8,6 +8,11 @@ import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.stream.StreamModuleContext; +import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; +import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; +import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; import org.skywalking.apm.network.proto.SpanObject; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -27,13 +32,13 @@ public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanLis if (spanObject.getPeerId() == 0) { peers = String.valueOf(spanObject.getPeerId()); } - String agg = spanObject.getComponent() + Const.ID_SPLIT + peers; + String agg = spanObject.getComponentId() + Const.ID_SPLIT + peers; nodeComponents.add(agg); } @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId) { String peers = String.valueOf(applicationId); - String agg = spanObject.getComponent() + Const.ID_SPLIT + peers; + String agg = spanObject.getComponentId() + Const.ID_SPLIT + peers; nodeComponents.add(agg); } @@ -42,11 +47,18 @@ public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanLis } @Override public void build() { + StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); for (String agg : nodeComponents) { NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); nodeComponent.setId(timeBucket + Const.ID_SPLIT + agg); nodeComponent.setAgg(agg); nodeComponent.setTimeBucket(timeBucket); + try { + logger.debug("send to node component aggregation worker, id: {}", nodeComponent.getId()); + context.getClusterWorkerContext().lookup(NodeComponentAggregationWorker.WorkerRole.INSTANCE).tell(nodeComponent.transform()); + } catch (WorkerInvokeException | WorkerNotFoundException e) { + logger.error(e.getMessage(), e); + } } } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/INodeComponentDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/INodeComponentDAO.java new file mode 100644 index 000000000..ffd9087fd --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/INodeComponentDAO.java @@ -0,0 +1,12 @@ +package org.skywalking.apm.collector.agentstream.worker.node.component.dao; + +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; + +/** + * @author pengys5 + */ +public interface INodeComponentDAO { + List prepareBatch(Map dataMap); +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java new file mode 100644 index 000000000..97840799f --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java @@ -0,0 +1,29 @@ +package org.skywalking.apm.collector.agentstream.worker.node.component.dao; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import org.elasticsearch.action.index.IndexRequestBuilder; +import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; + +/** + * @author pengys5 + */ +public class NodeComponentEsDAO extends EsDAO implements INodeComponentDAO { + + @Override public List prepareBatch(Map dataMap) { + List indexRequestBuilders = new ArrayList<>(); + dataMap.forEach((id, data) -> { + Map source = new HashMap(); + source.put(NodeComponentTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + IndexRequestBuilder builder = getClient().prepareIndex(NodeComponentTable.TABLE, id).setSource(); + indexRequestBuilders.add(builder); + }); + return indexRequestBuilders; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java new file mode 100644 index 000000000..c581fc5ed --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java @@ -0,0 +1,15 @@ +package org.skywalking.apm.collector.agentstream.worker.node.component.dao; + +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.storage.h2.dao.H2DAO; + +/** + * @author pengys5 + */ +public class NodeComponentH2DAO extends H2DAO implements INodeComponentDAO { + + @Override public List prepareBatch(Map map) { + return null; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java index 61941347b..caa691ce9 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java @@ -3,7 +3,9 @@ package org.skywalking.apm.collector.agentstream.worker.node.component.define; import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.impl.data.TransformToData; import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; @@ -34,7 +36,7 @@ public class NodeComponentDataDefine extends DataDefine { return null; } - public static class NodeComponent { + public static class NodeComponent implements TransformToData { private String id; private String agg; private long timeBucket; @@ -48,6 +50,15 @@ public class NodeComponentDataDefine extends DataDefine { public NodeComponent() { } + @Override public Data transform() { + NodeComponentDataDefine define = new NodeComponentDataDefine(); + Data data = define.build(id); + data.setDataString(0, this.id); + data.setDataString(1, this.agg); + data.setDataLong(0, this.timeBucket); + return data; + } + public String getId() { return id; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java index 700393edd..e4e8e2809 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java @@ -13,7 +13,7 @@ public class NodeComponentEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 0; + return 2; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java index 6470d4520..e14afa7b2 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java @@ -19,7 +19,7 @@ public class InstanceDataDefine extends DataDefine { } @Override protected int initialCapacity() { - return 3; + return 6; } @Override protected void attributeDefine() { @@ -28,7 +28,7 @@ public class InstanceDataDefine extends DataDefine { addAttribute(2, new Attribute(InstanceTable.COLUMN_AGENTUUID, AttributeType.STRING, new CoverOperation())); addAttribute(3, new Attribute(InstanceTable.COLUMN_REGISTER_TIME, AttributeType.LONG, new CoverOperation())); addAttribute(4, new Attribute(InstanceTable.COLUMN_INSTANCE_ID, AttributeType.INTEGER, new CoverOperation())); - addAttribute(4, new Attribute(InstanceTable.COLUMN_HEARTBEAT_TIME, AttributeType.LONG, new CoverOperation())); + addAttribute(5, new Attribute(InstanceTable.COLUMN_HEARTBEAT_TIME, AttributeType.LONG, new CoverOperation())); } @Override public Object deserialize(RemoteData remoteData) { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java index d3bda6437..64070090e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java @@ -22,4 +22,8 @@ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { @Override public void save(InstanceDataDefine.Instance instance) { } + + @Override public void updateHeartbeatTime(int instanceId, long heartbeatTime) { + + } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java index 2e6168b8e..7ddebb9d6 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java @@ -63,6 +63,13 @@ public class SegmentParse { } } } + + notifyListenerToBuild(); + } + + private void notifyListenerToBuild() { + spanListeners.forEach(listener -> listener.build()); + refsListeners.forEach(listener -> listener.build()); } private void notifyExitListener(SpanObject spanObject, int applicationId, int applicationInstanceId) { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java new file mode 100644 index 000000000..8e72282dd --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java @@ -0,0 +1,55 @@ +package org.skywalking.apm.collector.agentstream.worker.storage; + +import java.util.List; +import org.skywalking.apm.collector.core.framework.Starter; +import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.storage.dao.IBatchDAO; +import org.skywalking.apm.collector.stream.worker.WorkerException; +import org.skywalking.apm.collector.stream.worker.impl.FlushAndSwitch; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorker; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorkerContainer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class PersistenceTimer implements Starter { + + private final Logger logger = LoggerFactory.getLogger(PersistenceTimer.class); + + public void start() { + logger.info("persistence timer start"); + //TODO timer value config +// final long timeInterval = EsConfig.Es.Persistence.Timer.VALUE * 1000; + final long timeInterval = 3 * 1000; + + Runnable runnable = () -> { + while (true) { + try { + extractDataAndSave(); + Thread.sleep(timeInterval); + } catch (Throwable e) { + logger.error(e.getMessage(), e); + } + } + }; + Thread persistenceThread = new Thread(runnable); + persistenceThread.setName("timerPersistence"); + persistenceThread.start(); + } + + private void extractDataAndSave() { + List workers = PersistenceWorkerContainer.INSTANCE.getPersistenceWorkers(); + workers.forEach(worker -> { + try { + worker.allocateJob(new FlushAndSwitch()); + List batchCollection = worker.buildBatchCollection(); + IBatchDAO dao = (IBatchDAO)DAOContainer.INSTANCE.get(IBatchDAO.class.getName()); + dao.batchPersistence(batchCollection); + } catch (WorkerException e) { + logger.error(e.getMessage(), e); + } + }); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define index ad40a6443..c35382f07 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define @@ -1,3 +1,4 @@ org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationEsDAO org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceEsDAO -org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameEsDAO \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameEsDAO +org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentEsDAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define index 1366f2d92..1fda64fb7 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define @@ -1,3 +1,4 @@ org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationH2DAO org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceH2DAO -org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameH2DAO \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameH2DAO +org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentH2DAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define index 7acdd7e09..5e84c08ca 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define @@ -1,4 +1,5 @@ -org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentAggWorker$Factory +org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentAggregationWorker$Factory +org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentPersistenceWorker$Factory org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterSerialWorker$Factory org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterSerialWorker$Factory org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterSerialWorker$Factory \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define index 2b5f3709e..9fb6fc70c 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define @@ -1,3 +1,5 @@ org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterRemoteWorker$Factory org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterRemoteWorker$Factory -org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterRemoteWorker$Factory \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterRemoteWorker$Factory + +org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentRemoteWorker$Factory \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandlerTestCase.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandlerTestCase.java new file mode 100644 index 000000000..4ecfbea07 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandlerTestCase.java @@ -0,0 +1,92 @@ +package org.skywalking.apm.collector.agentstream.grpc.handler; + +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import io.grpc.stub.StreamObserver; +import org.junit.Test; +import org.skywalking.apm.network.proto.Downstream; +import org.skywalking.apm.network.proto.SpanLayer; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.SpanType; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.skywalking.apm.network.proto.TraceSegmentServiceGrpc; +import org.skywalking.apm.network.proto.UniqueId; +import org.skywalking.apm.network.proto.UpstreamSegment; +import org.skywalking.apm.network.trace.component.ComponentsDefine; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class TraceSegmentServiceHandlerTestCase { + + private final Logger logger = LoggerFactory.getLogger(TraceSegmentServiceHandlerTestCase.class); + + private TraceSegmentServiceGrpc.TraceSegmentServiceStub stub; + + @Test + public void testCollect() { + ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).usePlaintext(true).build(); + stub = TraceSegmentServiceGrpc.newStub(channel); + + StreamObserver streamObserver = stub.collect(new StreamObserver() { + @Override public void onNext(Downstream downstream) { + } + + @Override public void onError(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + + @Override public void onCompleted() { + + } + }); + + UpstreamSegment.Builder builder = UpstreamSegment.newBuilder(); + buildGlobalTraceIds(builder); + buildSegment(builder); + + streamObserver.onNext(builder.build()); + streamObserver.onCompleted(); + + try { + Thread.sleep(30000); + } catch (InterruptedException e) { + } + } + + private void buildGlobalTraceIds(UpstreamSegment.Builder builder) { + UniqueId.Builder builder1 = UniqueId.newBuilder(); + builder1.addIdParts(100); + builder1.addIdParts(100); + builder1.addIdParts(100); + builder.addGlobalTraceIds(builder1.build()); + } + + private void buildSegment(UpstreamSegment.Builder builder) { + long now = System.currentTimeMillis(); + + TraceSegmentObject.Builder segmentBuilder = TraceSegmentObject.newBuilder(); + segmentBuilder.setApplicationId(1); + segmentBuilder.setApplicationInstanceId(1); + segmentBuilder.setTraceSegmentId(UniqueId.newBuilder().addIdParts(200).addIdParts(200).addIdParts(200).build()); + + SpanObject.Builder span_0 = SpanObject.newBuilder(); + span_0.setSpanId(0); + span_0.setOperationName("/dubbox-case/case/dubbox-rest"); + span_0.setOperationNameId(0); + span_0.setParentSpanId(-1); + span_0.setSpanLayer(SpanLayer.Http); + span_0.setStartTime(now); + span_0.setEndTime(now + 100000); + span_0.setComponentId(ComponentsDefine.TOMCAT.getId()); + span_0.setIsError(false); + span_0.setSpanType(SpanType.Entry); + span_0.setPeerId(0); + span_0.setPeer("localhost:8080"); + segmentBuilder.addSpans(span_0); + + builder.setSegment(segmentBuilder.build().toByteString()); + } +} diff --git a/apm-collector/apm-collector-client/pom.xml b/apm-collector/apm-collector-client/pom.xml index 5fe914473..1f6b17431 100644 --- a/apm-collector/apm-collector-client/pom.xml +++ b/apm-collector/apm-collector-client/pom.xml @@ -52,10 +52,6 @@ slf4j-log4j12 org.slf4j - - log4j - log4j - diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java index 27218ce35..847f5e7cf 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java @@ -8,6 +8,7 @@ import java.util.concurrent.ExecutionException; import org.elasticsearch.action.admin.indices.create.CreateIndexResponse; import org.elasticsearch.action.admin.indices.delete.DeleteIndexResponse; import org.elasticsearch.action.admin.indices.exists.indices.IndicesExistsResponse; +import org.elasticsearch.action.bulk.BulkRequestBuilder; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.update.UpdateRequest; @@ -112,6 +113,10 @@ public class ElasticSearchClient implements Client { return client.prepareIndex(indexName, "type", id); } + public BulkRequestBuilder prepareBulk() { + return client.prepareBulk(); + } + public void update(UpdateRequest updateRequest) { try { client.update(updateRequest).get(); diff --git a/apm-collector/apm-collector-core/src/main/resources/logback.xml b/apm-collector/apm-collector-core/src/main/resources/logback.xml index 687182620..78debf076 100644 --- a/apm-collector/apm-collector-core/src/main/resources/logback.xml +++ b/apm-collector/apm-collector-core/src/main/resources/logback.xml @@ -2,7 +2,7 @@ - %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{70} - %msg%n diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/IBatchDAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/IBatchDAO.java new file mode 100644 index 000000000..afcac1bf1 --- /dev/null +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/IBatchDAO.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.storage.dao; + +import java.util.List; + +/** + * @author pengys5 + */ +public interface IBatchDAO { + void batchPersistence(List batchCollection); +} diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java new file mode 100644 index 000000000..74c06e634 --- /dev/null +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java @@ -0,0 +1,35 @@ +package org.skywalking.apm.collector.storage.elasticsearch.dao; + +import java.util.List; +import org.elasticsearch.action.bulk.BulkRequestBuilder; +import org.elasticsearch.action.bulk.BulkResponse; +import org.elasticsearch.action.index.IndexRequestBuilder; +import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.skywalking.apm.collector.storage.dao.IBatchDAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class BatchEsDAO extends EsDAO implements IBatchDAO { + + private final Logger logger = LoggerFactory.getLogger(BatchEsDAO.class); + + @Override public void batchPersistence(List batchCollection) { + BulkRequestBuilder bulkRequest = getClient().prepareBulk(); + + logger.info("bulk data size: {}", batchCollection.size()); + if (CollectionUtils.isNotEmpty(batchCollection)) { + for (int i = 0; i < batchCollection.size(); i++) { + IndexRequestBuilder builder = (IndexRequestBuilder)batchCollection.get(i); + bulkRequest.add(builder); + } + + BulkResponse bulkResponse = bulkRequest.execute().actionGet(); + if (bulkResponse.hasFailures()) { + logger.error(bulkResponse.buildFailureMessage()); + } + } + } +} diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java new file mode 100644 index 000000000..27b69f6d4 --- /dev/null +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java @@ -0,0 +1,14 @@ +package org.skywalking.apm.collector.storage.h2.dao; + +import java.util.List; +import org.skywalking.apm.collector.storage.dao.IBatchDAO; + +/** + * @author pengys5 + */ +public class BatchH2DAO extends H2DAO implements IBatchDAO { + + @Override public void batchPersistence(List batchCollection) { + + } +} diff --git a/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/es_dao.define new file mode 100644 index 000000000..1fcbc9a03 --- /dev/null +++ b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/es_dao.define @@ -0,0 +1 @@ +org.skywalking.apm.collector.storage.elasticsearch.dao.BatchEsDAO \ No newline at end of file diff --git a/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/h2_dao.define new file mode 100644 index 000000000..de143481d --- /dev/null +++ b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/h2_dao.define @@ -0,0 +1 @@ +org.skywalking.apm.collector.storage.h2.dao.BatchH2DAO \ No newline at end of file diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleInstaller.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleInstaller.java index a714b6fdd..908eecbb0 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleInstaller.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleInstaller.java @@ -66,7 +66,6 @@ public class StreamModuleInstaller implements ModuleInstaller { List remoteProviders = remoteProviderLoader.load(); for (AbstractRemoteWorkerProvider provider : remoteProviders) { provider.setClusterContext(clusterWorkerContext); -// provider.create(); clusterWorkerContext.putRole(provider.role()); clusterWorkerContext.putProvider(provider); } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorkerProvider.java index c7e4ff5c2..094b960cb 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorkerProvider.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorkerProvider.java @@ -6,6 +6,8 @@ import org.skywalking.apm.collector.core.queue.QueueEventHandler; import org.skywalking.apm.collector.core.queue.QueueExecutor; import org.skywalking.apm.collector.queue.QueueModuleContext; import org.skywalking.apm.collector.queue.QueueModuleGroupDefine; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorker; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorkerContainer; /** * @author pengys5 @@ -18,6 +20,10 @@ public abstract class AbstractLocalAsyncWorkerProvider remoteWorkerRefs; private Map> roleWorkers; private Map roles; @@ -56,6 +60,7 @@ public abstract class WorkerContext implements Context { } @Override final public void put(WorkerRef workerRef) { + logger.debug("put worker reference into context, role name: {}", workerRef.getRole().roleName()); if (!getRoleWorkers().containsKey(workerRef.getRole().roleName())) { getRoleWorkers().putIfAbsent(workerRef.getRole().roleName(), new ArrayList<>()); } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java index 222caa4b6..6d9dff63f 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java @@ -6,23 +6,29 @@ import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; import org.skywalking.apm.collector.stream.worker.Role; import org.skywalking.apm.collector.stream.worker.WorkerException; +import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; +import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; +import org.skywalking.apm.collector.stream.worker.WorkerRefs; import org.skywalking.apm.collector.stream.worker.impl.data.Data; import org.skywalking.apm.collector.stream.worker.impl.data.DataCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public abstract class AggregationWorker extends AbstractLocalAsyncWorker { + private final Logger logger = LoggerFactory.getLogger(AggregationWorker.class); + private DataCache dataCache; + private int messageNum; public AggregationWorker(Role role, ClusterWorkerContext clusterContext) { super(role, clusterContext); dataCache = new DataCache(); } - private int messageNum; - @Override public void preStart() throws ProviderNotFoundException { super.preStart(); } @@ -41,14 +47,34 @@ public abstract class AggregationWorker extends AbstractLocalAsyncWorker { } } - protected abstract void sendToNext(); + protected abstract WorkerRefs nextWorkRef(String id) throws WorkerNotFoundException; + + private void sendToNext() throws WorkerException { + dataCache.switchPointer(); + while (dataCache.getLast().isHolding()) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + throw new WorkerException(e.getMessage(), e); + } + } + dataCache.getLast().asMap().forEach((id, data) -> { + try { + nextWorkRef(id).tell(data); + } catch (WorkerNotFoundException | WorkerInvokeException e) { + logger.error(e.getMessage(), e); + } + }); + } protected final void aggregate(Object message) { Data data = (Data)message; + dataCache.hold(); if (dataCache.containsKey(data.id())) { getClusterContext().getDataDefine(data.getDefineId()).mergeData(data, dataCache.get(data.id())); } else { dataCache.put(data.id(), data); } + dataCache.release(); } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java new file mode 100644 index 000000000..15b8ee037 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/FlushAndSwitch.java @@ -0,0 +1,7 @@ +package org.skywalking.apm.collector.stream.worker.impl; + +/** + * @author pengys5 + */ +public class FlushAndSwitch { +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java new file mode 100644 index 000000000..c32e676de --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorker.java @@ -0,0 +1,82 @@ +package org.skywalking.apm.collector.stream.worker.impl; + +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.core.queue.EndOfBatchCommand; +import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorker; +import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.Role; +import org.skywalking.apm.collector.stream.worker.WorkerException; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { + + private final Logger logger = LoggerFactory.getLogger(PersistenceWorker.class); + + private DataCache dataCache; + + public PersistenceWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + dataCache = new DataCache(); + } + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected final void onWork(Object message) throws WorkerException { + if (message instanceof EndOfBatchCommand || message instanceof FlushAndSwitch) { + if (dataCache.trySwitchPointer()) { + dataCache.switchPointer(); + } + } else { + if (dataCache.currentCollectionSize() >= 1000) { + if (dataCache.trySwitchPointer()) { + dataCache.switchPointer(); + } + } + aggregate(message); + } + } + + public List buildBatchCollection() throws WorkerException { + List batchCollection; + try { + while (dataCache.getLast().isHolding()) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + logger.warn("thread wake up"); + } + } + batchCollection = prepareBatch(dataCache.getLast().asMap()); + } finally { + dataCache.releaseLast(); + } + return batchCollection; + } + + protected abstract List prepareBatch(Map dataMap); + + private void aggregate(Object message) { + dataCache.hold(); + Data data = (Data)message; + + if (dataCache.containsKey(data.id())) { + getClusterContext().getDataDefine(data.getDefineId()).mergeData(data, dataCache.get(data.id())); + } else { + if (dataCache.currentCollectionSize() < 1000) { + dataCache.put(data.id(), data); + } + } + + dataCache.release(); + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java new file mode 100644 index 000000000..1a71b1153 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/PersistenceWorkerContainer.java @@ -0,0 +1,21 @@ +package org.skywalking.apm.collector.stream.worker.impl; + +import java.util.ArrayList; +import java.util.List; + +/** + * @author pengys5 + */ +public enum PersistenceWorkerContainer { + INSTANCE; + + private List persistenceWorkers = new ArrayList<>(); + + public void addWorker(PersistenceWorker worker) { + persistenceWorkers.add(worker); + } + + public List getPersistenceWorkers() { + return persistenceWorkers; + } +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java index 28a2f3074..42ac82c28 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java @@ -1,11 +1,12 @@ package org.skywalking.apm.collector.stream.worker.impl.data; import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.stream.worker.selector.AbstractHashMessage; /** * @author pengys5 */ -public class Data { +public class Data extends AbstractHashMessage { private int defineId; private final int stringCapacity; private final int longCapacity; @@ -16,7 +17,8 @@ public class Data { private Float[] dataFloats; private Integer[] dataIntegers; - public Data(int defineId, int stringCapacity, int longCapacity, int floatCapacity, int integerCapacity) { + public Data(String id, int defineId, int stringCapacity, int longCapacity, int floatCapacity, int integerCapacity) { + super(id); this.defineId = defineId; this.dataStrings = new String[stringCapacity]; this.dataLongs = new Long[longCapacity]; diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java index 8db3ba70f..c22189fb5 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java @@ -23,6 +23,10 @@ public class DataCache extends Window { lockedDataCollection = getCurrentAndHold(); } + public int currentCollectionSize() { + return getCurrentAndHold().size(); + } + public void release() { lockedDataCollection.release(); lockedDataCollection = null; diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java index ebf571721..0c9174d6b 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java @@ -13,13 +13,10 @@ public abstract class DataDefine { private int integerCapacity; public DataDefine() { - stringCapacity = 0; - longCapacity = 0; - floatCapacity = 0; - integerCapacity = 0; + initial(); } - public final void initial() { + private void initial() { attributes = new Attribute[initialCapacity()]; attributeDefine(); for (Attribute attribute : attributes) { @@ -45,8 +42,8 @@ public abstract class DataDefine { protected abstract void attributeDefine(); - public final Data build() { - return new Data(defineId(), stringCapacity, longCapacity, floatCapacity, integerCapacity); + public final Data build(String id) { + return new Data(id, defineId(), stringCapacity, longCapacity, floatCapacity, integerCapacity); } public void mergeData(Data newData, Data oldData) { diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefineLoader.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefineLoader.java index cf39e6dd9..425fcb427 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefineLoader.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefineLoader.java @@ -22,7 +22,6 @@ public class DataDefineLoader implements Loader> { DefinitionLoader definitionLoader = DefinitionLoader.load(DataDefine.class, definitionFile); for (DataDefine dataDefine : definitionLoader) { logger.info("loaded data definition class: {}", dataDefine.getClass().getName()); - dataDefine.initial(); dataDefineMap.put(dataDefine.defineId(), dataDefine); } return dataDefineMap; diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/TransformToData.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/TransformToData.java new file mode 100644 index 000000000..1ce257f78 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/TransformToData.java @@ -0,0 +1,8 @@ +package org.skywalking.apm.collector.stream.worker.impl.data; + +/** + * @author pengys5 + */ +public interface TransformToData { + Data transform(); +} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java index d292dd5ef..887f24809 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java @@ -1,10 +1,14 @@ package org.skywalking.apm.collector.stream.worker.impl.data; +import java.util.concurrent.atomic.AtomicInteger; + /** * @author pengys5 */ public abstract class Window { + private AtomicInteger windowSwitch = new AtomicInteger(0); + private DataCollection pointer; private DataCollection windowDataA; @@ -16,6 +20,15 @@ public abstract class Window { pointer = windowDataA; } + public boolean trySwitchPointer() { + if (windowSwitch.incrementAndGet() == 1) { + return true; + } else { + windowSwitch.addAndGet(-1); + return false; + } + } + public void switchPointer() { if (pointer == windowDataA) { pointer = windowDataB; @@ -41,4 +54,9 @@ public abstract class Window { return windowDataA; } } + + public void releaseLast() { + getLast().clear(); + windowSwitch.addAndGet(-1); + } } diff --git a/apm-collector/pom.xml b/apm-collector/pom.xml index 7eda101c6..aa077d687 100644 --- a/apm-collector/pom.xml +++ b/apm-collector/pom.xml @@ -39,5 +39,15 @@ logback-classic 1.2.3 + + org.slf4j + log4j-over-slf4j + 1.7.25 + + + org.apache.logging.log4j + log4j-core + 2.8.2 +