diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java index 2cf836e90..fd169c8f7 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java @@ -1,10 +1,9 @@ package org.skywalking.apm.collector.agentregister.application; +import org.skywalking.apm.collector.agentstream.worker.cache.ApplicationCache; import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationDataDefine; import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterRemoteWorker; -import org.skywalking.apm.collector.agentstream.worker.register.application.dao.IApplicationDAO; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; -import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; @@ -20,8 +19,7 @@ public class ApplicationIDService { private final Logger logger = LoggerFactory.getLogger(ApplicationIDService.class); public int getOrCreate(String applicationCode) { - IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName()); - int applicationId = dao.getApplicationId(applicationCode); + int applicationId = ApplicationCache.get(applicationCode); if (applicationId == 0) { StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); 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 2b53fa68a..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,7 +2,6 @@ package org.skywalking.apm.collector.agentstream; import java.util.Iterator; import java.util.Map; -import org.skywalking.apm.collector.agentstream.worker.storage.IDNameExchangeTimer; 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; @@ -36,6 +35,5 @@ public class AgentStreamModuleInstaller implements ModuleInstaller { } new PersistenceTimer().start(); - new IDNameExchangeTimer().start(); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServiceHandler.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServiceHandler.java new file mode 100644 index 000000000..a92102319 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServiceHandler.java @@ -0,0 +1,34 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler; + +import com.google.gson.JsonElement; +import java.io.BufferedReader; +import java.io.IOException; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class TraceSegmentServiceHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(TraceSegmentServiceHandler.class); + + @Override public String pathSpec() { + return null; + } + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } + + @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + try { + BufferedReader reader = req.getReader(); + } catch (IOException e) { + logger.error(e.getMessage(), e); + } + } +} 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 9ddc733b8..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 @@ -7,6 +7,5 @@ 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_EXCHANGE_TIMES = "exchange_times"; 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/cache/ApplicationCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java new file mode 100644 index 000000000..ff5574029 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java @@ -0,0 +1,25 @@ +package org.skywalking.apm.collector.agentstream.worker.cache; + +import com.google.common.cache.Cache; +import com.google.common.cache.CacheBuilder; +import org.skywalking.apm.collector.agentstream.worker.register.application.dao.IApplicationDAO; +import org.skywalking.apm.collector.storage.dao.DAOContainer; + +/** + * @author pengys5 + */ +public class ApplicationCache { + + private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(1000).build(); + + public static int get(String applicationCode) { + try { + return CACHE.get(applicationCode, () -> { + IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName()); + return dao.getApplicationId(applicationCode); + }); + } catch (Throwable e) { + return 0; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ComponentCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ComponentCache.java deleted file mode 100644 index fe472ce6e..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ComponentCache.java +++ /dev/null @@ -1,26 +0,0 @@ -package org.skywalking.apm.collector.agentstream.worker.cache; - -import com.google.common.cache.Cache; -import com.google.common.cache.CacheBuilder; -import org.skywalking.apm.collector.agentstream.worker.Const; -import org.skywalking.apm.collector.agentstream.worker.node.component.dao.INodeComponentDAO; -import org.skywalking.apm.collector.storage.dao.DAOContainer; - -/** - * @author pengys5 - */ -public class ComponentCache { - - private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(1000).build(); - - public static int get(int applicationId, String componentName) { - try { - return CACHE.get(applicationId + Const.ID_SPLIT + componentName, () -> { - INodeComponentDAO dao = (INodeComponentDAO)DAOContainer.INSTANCE.get(INodeComponentDAO.class.getName()); - return dao.getComponentId(applicationId, componentName); - }); - } catch (Throwable e) { - return 0; - } - } -} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceNameCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceNameCache.java new file mode 100644 index 000000000..f6a88b980 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceNameCache.java @@ -0,0 +1,26 @@ +package org.skywalking.apm.collector.agentstream.worker.cache; + +import com.google.common.cache.Cache; +import com.google.common.cache.CacheBuilder; +import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO; +import org.skywalking.apm.collector.storage.dao.DAOContainer; + +/** + * @author pengys5 + */ +public class ServiceNameCache { + + private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(2000).build(); + + public static int get(int applicationId, String serviceName) { + try { + return CACHE.get(applicationId + Const.ID_SPLIT + serviceName, () -> { + IServiceNameDAO dao = (IServiceNameDAO)DAOContainer.INSTANCE.get(IServiceNameDAO.class.getName()); + return dao.getServiceId(applicationId, serviceName); + }); + } catch (Throwable e) { + return 0; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTracePersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTracePersistenceWorker.java index ba58718b6..8998c71ee 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTracePersistenceWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/GlobalTracePersistenceWorker.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.global; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.agentstream.worker.global.dao.IGlobalTraceDAO; import org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceDataDefine; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -10,7 +8,7 @@ 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.dao.IPersistenceDAO; 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.WorkerSelector; @@ -28,9 +26,12 @@ public class GlobalTracePersistenceWorker extends PersistenceWorker { super.preStart(); } - @Override protected List prepareBatch(Map dataMap) { - IGlobalTraceDAO dao = (IGlobalTraceDAO)DAOContainer.INSTANCE.get(IGlobalTraceDAO.class.getName()); - return dao.prepareBatch(dataMap); + @Override protected boolean needMergeDBData() { + return false; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(IGlobalTraceDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceEsDAO.java index 4c1c49707..c36ca1108 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceEsDAO.java @@ -1,35 +1,38 @@ package org.skywalking.apm.collector.agentstream.worker.global.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.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public class GlobalTraceEsDAO extends EsDAO implements IGlobalTraceDAO { +public class GlobalTraceEsDAO extends EsDAO implements IGlobalTraceDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(GlobalTraceEsDAO.class); - @Override public List prepareBatch(Map dataMap) { - List indexRequestBuilders = new ArrayList<>(); - dataMap.forEach((id, data) -> { - logger.debug("global trace prepareBatch, id: {}", id); - Map source = new HashMap(); - source.put(GlobalTraceTable.COLUMN_SEGMENT_ID, data.getDataString(1)); - source.put(GlobalTraceTable.COLUMN_GLOBAL_TRACE_ID, data.getDataString(2)); - source.put(GlobalTraceTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); - logger.debug("global trace source: {}", source.toString()); - IndexRequestBuilder builder = getClient().prepareIndex(GlobalTraceTable.TABLE, id).setSource(source); - indexRequestBuilders.add(builder); - }); - return indexRequestBuilders; + @Override public Data get(String id, DataDefine dataDefine) { + return null; + } + + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + return null; + } + + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + source.put(GlobalTraceTable.COLUMN_SEGMENT_ID, data.getDataString(1)); + source.put(GlobalTraceTable.COLUMN_GLOBAL_TRACE_ID, data.getDataString(2)); + source.put(GlobalTraceTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + logger.debug("global trace source: {}", source.toString()); + return getClient().prepareIndex(GlobalTraceTable.TABLE, data.getDataString(0)).setSource(source); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java index 1a34c616a..ea48e612f 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java @@ -1,15 +1,9 @@ package org.skywalking.apm.collector.agentstream.worker.global.dao; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; /** * @author pengys5 */ public class GlobalTraceH2DAO extends H2DAO implements IGlobalTraceDAO { - - @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/global/dao/IGlobalTraceDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/IGlobalTraceDAO.java index 4138f423c..1b73837c8 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/IGlobalTraceDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/IGlobalTraceDAO.java @@ -1,12 +1,7 @@ package org.skywalking.apm.collector.agentstream.worker.global.dao; -import java.util.List; -import java.util.Map; -import org.skywalking.apm.collector.stream.worker.impl.data.Data; - /** * @author pengys5 */ public interface IGlobalTraceDAO { - List prepareBatch(Map dataMap); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentExchangeWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentExchangeWorker.java deleted file mode 100644 index 056847430..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentExchangeWorker.java +++ /dev/null @@ -1,83 +0,0 @@ -package org.skywalking.apm.collector.agentstream.worker.node.component; - -import org.skywalking.apm.collector.agentstream.worker.cache.ComponentCache; -import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine; -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.WorkerInvokeException; -import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; -import org.skywalking.apm.collector.stream.worker.impl.ExchangeWorker; -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; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * @author pengys5 - */ -public class NodeComponentExchangeWorker extends ExchangeWorker { - - private final Logger logger = LoggerFactory.getLogger(NodeComponentExchangeWorker.class); - - public NodeComponentExchangeWorker(Role role, ClusterWorkerContext clusterContext) { - super(role, clusterContext); - } - - @Override public void preStart() throws ProviderNotFoundException { - super.preStart(); - } - - @Override protected void exchange(Data data) { - NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); - nodeComponent.toSelf(data); - - int componentId = ComponentCache.get(nodeComponent.getApplicationId(), nodeComponent.getComponentName()); - if (componentId == 0 && nodeComponent.getTimes() < 10) { - try { - nodeComponent.increase(); - getClusterContext().lookup(NodeComponentExchangeWorker.WorkerRole.INSTANCE).tell(nodeComponent.toData()); - } catch (WorkerNotFoundException | WorkerInvokeException e) { - logger.error(e.getMessage(), e); - } - } - } - - public static class Factory extends AbstractLocalAsyncWorkerProvider { - @Override - public Role role() { - return WorkerRole.INSTANCE; - } - - @Override - public NodeComponentExchangeWorker workerInstance(ClusterWorkerContext clusterContext) { - return new NodeComponentExchangeWorker(role(), clusterContext); - } - - @Override - public int queueSize() { - return 1024; - } - } - - public enum WorkerRole implements Role { - INSTANCE; - - @Override - public String roleName() { - return NodeComponentExchangeWorker.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/NodeComponentPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentPersistenceWorker.java index 2e7ccf0b9..5ff722b86 100644 --- 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 @@ -1,7 +1,5 @@ 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; @@ -10,7 +8,7 @@ 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.dao.IPersistenceDAO; 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; @@ -28,9 +26,12 @@ public class NodeComponentPersistenceWorker extends PersistenceWorker { super.preStart(); } - @Override protected List prepareBatch(Map dataMap) { - INodeComponentDAO dao = (INodeComponentDAO)DAOContainer.INSTANCE.get(INodeComponentDAO.class.getName()); - return dao.prepareBatch(dataMap); + @Override protected boolean needMergeDBData() { + return true; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(INodeComponentDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { 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 3bc4be924..4f4691802 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 @@ -3,85 +3,89 @@ package org.skywalking.apm.collector.agentstream.worker.node.component; import java.util.ArrayList; import java.util.List; import org.skywalking.apm.collector.agentstream.worker.Const; -import org.skywalking.apm.collector.agentstream.worker.cache.ComponentCache; import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine; 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.segment.LocalSpanListener; +import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; +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.skywalking.apm.network.trace.component.ComponentsDefine; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanListener, LocalSpanListener { +public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanListener, FirstSpanListener, LocalSpanListener { private final Logger logger = LoggerFactory.getLogger(NodeComponentSpanListener.class); - private List nodeComponents = new ArrayList<>(); + private List nodeComponents = new ArrayList<>(); + private long timeBucket; @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String componentName = ComponentsDefine.getInstance().getComponentName(spanObject.getComponentId()); - createNodeComponent(spanObject, applicationId, componentName); + String componentName = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getComponentId()); + if (spanObject.getComponentId() == 0) { + componentName = spanObject.getComponent(); + } + String peer = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getPeerId()); + if (spanObject.getPeerId() == 0) { + peer = spanObject.getPeer(); + } + + String agg = componentName + Const.ID_SPLIT + peer; + nodeComponents.add(agg); } @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String componentName = ComponentsDefine.getInstance().getComponentName(spanObject.getComponentId()); - createNodeComponent(spanObject, applicationId, componentName); + buildEntryOrLocal(spanObject, applicationId); } @Override public void parseLocal(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - int componentId = ComponentCache.get(applicationId, spanObject.getComponent()); - - NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); - nodeComponent.setApplicationId(applicationId); - nodeComponent.setComponentId(componentId); - nodeComponent.setComponentName(spanObject.getComponent()); - - if (componentId == 0) { - StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); - - logger.debug("send to node component exchange worker, id: {}", nodeComponent.getId()); - nodeComponent.setId(applicationId + Const.ID_SPLIT + spanObject.getComponent()); - try { - context.getClusterWorkerContext().lookup(NodeComponentExchangeWorker.WorkerRole.INSTANCE).tell(nodeComponent.toData()); - } catch (WorkerInvokeException | WorkerNotFoundException e) { - logger.error(e.getMessage(), e); - } - } else { - nodeComponent.setId(applicationId + Const.ID_SPLIT + componentId); - nodeComponents.add(nodeComponent); - } + buildEntryOrLocal(spanObject, applicationId); } - private void createNodeComponent(SpanObject spanObject, int applicationId, String componentName) { - NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); - nodeComponent.setApplicationId(applicationId); - nodeComponent.setComponentId(spanObject.getComponentId()); - nodeComponent.setComponentName(componentName); - nodeComponent.setId(applicationId + Const.ID_SPLIT + spanObject.getComponentId()); - nodeComponents.add(nodeComponent); + private void buildEntryOrLocal(SpanObject spanObject, int applicationId) { + String componentName = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getComponentId()); + + if (spanObject.getComponentId() == 0) { + componentName = spanObject.getComponent(); + } + + String peer = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); + String agg = componentName + Const.ID_SPLIT + peer; + nodeComponents.add(agg); + } + + @Override + public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { + timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); } @Override public void build() { StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); - for (NodeComponentDataDefine.NodeComponent nodeComponent : nodeComponents) { + + nodeComponents.forEach(agg -> { + 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.toData()); } 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 index 3ae17bd74..79f44ccb5 100644 --- 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 @@ -1,14 +1,7 @@ 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); - - int getComponentId(int applicationId, String componentName); } 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 index 54e2e1442..a3e114f90 100644 --- 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 @@ -1,58 +1,47 @@ 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.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; -import org.elasticsearch.action.search.SearchRequestBuilder; -import org.elasticsearch.action.search.SearchResponse; -import org.elasticsearch.action.search.SearchType; -import org.elasticsearch.index.query.BoolQueryBuilder; -import org.elasticsearch.index.query.QueryBuilders; -import org.elasticsearch.search.SearchHit; +import org.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentTable; -import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; /** * @author pengys5 */ -public class NodeComponentEsDAO extends EsDAO implements INodeComponentDAO { +public class NodeComponentEsDAO extends EsDAO implements INodeComponentDAO, IPersistenceDAO { - @Override public List prepareBatch(Map dataMap) { - List indexRequestBuilders = new ArrayList<>(); - dataMap.forEach((id, data) -> { - Map source = new HashMap(); - source.put(NodeComponentTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); - source.put(NodeComponentTable.COLUMN_COMPONENT_NAME, data.getDataString(1)); - source.put(NodeComponentTable.COLUMN_COMPONENT_ID, data.getDataInteger(1)); - - IndexRequestBuilder builder = getClient().prepareIndex(NodeComponentTable.TABLE, id).setSource(source); - indexRequestBuilders.add(builder); - }); - return indexRequestBuilders; + @Override public Data get(String id, DataDefine dataDefine) { + GetResponse getResponse = getClient().prepareGet(NodeComponentTable.TABLE, id).get(); + if (getResponse.isExists()) { + Data data = dataDefine.build(id); + Map source = getResponse.getSource(); + data.setDataString(1, (String)source.get(NodeComponentTable.COLUMN_AGG)); + data.setDataLong(0, (Long)source.get(NodeComponentTable.COLUMN_TIME_BUCKET)); + return data; + } else { + return null; + } } - public int getComponentId(int applicationId, String componentName) { - ElasticSearchClient client = getClient(); + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + source.put(NodeComponentTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); - SearchRequestBuilder searchRequestBuilder = client.prepareSearch(NodeComponentTable.TABLE); - searchRequestBuilder.setTypes("type"); - searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH); - BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery(); - boolQueryBuilder.must(QueryBuilders.termQuery(NodeComponentTable.COLUMN_APPLICATION_ID, applicationId)); - boolQueryBuilder.must(QueryBuilders.termQuery(NodeComponentTable.COLUMN_COMPONENT_NAME, componentName)); - searchRequestBuilder.setQuery(boolQueryBuilder); - searchRequestBuilder.setSize(1); + return getClient().prepareIndex(NodeComponentTable.TABLE, data.getDataString(0)).setSource(source); + } - SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); - if (searchResponse.getHits().totalHits > 0) { - SearchHit searchHit = searchResponse.getHits().iterator().next(); - int componentId = (int)searchHit.getSource().get(NodeComponentTable.COLUMN_COMPONENT_ID); - return componentId; - } - return 0; + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + source.put(NodeComponentTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + return getClient().prepareUpdate(NodeComponentTable.TABLE, data.getDataString(0)).setDoc(source); } } 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 index 3277e059d..5e08b3445 100644 --- 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 @@ -1,7 +1,5 @@ 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; /** @@ -9,11 +7,4 @@ import org.skywalking.apm.collector.storage.h2.dao.H2DAO; */ public class NodeComponentH2DAO extends H2DAO implements INodeComponentDAO { - @Override public List prepareBatch(Map map) { - return null; - } - - @Override public int getComponentId(int applicationId, String componentName) { - return 0; - } } 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 7e2342d22..0f08484ce 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 @@ -5,7 +5,6 @@ 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.Exchange; import org.skywalking.apm.collector.stream.worker.impl.data.Transform; import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; @@ -20,15 +19,13 @@ public class NodeComponentDataDefine extends DataDefine { } @Override protected int initialCapacity() { - return 5; + return 3; } @Override protected void attributeDefine() { addAttribute(0, new Attribute(NodeComponentTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); - addAttribute(1, new Attribute(NodeComponentTable.COLUMN_APPLICATION_ID, AttributeType.INTEGER, new CoverOperation())); - addAttribute(2, new Attribute(NodeComponentTable.COLUMN_COMPONENT_NAME, AttributeType.STRING, new CoverOperation())); - addAttribute(3, new Attribute(NodeComponentTable.COLUMN_COMPONENT_ID, AttributeType.INTEGER, new CoverOperation())); - addAttribute(4, new Attribute(NodeComponentTable.COLUMN_EXCHANGE_TIMES, AttributeType.INTEGER, new NonOperation())); + addAttribute(1, new Attribute(NodeComponentTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); + addAttribute(2, new Attribute(NodeComponentTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); } @Override public Object deserialize(RemoteData remoteData) { @@ -39,41 +36,33 @@ public class NodeComponentDataDefine extends DataDefine { return null; } - public static class NodeComponent extends Exchange implements Transform { + public static class NodeComponent implements Transform { private String id; - private int applicationId; - private String componentName; - private int componentId; + private String agg; + private long timeBucket; - public NodeComponent(String id, int applicationId, String componentName, int componentId) { - super(0); + NodeComponent(String id, String agg, long timeBucket) { this.id = id; - this.applicationId = applicationId; - this.componentName = componentName; - this.componentId = componentId; + this.agg = agg; + this.timeBucket = timeBucket; } public NodeComponent() { - super(0); } @Override public Data toData() { NodeComponentDataDefine define = new NodeComponentDataDefine(); Data data = define.build(id); data.setDataString(0, this.id); - data.setDataInteger(0, this.applicationId); - data.setDataString(1, this.componentName); - data.setDataInteger(1, this.componentId); - data.setDataInteger(2, this.getTimes()); + data.setDataString(1, this.agg); + data.setDataLong(0, this.timeBucket); return data; } @Override public NodeComponent toSelf(Data data) { this.id = data.getDataString(0); - this.applicationId = data.getDataInteger(0); - this.componentName = data.getDataString(1); - this.componentId = data.getDataInteger(1); - this.setTimes(data.getDataInteger(2)); + this.agg = data.getDataString(1); + this.timeBucket = data.getDataLong(0); return this; } @@ -81,32 +70,24 @@ public class NodeComponentDataDefine extends DataDefine { return id; } + public String getAgg() { + return agg; + } + + public long getTimeBucket() { + return timeBucket; + } + public void setId(String id) { this.id = id; } - public String getComponentName() { - return componentName; + public void setAgg(String agg) { + this.agg = agg; } - public void setComponentName(String componentName) { - this.componentName = componentName; - } - - public int getComponentId() { - return componentId; - } - - public void setComponentId(int componentId) { - this.componentId = componentId; - } - - public int getApplicationId() { - return applicationId; - } - - public void setApplicationId(int applicationId) { - this.applicationId = applicationId; + public void setTimeBucket(long timeBucket) { + this.timeBucket = timeBucket; } } } 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 9085e2b83..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 @@ -25,8 +25,7 @@ public class NodeComponentEsTableDefine extends ElasticSearchTableDefine { } @Override public void initialize() { - addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); - addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_COMPONENT_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); - addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_COMPONENT_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java index b85787ee7..62b77b04e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java @@ -14,8 +14,7 @@ public class NodeComponentH2TableDefine extends H2TableDefine { @Override public void initialize() { addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_ID, H2ColumnDefine.Type.Varchar.name())); - addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); - addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_COMPONENT_ID, H2ColumnDefine.Type.Int.name())); - addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_COMPONENT_NAME, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentTable.java index e95365287..1bce8df83 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentTable.java @@ -7,7 +7,4 @@ import org.skywalking.apm.collector.agentstream.worker.CommonTable; */ public class NodeComponentTable extends CommonTable { public static final String TABLE = "node_component"; - public static final String COLUMN_APPLICATION_ID = "application_id"; - public static final String COLUMN_COMPONENT_NAME = "component_name"; - public static final String COLUMN_COMPONENT_ID = "component_id"; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingPersistenceWorker.java index 48a882d5a..57edf24c5 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingPersistenceWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingPersistenceWorker.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.INodeMappingDAO; import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -10,7 +8,7 @@ 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.dao.IPersistenceDAO; 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; @@ -28,9 +26,12 @@ public class NodeMappingPersistenceWorker extends PersistenceWorker { super.preStart(); } - @Override protected List prepareBatch(Map dataMap) { - INodeMappingDAO dao = (INodeMappingDAO)DAOContainer.INSTANCE.get(INodeMappingDAO.class.getName()); - return dao.prepareBatch(dataMap); + @Override protected boolean needMergeDBData() { + return true; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(INodeMappingDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java index 4b8e4af71..2c5ebe600 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java @@ -6,6 +6,7 @@ import org.skywalking.apm.collector.agentstream.worker.Const; import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine; import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; +import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; @@ -30,9 +31,9 @@ public class NodeMappingSpanListener implements RefsListener, FirstSpanListener @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId, String segmentId) { logger.debug("node mapping listener parse reference"); - String peers = Const.PEERS_FRONT_SPLIT + reference.getNetworkAddressId() + Const.PEERS_BEHIND_SPLIT; + String peers = reference.getNetworkAddress(); if (reference.getNetworkAddressId() == 0) { - peers = Const.PEERS_FRONT_SPLIT + reference.getNetworkAddress() + Const.PEERS_BEHIND_SPLIT; + peers = ExchangeMarkUtils.INSTANCE.buildMarkedID(reference.getNetworkAddressId()); } String agg = applicationId + Const.ID_SPLIT + peers; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/INodeMappingDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/INodeMappingDAO.java index f864617a2..450cc5451 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/INodeMappingDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/INodeMappingDAO.java @@ -1,12 +1,7 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.dao; -import java.util.List; -import java.util.Map; -import org.skywalking.apm.collector.stream.worker.impl.data.Data; - /** * @author pengys5 */ public interface INodeMappingDAO { - List prepareBatch(Map dataMap); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java index 4a36d0378..66a22211d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java @@ -1,29 +1,47 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.dao; -import java.util.ArrayList; import java.util.HashMap; -import java.util.List; import java.util.Map; +import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; +import org.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; /** * @author pengys5 */ -public class NodeMappingEsDAO extends EsDAO implements INodeMappingDAO { +public class NodeMappingEsDAO extends EsDAO implements INodeMappingDAO, IPersistenceDAO { - @Override public List prepareBatch(Map dataMap) { - List indexRequestBuilders = new ArrayList<>(); - dataMap.forEach((id, data) -> { - Map source = new HashMap(); - source.put(NodeMappingTable.COLUMN_AGG, data.getDataString(1)); - source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + @Override public Data get(String id, DataDefine dataDefine) { + GetResponse getResponse = getClient().prepareGet(NodeMappingTable.TABLE, id).get(); + if (getResponse.isExists()) { + Data data = dataDefine.build(id); + Map source = getResponse.getSource(); + data.setDataString(1, (String)source.get(NodeMappingTable.COLUMN_AGG)); + data.setDataLong(0, (Long)source.get(NodeMappingTable.COLUMN_TIME_BUCKET)); + return data; + } else { + return null; + } + } - IndexRequestBuilder builder = getClient().prepareIndex(NodeMappingTable.TABLE, id).setSource(source); - indexRequestBuilders.add(builder); - }); - return indexRequestBuilders; + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + source.put(NodeMappingTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + return getClient().prepareIndex(NodeMappingTable.TABLE, data.getDataString(0)).setSource(source); + } + + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + source.put(NodeMappingTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + return getClient().prepareUpdate(NodeMappingTable.TABLE, data.getDataString(0)).setDoc(source); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java index 8ad24c091..045d0738e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java @@ -1,15 +1,9 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.dao; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; /** * @author pengys5 */ public class NodeMappingH2DAO extends H2DAO implements INodeMappingDAO { - - @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/mapping/define/NodeMappingDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java index 4c949355b..e227b291e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java @@ -36,12 +36,12 @@ public class NodeMappingDataDefine extends DataDefine { return null; } - public static class NodeMapping implements Transform { + public static class NodeMapping implements Transform { private String id; private String agg; private long timeBucket; - public NodeMapping(String id, String agg, long timeBucket) { + NodeMapping(String id, String agg, long timeBucket) { this.id = id; this.agg = agg; this.timeBucket = timeBucket; @@ -59,8 +59,11 @@ public class NodeMappingDataDefine extends DataDefine { return data; } - @Override public Object toSelf(Data data) { - return null; + @Override public NodeMapping toSelf(Data data) { + this.id = data.getDataString(0); + this.agg = data.getDataString(1); + this.timeBucket = data.getDataLong(0); + return this; } public String getId() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefPersistenceWorker.java index 06e363a1a..ec48c3a8c 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefPersistenceWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefPersistenceWorker.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao.INodeReferenceDAO; import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefDataDefine; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -10,7 +8,7 @@ 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.dao.IPersistenceDAO; 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; @@ -28,9 +26,12 @@ public class NodeRefPersistenceWorker extends PersistenceWorker { super.preStart(); } - @Override protected List prepareBatch(Map dataMap) { - INodeReferenceDAO dao = (INodeReferenceDAO)DAOContainer.INSTANCE.get(INodeReferenceDAO.class.getName()); - return dao.prepareBatch(dataMap); + @Override protected boolean needMergeDBData() { + return true; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(INodeReferenceDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java index 0e5da0430..a57e678b1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java @@ -8,6 +8,7 @@ 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.segment.RefsListener; +import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; @@ -34,9 +35,9 @@ public class NodeRefSpanListener implements EntrySpanListener, ExitSpanListener, @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { String front = String.valueOf(applicationId); - String behind = String.valueOf(spanObject.getPeerId()); + String behind = spanObject.getPeer(); if (spanObject.getPeerId() == 0) { - behind = spanObject.getPeer(); + behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getPeerId()); } String agg = front + Const.ID_SPLIT + behind; @@ -45,9 +46,8 @@ public class NodeRefSpanListener implements EntrySpanListener, ExitSpanListener, @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String behind = String.valueOf(applicationId); + String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); String front = Const.USER_CODE; - String agg = front + Const.ID_SPLIT + behind; nodeEntryReferences.add(agg); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/INodeReferenceDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/INodeReferenceDAO.java index 122829dfb..7f2406d1e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/INodeReferenceDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/INodeReferenceDAO.java @@ -1,12 +1,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao; -import java.util.List; -import java.util.Map; -import org.skywalking.apm.collector.stream.worker.impl.data.Data; - /** * @author pengys5 */ public interface INodeReferenceDAO { - List prepareBatch(Map dataMap); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceEsDAO.java index d1334bec6..84d572a70 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceEsDAO.java @@ -1,29 +1,47 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao; -import java.util.ArrayList; import java.util.HashMap; -import java.util.List; import java.util.Map; +import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; +import org.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; /** * @author pengys5 */ -public class NodeReferenceEsDAO extends EsDAO implements INodeReferenceDAO { +public class NodeReferenceEsDAO extends EsDAO implements INodeReferenceDAO, IPersistenceDAO { - @Override public List prepareBatch(Map dataMap) { - List indexRequestBuilders = new ArrayList<>(); - dataMap.forEach((id, data) -> { - Map source = new HashMap(); - source.put(NodeRefTable.COLUMN_AGG, data.getDataString(1)); - source.put(NodeRefTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + @Override public Data get(String id, DataDefine dataDefine) { + GetResponse getResponse = getClient().prepareGet(NodeRefTable.TABLE, id).get(); + if (getResponse.isExists()) { + Data data = dataDefine.build(id); + Map source = getResponse.getSource(); + data.setDataString(1, (String)source.get(NodeRefTable.COLUMN_AGG)); + data.setDataLong(0, (Long)source.get(NodeRefTable.COLUMN_TIME_BUCKET)); + return data; + } else { + return null; + } + } - IndexRequestBuilder builder = getClient().prepareIndex(NodeRefTable.TABLE, id).setSource(source); - indexRequestBuilders.add(builder); - }); - return indexRequestBuilders; + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + source.put(NodeRefTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeRefTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + return getClient().prepareIndex(NodeRefTable.TABLE, data.getDataString(0)).setSource(source); + } + + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + source.put(NodeRefTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeRefTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + return getClient().prepareUpdate(NodeRefTable.TABLE, data.getDataString(0)).setDoc(source); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceH2DAO.java index bc2d7a0a1..7820be366 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/dao/NodeReferenceH2DAO.java @@ -1,15 +1,24 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; /** * @author pengys5 */ -public class NodeReferenceH2DAO extends H2DAO implements INodeReferenceDAO { +public class NodeReferenceH2DAO extends H2DAO implements INodeReferenceDAO, IPersistenceDAO { - @Override public List prepareBatch(Map map) { + @Override public Data get(String id, DataDefine dataDefine) { + return null; + } + + @Override public String prepareBatchInsert(Data data) { + return null; + } + + @Override public String prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java index 10c2c969f..079062861 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java @@ -14,10 +14,8 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class NodeRefDataDefine extends DataDefine { - public static final int DEFINE_ID = 201; - @Override public int defineId() { - return DEFINE_ID; + return 201; } @Override protected int initialCapacity() { @@ -51,7 +49,7 @@ public class NodeRefDataDefine extends DataDefine { private String agg; private long timeBucket; - public NodeReference(String id, String agg, long timeBucket) { + NodeReference(String id, String agg, long timeBucket) { this.id = id; this.agg = agg; this.timeBucket = timeBucket; @@ -70,7 +68,10 @@ public class NodeRefDataDefine extends DataDefine { } @Override public Object toSelf(Data data) { - return null; + this.id = data.getDataString(0); + this.agg = data.getDataString(1); + this.timeBucket = data.getDataLong(0); + return this; } public String getId() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumPersistenceWorker.java index 96dead7e0..6fcb36eb7 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumPersistenceWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumPersistenceWorker.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.INodeRefSumDAO; import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -10,7 +8,7 @@ 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.dao.IPersistenceDAO; 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; @@ -28,9 +26,12 @@ public class NodeRefSumPersistenceWorker extends PersistenceWorker { super.preStart(); } - @Override protected List prepareBatch(Map dataMap) { - INodeRefSumDAO dao = (INodeRefSumDAO)DAOContainer.INSTANCE.get(INodeRefSumDAO.class.getName()); - return dao.prepareBatch(dataMap); + @Override protected boolean needMergeDBData() { + return true; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(INodeRefSumDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java index 69bffe8ea..f3251d032 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java @@ -8,6 +8,7 @@ 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.segment.RefsListener; +import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; @@ -34,9 +35,9 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { String front = String.valueOf(applicationId); - String behind = String.valueOf(spanObject.getPeerId()); + String behind = spanObject.getPeer(); if (spanObject.getPeerId() == 0) { - behind = spanObject.getPeer(); + behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getPeerId()); } String agg = front + Const.ID_SPLIT + behind; @@ -45,7 +46,7 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String behind = String.valueOf(applicationId); + String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); String front = Const.USER_CODE; String agg = front + Const.ID_SPLIT + behind; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/INodeRefSumDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/INodeRefSumDAO.java index ef4f69737..6be979957 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/INodeRefSumDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/INodeRefSumDAO.java @@ -1,12 +1,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao; -import java.util.List; -import java.util.Map; -import org.skywalking.apm.collector.stream.worker.impl.data.Data; - /** * @author pengys5 */ public interface INodeRefSumDAO { - List prepareBatch(Map dataMap); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java index a20099784..c117e70a9 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java @@ -1,35 +1,66 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao; -import java.util.ArrayList; import java.util.HashMap; -import java.util.List; import java.util.Map; +import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; +import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefTable; import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; /** * @author pengys5 */ -public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO { +public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO, IPersistenceDAO { - @Override public List prepareBatch(Map dataMap) { - List indexRequestBuilders = new ArrayList<>(); - dataMap.forEach((id, data) -> { - Map source = new HashMap(); - source.put(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, data.getDataLong(0)); - source.put(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, data.getDataLong(1)); - source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, data.getDataLong(2)); - source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, data.getDataLong(3)); - source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataLong(4)); - source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataLong(5)); - source.put(NodeRefSumTable.COLUMN_AGG, data.getDataString(1)); - source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(6)); + @Override public Data get(String id, DataDefine dataDefine) { + GetResponse getResponse = getClient().prepareGet(NodeRefSumTable.TABLE, id).get(); + if (getResponse.isExists()) { + Data data = dataDefine.build(id); + Map source = getResponse.getSource(); + data.setDataLong(0, (Long)source.get(NodeRefSumTable.COLUMN_ONE_SECOND_LESS)); + data.setDataLong(1, (Long)source.get(NodeRefSumTable.COLUMN_THREE_SECOND_LESS)); + data.setDataLong(2, (Long)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS)); + data.setDataLong(3, (Long)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER)); + data.setDataLong(4, (Long)source.get(NodeRefSumTable.COLUMN_ERROR)); + data.setDataLong(5, (Long)source.get(NodeRefSumTable.COLUMN_SUMMARY)); + data.setDataLong(6, (Long)source.get(NodeRefSumTable.COLUMN_TIME_BUCKET)); + data.setDataString(1, (String)source.get(NodeRefSumTable.COLUMN_AGG)); + return data; + } else { + return null; + } + } - IndexRequestBuilder builder = getClient().prepareIndex(NodeRefSumTable.TABLE, id).setSource(source); - indexRequestBuilders.add(builder); - }); - return indexRequestBuilders; + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + source.put(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, data.getDataLong(0)); + source.put(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, data.getDataLong(1)); + source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, data.getDataLong(2)); + source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, data.getDataLong(3)); + source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataLong(4)); + source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataLong(5)); + source.put(NodeRefSumTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(6)); + + return getClient().prepareIndex(NodeRefSumTable.TABLE, data.getDataString(0)).setSource(source); + } + + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + source.put(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, data.getDataLong(0)); + source.put(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, data.getDataLong(1)); + source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, data.getDataLong(2)); + source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, data.getDataLong(3)); + source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataLong(4)); + source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataLong(5)); + source.put(NodeRefSumTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(6)); + + return getClient().prepareUpdate(NodeRefTable.TABLE, data.getDataString(0)).setDoc(source); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumH2DAO.java index 2a3b83501..7483c8a82 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumH2DAO.java @@ -1,15 +1,9 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; /** * @author pengys5 */ public class NodeRefSumH2DAO extends H2DAO implements INodeRefSumDAO { - - @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/segment/SegmentParse.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java index 3d3738cd2..6558a2615 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 @@ -38,7 +38,6 @@ public class SegmentParse { spanListeners.add(new NodeRefSpanListener()); spanListeners.add(new NodeComponentSpanListener()); spanListeners.add(new NodeMappingSpanListener()); - spanListeners.add(new NodeRefSpanListener()); spanListeners.add(new NodeRefSumSpanListener()); spanListeners.add(new SegmentCostSpanListener()); spanListeners.add(new GlobalTraceSpanListener()); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostPersistenceWorker.java index d6b72f8a2..036a291fa 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostPersistenceWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostPersistenceWorker.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.agentstream.worker.segment.cost.dao.ISegmentCostDAO; import org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostDataDefine; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -10,7 +8,7 @@ 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.dao.IPersistenceDAO; 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.WorkerSelector; @@ -28,9 +26,12 @@ public class SegmentCostPersistenceWorker extends PersistenceWorker { super.preStart(); } - @Override protected List prepareBatch(Map dataMap) { - ISegmentCostDAO dao = (ISegmentCostDAO)DAOContainer.INSTANCE.get(ISegmentCostDAO.class.getName()); - return dao.prepareBatch(dataMap); + @Override protected boolean needMergeDBData() { + return false; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(ISegmentCostDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java index e5934792f..c49e86183 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/SegmentCostSpanListener.java @@ -7,6 +7,7 @@ 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.segment.LocalSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostDataDefine; +import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.stream.StreamModuleContext; @@ -40,10 +41,14 @@ public class SegmentCostSpanListener implements EntrySpanListener, ExitSpanListe segmentCost.setCost(spanObject.getEndTime() - spanObject.getStartTime()); segmentCost.setStartTime(spanObject.getStartTime()); segmentCost.setEndTime(spanObject.getEndTime()); - segmentCost.setOperationName(spanObject.getOperationName()); segmentCost.setId(segmentId); - segmentCosts.add(segmentCost); + if (spanObject.getOperationNameId() == 0) { + segmentCost.setServiceName(spanObject.getOperationName()); + } else { + segmentCost.setServiceName(ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getOperationNameId())); + } + segmentCosts.add(segmentCost); isError = isError || spanObject.getIsError(); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/ISegmentCostDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/ISegmentCostDAO.java index 03787ca00..03a78dd9a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/ISegmentCostDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/ISegmentCostDAO.java @@ -1,12 +1,7 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost.dao; -import java.util.List; -import java.util.Map; -import org.skywalking.apm.collector.stream.worker.impl.data.Data; - /** * @author pengys5 */ public interface ISegmentCostDAO { - List prepareBatch(Map dataMap); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostEsDAO.java index d1b9e19b2..794b0d32a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostEsDAO.java @@ -1,39 +1,43 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost.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.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public class SegmentCostEsDAO extends EsDAO implements ISegmentCostDAO { +public class SegmentCostEsDAO extends EsDAO implements ISegmentCostDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(SegmentCostEsDAO.class); - @Override public List prepareBatch(Map dataMap) { - List indexRequestBuilders = new ArrayList<>(); - dataMap.forEach((id, data) -> { - logger.debug("segment cost prepareBatch, id: {}", id); - Map source = new HashMap(); - source.put(SegmentCostTable.COLUMN_SEGMENT_ID, data.getDataString(1)); - source.put(SegmentCostTable.COLUMN_OPERATION_NAME, data.getDataString(2)); - source.put(SegmentCostTable.COLUMN_COST, data.getDataLong(0)); - source.put(SegmentCostTable.COLUMN_START_TIME, data.getDataLong(1)); - source.put(SegmentCostTable.COLUMN_END_TIME, data.getDataLong(2)); - source.put(SegmentCostTable.COLUMN_IS_ERROR, data.getDataBoolean(0)); - source.put(SegmentCostTable.COLUMN_TIME_BUCKET, data.getDataLong(3)); - logger.debug("segment cost source: {}", source.toString()); - IndexRequestBuilder builder = getClient().prepareIndex(SegmentCostTable.TABLE, id).setSource(source); - indexRequestBuilders.add(builder); - }); - return indexRequestBuilders; + @Override public Data get(String id, DataDefine dataDefine) { + return null; + } + + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + return null; + } + + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + logger.debug("segment cost prepareBatchInsert, id: {}", data.getDataString(0)); + Map source = new HashMap<>(); + source.put(SegmentCostTable.COLUMN_SEGMENT_ID, data.getDataString(1)); + source.put(SegmentCostTable.COLUMN_SERVICE_NAME, data.getDataString(2)); + source.put(SegmentCostTable.COLUMN_COST, data.getDataLong(0)); + source.put(SegmentCostTable.COLUMN_START_TIME, data.getDataLong(1)); + source.put(SegmentCostTable.COLUMN_END_TIME, data.getDataLong(2)); + source.put(SegmentCostTable.COLUMN_IS_ERROR, data.getDataBoolean(0)); + source.put(SegmentCostTable.COLUMN_TIME_BUCKET, data.getDataLong(3)); + logger.debug("segment cost source: {}", source.toString()); + return getClient().prepareIndex(SegmentCostTable.TABLE, data.getDataString(0)).setSource(source); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java index 90fd044b9..63a6c78c5 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost.dao; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; /** @@ -9,7 +7,4 @@ import org.skywalking.apm.collector.storage.h2.dao.H2DAO; */ public class SegmentCostH2DAO extends H2DAO implements ISegmentCostDAO { - @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/segment/cost/define/SegmentCostDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostDataDefine.java index c86de4fcd..51f57fd43 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostDataDefine.java @@ -25,7 +25,7 @@ public class SegmentCostDataDefine extends DataDefine { @Override protected void attributeDefine() { addAttribute(0, new Attribute(SegmentCostTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); addAttribute(1, new Attribute(SegmentCostTable.COLUMN_SEGMENT_ID, AttributeType.STRING, new CoverOperation())); - addAttribute(2, new Attribute(SegmentCostTable.COLUMN_OPERATION_NAME, AttributeType.STRING, new CoverOperation())); + addAttribute(2, new Attribute(SegmentCostTable.COLUMN_SERVICE_NAME, AttributeType.STRING, new CoverOperation())); addAttribute(3, new Attribute(SegmentCostTable.COLUMN_COST, AttributeType.LONG, new CoverOperation())); addAttribute(4, new Attribute(SegmentCostTable.COLUMN_START_TIME, AttributeType.LONG, new CoverOperation())); addAttribute(5, new Attribute(SegmentCostTable.COLUMN_END_TIME, AttributeType.LONG, new CoverOperation())); @@ -36,13 +36,13 @@ public class SegmentCostDataDefine extends DataDefine { @Override public Object deserialize(RemoteData remoteData) { String id = remoteData.getDataStrings(0); String segmentId = remoteData.getDataStrings(1); - String operationName = remoteData.getDataStrings(2); + String serviceName = remoteData.getDataStrings(2); Long cost = remoteData.getDataLongs(0); Long startTime = remoteData.getDataLongs(1); Long endTime = remoteData.getDataLongs(2); Boolean isError = remoteData.getDataBooleans(0); Long timeBucket = remoteData.getDataLongs(2); - return new SegmentCost(id, segmentId, operationName, cost, startTime, endTime, isError, timeBucket); + return new SegmentCost(id, segmentId, serviceName, cost, startTime, endTime, isError, timeBucket); } @Override public RemoteData serialize(Object object) { @@ -50,7 +50,7 @@ public class SegmentCostDataDefine extends DataDefine { RemoteData.Builder builder = RemoteData.newBuilder(); builder.addDataStrings(segmentCost.getId()); builder.addDataStrings(segmentCost.getSegmentId()); - builder.addDataStrings(segmentCost.getOperationName()); + builder.addDataStrings(segmentCost.getServiceName()); builder.addDataLongs(segmentCost.getCost()); builder.addDataLongs(segmentCost.getStartTime()); builder.addDataLongs(segmentCost.getEndTime()); @@ -62,18 +62,18 @@ public class SegmentCostDataDefine extends DataDefine { public static class SegmentCost implements Transform { private String id; private String segmentId; - private String operationName; + private String serviceName; private Long cost; private Long startTime; private Long endTime; private boolean isError; private long timeBucket; - SegmentCost(String id, String segmentId, String operationName, Long cost, + SegmentCost(String id, String segmentId, String serviceName, Long cost, Long startTime, Long endTime, boolean isError, long timeBucket) { this.id = id; this.segmentId = segmentId; - this.operationName = operationName; + this.serviceName = serviceName; this.cost = cost; this.startTime = startTime; this.endTime = endTime; @@ -89,7 +89,7 @@ public class SegmentCostDataDefine extends DataDefine { Data data = define.build(id); data.setDataString(0, this.id); data.setDataString(1, this.segmentId); - data.setDataString(2, this.operationName); + data.setDataString(2, this.serviceName); data.setDataLong(0, this.cost); data.setDataLong(1, this.startTime); data.setDataLong(2, this.endTime); @@ -118,12 +118,12 @@ public class SegmentCostDataDefine extends DataDefine { this.segmentId = segmentId; } - public String getOperationName() { - return operationName; + public String getServiceName() { + return serviceName; } - public void setOperationName(String operationName) { - this.operationName = operationName; + public void setServiceName(String serviceName) { + this.serviceName = serviceName; } public Long getCost() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java index 4f493da99..e44e40829 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java @@ -26,7 +26,7 @@ public class SegmentCostEsTableDefine extends ElasticSearchTableDefine { @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_SEGMENT_ID, ElasticSearchColumnDefine.Type.Keyword.name())); - addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_OPERATION_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_SERVICE_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_COST, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_START_TIME, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_END_TIME, ElasticSearchColumnDefine.Type.Long.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostH2TableDefine.java index c4ca90226..f1d896919 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostH2TableDefine.java @@ -15,7 +15,7 @@ public class SegmentCostH2TableDefine extends H2TableDefine { @Override public void initialize() { addColumn(new H2ColumnDefine(SegmentCostTable.COLUMN_ID, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(SegmentCostTable.COLUMN_SEGMENT_ID, H2ColumnDefine.Type.Varchar.name())); - addColumn(new H2ColumnDefine(SegmentCostTable.COLUMN_OPERATION_NAME, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(SegmentCostTable.COLUMN_SERVICE_NAME, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(SegmentCostTable.COLUMN_COST, H2ColumnDefine.Type.Bigint.name())); addColumn(new H2ColumnDefine(SegmentCostTable.COLUMN_START_TIME, H2ColumnDefine.Type.Bigint.name())); addColumn(new H2ColumnDefine(SegmentCostTable.COLUMN_END_TIME, H2ColumnDefine.Type.Bigint.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostTable.java index 0bc0ee947..fb9c7f83d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostTable.java @@ -10,7 +10,7 @@ public class SegmentCostTable extends CommonTable { public static final String COLUMN_SEGMENT_ID = "segment_id"; public static final String COLUMN_START_TIME = "start_time"; public static final String COLUMN_END_TIME = "end_time"; - public static final String COLUMN_OPERATION_NAME = "operation_name"; + public static final String COLUMN_SERVICE_NAME = "service_name"; public static final String COLUMN_COST = "cost"; public static final String COLUMN_IS_ERROR = "is_error"; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/SegmentPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/SegmentPersistenceWorker.java index a76fbd11d..03bbf0351 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/SegmentPersistenceWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/SegmentPersistenceWorker.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.agentstream.worker.segment.origin.dao.ISegmentDAO; import org.skywalking.apm.collector.agentstream.worker.segment.origin.define.SegmentDataDefine; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -10,7 +8,7 @@ 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.dao.IPersistenceDAO; 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.WorkerSelector; @@ -28,9 +26,12 @@ public class SegmentPersistenceWorker extends PersistenceWorker { super.preStart(); } - @Override protected List prepareBatch(Map dataMap) { - ISegmentDAO dao = (ISegmentDAO)DAOContainer.INSTANCE.get(ISegmentDAO.class.getName()); - return dao.prepareBatch(dataMap); + @Override protected boolean needMergeDBData() { + return false; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(ISegmentDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/ISegmentDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/ISegmentDAO.java index 1c1d67e03..4e8bb0353 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/ISegmentDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/ISegmentDAO.java @@ -1,12 +1,7 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin.dao; -import java.util.List; -import java.util.Map; -import org.skywalking.apm.collector.stream.worker.impl.data.Data; - /** * @author pengys5 */ public interface ISegmentDAO { - List prepareBatch(Map dataMap); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentEsDAO.java index 11ced7431..edfb0cd12 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentEsDAO.java @@ -1,34 +1,37 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin.dao; -import java.util.ArrayList; import java.util.Base64; import java.util.HashMap; -import java.util.List; import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; +import org.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.agentstream.worker.segment.origin.define.SegmentTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public class SegmentEsDAO extends EsDAO implements ISegmentDAO { +public class SegmentEsDAO extends EsDAO implements ISegmentDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(SegmentEsDAO.class); - @Override public List prepareBatch(Map dataMap) { - List indexRequestBuilders = new ArrayList<>(); - dataMap.forEach((id, data) -> { - logger.debug("segment prepareBatch, id: {}", id); - Map source = new HashMap(); - source.put(SegmentTable.COLUMN_DATA_BINARY, new String(Base64.getEncoder().encode(data.getDataBytes(0)))); - logger.debug("segment source: {}", source.toString()); - IndexRequestBuilder builder = getClient().prepareIndex(SegmentTable.TABLE, id).setSource(source); - indexRequestBuilders.add(builder); - }); - return indexRequestBuilders; + @Override public Data get(String id, DataDefine dataDefine) { + return null; + } + + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + return null; + } + + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + source.put(SegmentTable.COLUMN_DATA_BINARY, new String(Base64.getEncoder().encode(data.getDataBytes(0)))); + logger.debug("segment source: {}", source.toString()); + return getClient().prepareIndex(SegmentTable.TABLE, data.getDataString(0)).setSource(source); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java index 5e9918fe8..363088b87 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java @@ -1,15 +1,9 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin.dao; -import java.util.List; -import java.util.Map; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; /** * @author pengys5 */ public class SegmentH2DAO extends H2DAO implements ISegmentDAO { - - @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/storage/IDNameExchangeTimer.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/IDNameExchangeTimer.java deleted file mode 100644 index c1880a589..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/IDNameExchangeTimer.java +++ /dev/null @@ -1,41 +0,0 @@ -package org.skywalking.apm.collector.agentstream.worker.storage; - -import java.util.List; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; -import org.skywalking.apm.collector.core.framework.Starter; -import org.skywalking.apm.collector.stream.worker.WorkerException; -import org.skywalking.apm.collector.stream.worker.impl.ExchangeWorker; -import org.skywalking.apm.collector.stream.worker.impl.ExchangeWorkerContainer; -import org.skywalking.apm.collector.stream.worker.impl.FlushAndSwitch; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * @author pengys5 - */ -public class IDNameExchangeTimer implements Starter { - - private final Logger logger = LoggerFactory.getLogger(IDNameExchangeTimer.class); - - public void start() { - logger.info("id and name exchange timer start"); - //TODO timer value config -// final long timeInterval = EsConfig.Es.Persistence.Timer.VALUE * 1000; - final long timeInterval = 3; - - Executors.newSingleThreadScheduledExecutor().schedule(() -> exchangeLastData(), timeInterval, TimeUnit.SECONDS); - } - - private void exchangeLastData() { - List workers = ExchangeWorkerContainer.INSTANCE.getExchangeWorkers(); - workers.forEach((ExchangeWorker worker) -> { - try { - worker.allocateJob(new FlushAndSwitch()); - worker.exchangeLastData(); - } catch (WorkerException e) { - logger.error(e.getMessage(), e); - } - }); - } -} 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 index ac02fee13..aeb38e243 100644 --- 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 @@ -26,7 +26,7 @@ public class PersistenceTimer implements Starter { //TODO timer value config // final long timeInterval = EsConfig.Es.Persistence.Timer.VALUE * 1000; final long timeInterval = 3; - Executors.newSingleThreadScheduledExecutor().schedule(() -> extractDataAndSave(), timeInterval, TimeUnit.SECONDS); + Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(() -> extractDataAndSave(), 1, timeInterval, TimeUnit.SECONDS); } private void extractDataAndSave() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/util/ExchangeMarkUtils.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/util/ExchangeMarkUtils.java new file mode 100644 index 000000000..84faf171e --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/util/ExchangeMarkUtils.java @@ -0,0 +1,14 @@ +package org.skywalking.apm.collector.agentstream.worker.util; + +/** + * @author pengys5 + */ +public enum ExchangeMarkUtils { + INSTANCE; + + private static final String MARK_TAG = "M"; + + public String buildMarkedID(int id) { + return MARK_TAG + id; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define deleted file mode 100644 index c8abe9d72..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define +++ /dev/null @@ -1,4 +0,0 @@ -org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine -org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationDataDefine -org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceDataDefine -org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameDataDefine \ 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_worker_provider.define similarity index 90% rename from apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define rename to apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define index a1744eee4..fb133fbe0 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_worker_provider.define @@ -2,7 +2,6 @@ org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentAggr org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentPersistenceWorker$Factory org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingAggregationWorker$Factory -org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeComponentExchangeWorker$Factory org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingPersistenceWorker$Factory org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefAggregationWorker$Factory 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 index fc988f521..3949d8a15 100644 --- 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 @@ -109,7 +109,7 @@ public class TraceSegmentServiceHandlerTestCase { ref_0.setParentServiceName(""); ref_0.setParentSpanId(2); ref_0.setParentTraceSegmentId(UniqueId.newBuilder().addIdParts(100).addIdParts(100).addIdParts(100).build()); - segmentBuilder.addRefs(ref_0); +// segmentBuilder.addRefs(ref_0); builder.setSegment(segmentBuilder.build().toByteString()); } 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 58cf145ef..1e4560af3 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 @@ -13,6 +13,7 @@ import org.elasticsearch.action.get.GetRequestBuilder; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.update.UpdateRequest; +import org.elasticsearch.action.update.UpdateRequestBuilder; import org.elasticsearch.client.IndicesAdminClient; import org.elasticsearch.common.settings.Settings; import org.elasticsearch.common.transport.InetSocketTransportAddress; @@ -111,9 +112,14 @@ public class ElasticSearchClient implements Client { } public IndexRequestBuilder prepareIndex(String indexName, String id) { + client.prepareUpdate(); return client.prepareIndex(indexName, "type", id); } + public UpdateRequestBuilder prepareUpdate(String indexName, String id) { + return client.prepareUpdate(indexName, "type", id); + } + public GetRequestBuilder prepareGet(String indexName, String id) { return client.prepareGet(indexName, "type", id); } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java index 0ae185794..3cd7559af 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/storage/StorageInstaller.java @@ -23,8 +23,11 @@ public abstract class StorageInstaller { if (!isExists(client, tableDefine)) { logger.info("table: {} not exists", tableDefine.getName()); tableDefine.initialize(); - createTable(client, tableDefine); + } else { + logger.info("table: {} exists", tableDefine.getName()); + deleteTable(client, tableDefine); } + createTable(client, tableDefine); } } catch (DefineException e) { throw new StorageInstallException(e.getMessage(), e); @@ -35,7 +38,7 @@ public abstract class StorageInstaller { protected abstract boolean isExists(Client client, TableDefine tableDefine) throws StorageException; - protected abstract boolean deleteIndex(Client client, TableDefine tableDefine) throws StorageException; + protected abstract boolean deleteTable(Client client, TableDefine tableDefine) throws StorageException; protected abstract boolean createTable(Client client, TableDefine tableDefine) throws StorageException; } diff --git a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java index 7b7c39bc5..7921df426 100644 --- a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java +++ b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java @@ -34,9 +34,16 @@ public abstract class JettyHandler extends HttpServlet implements Handler { @Override protected final void doPost(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException { - super.doPost(req, resp); + try { + doPost(req); + reply(resp); + } catch (ArgumentsParseException e) { + replyError(resp, e.getMessage(), HttpServletResponse.SC_BAD_REQUEST); + } } + protected abstract void doPost(HttpServletRequest req) throws ArgumentsParseException; + @Override protected final void doHead(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException { super.doHead(req, resp); @@ -134,13 +141,23 @@ public abstract class JettyHandler extends HttpServlet implements Handler { out.close(); } + private void reply(HttpServletResponse response) throws IOException { + response.setContentType("text/json"); + response.setCharacterEncoding("utf-8"); + response.setStatus(HttpServletResponse.SC_OK); + + PrintWriter out = response.getWriter(); + out.flush(); + out.close(); + } + private void replyError(HttpServletResponse response, String errorMessage, int status) throws IOException { response.setContentType("text/plain"); response.setCharacterEncoding("utf-8"); response.setStatus(status); + response.setHeader("error-message", errorMessage); PrintWriter out = response.getWriter(); - out.print(errorMessage); out.flush(); out.close(); } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java index 25dd33d61..26c618f13 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java @@ -81,7 +81,7 @@ public class ElasticSearchStorageInstaller extends StorageInstaller { return mappingBuilder; } - @Override protected boolean deleteIndex(Client client, TableDefine tableDefine) { + @Override protected boolean deleteTable(Client client, TableDefine tableDefine) { ElasticSearchClient esClient = (ElasticSearchClient)client; try { return esClient.deleteIndex(tableDefine.getName()); diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java index d27a8ca74..3244808a5 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java @@ -27,7 +27,7 @@ public class H2StorageInstaller extends StorageInstaller { return false; } - @Override protected boolean deleteIndex(Client client, TableDefine tableDefine) throws StorageException { + @Override protected boolean deleteTable(Client client, TableDefine tableDefine) throws StorageException { return false; } 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 908eecbb0..a1f22c2fa 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 @@ -16,8 +16,6 @@ import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.LocalAsyncWorkerProviderDefineLoader; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; import org.skywalking.apm.collector.stream.worker.RemoteWorkerProviderDefineLoader; -import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; -import org.skywalking.apm.collector.stream.worker.impl.data.DataDefineLoader; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -34,10 +32,6 @@ public class StreamModuleInstaller implements ModuleInstaller { StreamModuleContext context = new StreamModuleContext(StreamModuleGroupDefine.GROUP_NAME); CollectorContextHelper.INSTANCE.putContext(context); - DataDefineLoader dataDefineLoader = new DataDefineLoader(); - Map dataDefineMap = dataDefineLoader.load(); - context.putAllDataDefine(dataDefineMap); - initializeWorker(context); logger.info("could not configure cluster module, use the default"); diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorker.java index e43561bee..896a2ad51 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorker.java @@ -9,7 +9,9 @@ import org.skywalking.apm.collector.core.queue.QueueExecutor; * @author pengys5 * @since v3.0-2017 */ -public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker implements QueueExecutor { +public abstract class AbstractLocalAsyncWorker extends AbstractWorker implements QueueExecutor { + + private LocalAsyncWorkerRef workerRef; /** * Construct an AbstractLocalAsyncWorker with the worker role and context. @@ -32,6 +34,14 @@ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker imple public void preStart() throws ProviderNotFoundException { } + @Override protected final LocalAsyncWorkerRef getSelf() { + return workerRef; + } + + @Override protected final void putSelfRef(LocalAsyncWorkerRef workerRef) { + this.workerRef = workerRef; + } + /** * Receive message * 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 7f00cb13c..fe702b1b7 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,15 +6,13 @@ 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.ExchangeWorker; -import org.skywalking.apm.collector.stream.worker.impl.ExchangeWorkerContainer; import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorker; import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorkerContainer; /** * @author pengys5 */ -public abstract class AbstractLocalAsyncWorkerProvider extends AbstractLocalWorkerProvider { +public abstract class AbstractLocalAsyncWorkerProvider extends AbstractWorkerProvider { public abstract int queueSize(); @@ -24,8 +22,6 @@ public abstract class AbstractLocalAsyncWorkerProviderAbstractLocalSyncWorker defines workers who receive data from jvm inside call and response in real - * time. - * - * @author pengys5 - * @since v3.0-2017 - */ -public abstract class AbstractLocalSyncWorker extends AbstractLocalWorker { - public AbstractLocalSyncWorker(Role role, ClusterWorkerContext clusterContext) { - super(role, clusterContext); - } - - /** - * Called by the worker reference to execute the worker service. - * - * @param request {@link Object} is an input parameter - * @param response {@link Object} is an output parameter - */ - final public void allocateJob(Object request, Object response) throws WorkerInvokeException { - try { - onWork(request, response); - } catch (WorkerException e) { - throw new WorkerInvokeException(e.getMessage(), e.getCause()); - } - } - - /** - * Override this method to implement business logic. - * - * @param request {@link Object} is a in parameter - * @param response {@link Object} is a out parameter - */ - protected abstract void onWork(Object request, Object response) throws WorkerException; - - /** - * Prepare methods before this work starts to work. - *

Usually, create or find the workers reference should be call. - * - * @throws ProviderNotFoundException - */ - @Override - public void preStart() throws ProviderNotFoundException { - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorkerProvider.java deleted file mode 100644 index d3a542a61..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorkerProvider.java +++ /dev/null @@ -1,15 +0,0 @@ -package org.skywalking.apm.collector.stream.worker; - -/** - * @author pengys5 - */ -public abstract class AbstractLocalSyncWorkerProvider extends AbstractLocalWorkerProvider { - - @Override final public WorkerRef create() throws ProviderNotFoundException { - T localSyncWorker = workerInstance(getClusterContext()); - localSyncWorker.preStart(); - - LocalSyncWorkerRef workerRef = new LocalSyncWorkerRef(role(), localSyncWorker); - return workerRef; - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorker.java deleted file mode 100644 index c92a23fc7..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorker.java +++ /dev/null @@ -1,10 +0,0 @@ -package org.skywalking.apm.collector.stream.worker; - -/** - * @author pengys5 - */ -public abstract class AbstractLocalWorker extends AbstractWorker { - public AbstractLocalWorker(Role role, ClusterWorkerContext clusterContext) { - super(role, clusterContext); - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorkerProvider.java deleted file mode 100644 index 8cc96b71f..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorkerProvider.java +++ /dev/null @@ -1,7 +0,0 @@ -package org.skywalking.apm.collector.stream.worker; - -/** - * @author pengys5 - */ -public abstract class AbstractLocalWorkerProvider extends AbstractWorkerProvider { -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorker.java index 8dab25277..bba64efd2 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorker.java @@ -9,7 +9,9 @@ package org.skywalking.apm.collector.stream.worker; * @author pengys5 * @since v3.0-2017 */ -public abstract class AbstractRemoteWorker extends AbstractWorker { +public abstract class AbstractRemoteWorker extends AbstractWorker { + + private RemoteWorkerRef workerRef; /** * Construct an AbstractRemoteWorker with the worker role and context. @@ -35,4 +37,12 @@ public abstract class AbstractRemoteWorker extends AbstractWorker { throw new WorkerInvokeException(e.getMessage(), e.getCause()); } } + + @Override protected final RemoteWorkerRef getSelf() { + return workerRef; + } + + @Override protected final void putSelfRef(RemoteWorkerRef workerRef) { + this.workerRef = workerRef; + } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorker.java index 42d69415a..01918b5e4 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorker.java @@ -7,7 +7,7 @@ import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public abstract class AbstractWorker implements Executor { +public abstract class AbstractWorker implements Executor { private final Logger logger = LoggerFactory.getLogger(AbstractWorker.class); @@ -45,4 +45,8 @@ public abstract class AbstractWorker implements Executor { final public Role getRole() { return role; } + + protected abstract S getSelf(); + + protected abstract void putSelfRef(S workerRef); } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalAsyncWorkerProviderDefineLoader.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalAsyncWorkerProviderDefineLoader.java index 27fd2f1df..77ed5d4a5 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalAsyncWorkerProviderDefineLoader.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalAsyncWorkerProviderDefineLoader.java @@ -17,7 +17,7 @@ public class LocalAsyncWorkerProviderDefineLoader implements Loader load() throws DefineException { List providers = new ArrayList<>(); - LocalAsyncWorkerProviderDefinitionFile definitionFile = new LocalAsyncWorkerProviderDefinitionFile(); + LocalWorkerProviderDefinitionFile definitionFile = new LocalWorkerProviderDefinitionFile(); logger.info("local async worker provider definition file name: {}", definitionFile.fileName()); DefinitionLoader definitionLoader = DefinitionLoader.load(AbstractLocalAsyncWorkerProvider.class, definitionFile); diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalSyncWorkerRef.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalSyncWorkerRef.java deleted file mode 100644 index d24d07779..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalSyncWorkerRef.java +++ /dev/null @@ -1,23 +0,0 @@ -package org.skywalking.apm.collector.stream.worker; - -/** - * @author pengys5 - */ -public class LocalSyncWorkerRef extends WorkerRef { - - private AbstractLocalSyncWorker localSyncWorker; - - public LocalSyncWorkerRef(Role role, AbstractLocalSyncWorker localSyncWorker) { - super(role); - this.localSyncWorker = localSyncWorker; - } - - @Override - public void tell(Object message) throws WorkerInvokeException { - localSyncWorker.allocateJob(message, null); - } - - public void ask(Object request, Object response) throws WorkerInvokeException { - localSyncWorker.allocateJob(request, response); - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalAsyncWorkerProviderDefinitionFile.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalWorkerProviderDefinitionFile.java similarity index 60% rename from apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalAsyncWorkerProviderDefinitionFile.java rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalWorkerProviderDefinitionFile.java index 356cf9ffb..5dbfa09f7 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalAsyncWorkerProviderDefinitionFile.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/LocalWorkerProviderDefinitionFile.java @@ -5,8 +5,8 @@ import org.skywalking.apm.collector.core.framework.DefinitionFile; /** * @author pengys5 */ -public class LocalAsyncWorkerProviderDefinitionFile extends DefinitionFile { +public class LocalWorkerProviderDefinitionFile extends DefinitionFile { @Override protected String fileName() { - return "local_async_worker_provider.define"; + return "local_worker_provider.define"; } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/ExchangeWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/ExchangeWorker.java deleted file mode 100644 index 3a2e8a098..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/ExchangeWorker.java +++ /dev/null @@ -1,74 +0,0 @@ -package org.skywalking.apm.collector.stream.worker.impl; - -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 ExchangeWorker extends AbstractLocalAsyncWorker { - - private final Logger logger = LoggerFactory.getLogger(ExchangeWorker.class); - - private DataCache dataCache; - - public ExchangeWorker(Role role, ClusterWorkerContext clusterContext) { - super(role, clusterContext); - dataCache = new DataCache(); - } - - @Override public void preStart() throws ProviderNotFoundException { - super.preStart(); - } - - @Override protected final void onWork(Object message) throws WorkerException { - if (message instanceof FlushAndSwitch) { - if (dataCache.trySwitchPointer()) { - dataCache.switchPointer(); - } - } else if (message instanceof EndOfBatchCommand) { - } else { - if (dataCache.currentCollectionSize() <= 1000) { - aggregate(message); - } - } - } - - protected abstract void exchange(Data data); - - public final void exchangeLastData() { - try { - while (dataCache.getLast().isHolding()) { - try { - Thread.sleep(10); - } catch (InterruptedException e) { - logger.warn("thread wake up"); - } - } - dataCache.getLast().asMap().values().forEach(data -> { - exchange(data); - }); - } finally { - dataCache.releaseLast(); - } - } - - protected final void aggregate(Object message) { - Data data = (Data)message; - dataCache.hold(); - if (dataCache.containsKey(data.id())) { - getRole().dataDefine().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/ExchangeWorkerContainer.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/ExchangeWorkerContainer.java deleted file mode 100644 index 38ffea545..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/ExchangeWorkerContainer.java +++ /dev/null @@ -1,21 +0,0 @@ -package org.skywalking.apm.collector.stream.worker.impl; - -import java.util.ArrayList; -import java.util.List; - -/** - * @author pengys5 - */ -public enum ExchangeWorkerContainer { - INSTANCE; - - private List exchangeWorkers = new ArrayList<>(); - - public void addWorker(ExchangeWorker worker) { - exchangeWorkers.add(worker); - } - - public List getExchangeWorkers() { - return exchangeWorkers; - } -} 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 index f0f3bae3b..878c2c6e0 100644 --- 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 @@ -1,13 +1,16 @@ package org.skywalking.apm.collector.stream.worker.impl; +import java.util.ArrayList; import java.util.List; import java.util.Map; import org.skywalking.apm.collector.core.queue.EndOfBatchCommand; +import org.skywalking.apm.collector.core.util.ObjectUtils; 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.dao.IPersistenceDAO; import org.skywalking.apm.collector.stream.worker.impl.data.Data; import org.skywalking.apm.collector.stream.worker.impl.data.DataCache; import org.slf4j.Logger; @@ -46,7 +49,7 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { } } - public List buildBatchCollection() throws WorkerException { + public final List buildBatchCollection() throws WorkerException { List batchCollection; try { while (dataCache.getLast().isHolding()) { @@ -56,6 +59,7 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { logger.warn("thread wake up"); } } + batchCollection = prepareBatch(dataCache.getLast().asMap()); } finally { dataCache.releaseLast(); @@ -63,7 +67,26 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { return batchCollection; } - protected abstract List prepareBatch(Map dataMap); + protected final List prepareBatch(Map dataMap) { + List insertBatchCollection = new ArrayList<>(); + List updateBatchCollection = new ArrayList<>(); + dataMap.forEach((id, data) -> { + if (needMergeDBData()) { + Data dbData = persistenceDAO().get(id, getRole().dataDefine()); + if (ObjectUtils.isNotEmpty(dbData)) { + getRole().dataDefine().mergeData(data, dbData); + updateBatchCollection.add(persistenceDAO().prepareBatchUpdate(data)); + } else { + insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + } + } else { + insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data)); + } + }); + + insertBatchCollection.addAll(updateBatchCollection); + return insertBatchCollection; + } private void aggregate(Object message) { dataCache.hold(); @@ -79,4 +102,8 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { dataCache.release(); } + + protected abstract IPersistenceDAO persistenceDAO(); + + protected abstract boolean needMergeDBData(); } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/dao/IPersistenceDAO.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/dao/IPersistenceDAO.java new file mode 100644 index 000000000..956586265 --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/dao/IPersistenceDAO.java @@ -0,0 +1,15 @@ +package org.skywalking.apm.collector.stream.worker.impl.dao; + +import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; + +/** + * @author pengys5 + */ +public interface IPersistenceDAO { + Data get(String id, DataDefine dataDefine); + + I prepareBatchInsert(Data data); + + U prepareBatchUpdate(Data data); +} 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 deleted file mode 100644 index 425fcb427..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefineLoader.java +++ /dev/null @@ -1,29 +0,0 @@ -package org.skywalking.apm.collector.stream.worker.impl.data; - -import java.util.HashMap; -import java.util.Map; -import org.skywalking.apm.collector.core.framework.DefineException; -import org.skywalking.apm.collector.core.framework.Loader; -import org.skywalking.apm.collector.core.util.DefinitionLoader; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * @author pengys5 - */ -public class DataDefineLoader implements Loader> { - - private final Logger logger = LoggerFactory.getLogger(DataDefineLoader.class); - - @Override public Map load() throws DefineException { - Map dataDefineMap = new HashMap<>(); - - DataDefinitionFile definitionFile = new DataDefinitionFile(); - DefinitionLoader definitionLoader = DefinitionLoader.load(DataDefine.class, definitionFile); - for (DataDefine dataDefine : definitionLoader) { - logger.info("loaded data definition class: {}", dataDefine.getClass().getName()); - dataDefineMap.put(dataDefine.defineId(), dataDefine); - } - return dataDefineMap; - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefinitionFile.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefinitionFile.java deleted file mode 100644 index e7875d8b8..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefinitionFile.java +++ /dev/null @@ -1,13 +0,0 @@ -package org.skywalking.apm.collector.stream.worker.impl.data; - -import org.skywalking.apm.collector.core.framework.DefinitionFile; - -/** - * @author pengys5 - */ -public class DataDefinitionFile extends DefinitionFile { - - @Override protected String fileName() { - return "data.define"; - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Exchange.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Exchange.java deleted file mode 100644 index 3d37d4500..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Exchange.java +++ /dev/null @@ -1,24 +0,0 @@ -package org.skywalking.apm.collector.stream.worker.impl.data; - -/** - * @author pengys5 - */ -public abstract class Exchange { - private int times; - - public Exchange(int times) { - this.times = times; - } - - public void increase() { - times++; - } - - public int getTimes() { - return times; - } - - public void setTimes(int times) { - this.times = times; - } -} diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java index c06df8079..4e1d05699 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/SegmentCostEsDAO.java @@ -45,7 +45,7 @@ public class SegmentCostEsDAO extends EsDAO implements ISegmentCostDAO { boolQueryBuilder.must().add(rangeQueryBuilder); } if (!StringUtils.isEmpty(operationName)) { - mustQueryList.add(QueryBuilders.matchQuery(SegmentCostTable.COLUMN_OPERATION_NAME, operationName)); + mustQueryList.add(QueryBuilders.matchQuery(SegmentCostTable.COLUMN_SERVICE_NAME, operationName)); } searchRequestBuilder.addSort(SegmentCostTable.COLUMN_COST, SortOrder.DESC); @@ -77,7 +77,7 @@ public class SegmentCostEsDAO extends EsDAO implements ISegmentCostDAO { topSegmentJson.addProperty(GlobalTraceTable.COLUMN_GLOBAL_TRACE_ID, globalTraces.get(0)); } - topSegmentJson.addProperty(SegmentCostTable.COLUMN_OPERATION_NAME, (String)searchHit.getSource().get(SegmentCostTable.COLUMN_OPERATION_NAME)); + topSegmentJson.addProperty(SegmentCostTable.COLUMN_SERVICE_NAME, (String)searchHit.getSource().get(SegmentCostTable.COLUMN_SERVICE_NAME)); topSegmentJson.addProperty(SegmentCostTable.COLUMN_COST, (Number)searchHit.getSource().get(SegmentCostTable.COLUMN_COST)); topSegmentJson.addProperty(SegmentCostTable.COLUMN_IS_ERROR, (Boolean)searchHit.getSource().get(SegmentCostTable.COLUMN_IS_ERROR));