From dcad78c45c4f310f9a22496d62b0fc1c75e1f74f Mon Sep 17 00:00:00 2001 From: peng-yongsheng <8082209@qq.com> Date: Mon, 18 Dec 2017 19:11:45 +0800 Subject: [PATCH] Analysis segment parser module finished. --- .../stream/AgentStreamModuleProvider.java | 2 +- .../parser/define/graph/GraphIdDefine.java | 27 ++++++++++++ .../parser/define/graph/WorkerIdDefine.java | 27 ++++++++++++ .../segment-parser-provider/pom.xml | 12 +++++- .../AnalysisSegmentParserModuleProvider.java | 8 ++++ .../provider}/buffer/BufferFileConfig.java | 2 +- .../parser/provider}/buffer/Offset.java | 2 +- .../provider}/buffer/OffsetManager.java | 5 +-- .../buffer/SegmentBufferManager.java | 2 +- .../provider}/buffer/SegmentBufferReader.java | 16 +++++--- .../parser/provider/parser/SegmentParse.java | 7 ++-- .../parser/SegmentPersistenceGraph.java | 41 +++++++++++++++++++ .../parser}/SegmentPersistenceWorker.java | 15 ++++--- .../SegmentStandardizationGraph.java | 40 ++++++++++++++++++ .../SegmentStandardizationWorker.java | 13 +++--- .../AbstractLocalAsyncWorkerProvider.java | 2 +- .../collector/core/graph/GraphManager.java | 3 +- .../apm/collector/core}/util/FileUtils.java | 3 +- 18 files changed, 191 insertions(+), 36 deletions(-) create mode 100644 apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/GraphIdDefine.java create mode 100644 apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/WorkerIdDefine.java rename apm-collector/{apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream => apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider}/buffer/BufferFileConfig.java (97%) rename apm-collector/{apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream => apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider}/buffer/Offset.java (97%) rename apm-collector/{apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream => apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider}/buffer/OffsetManager.java (97%) rename apm-collector/{apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream => apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider}/buffer/SegmentBufferManager.java (97%) rename apm-collector/{apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream => apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider}/buffer/SegmentBufferReader.java (89%) create mode 100644 apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentPersistenceGraph.java rename apm-collector/{apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/worker/trace/segment => apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser}/SegmentPersistenceWorker.java (87%) create mode 100644 apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationGraph.java rename apm-collector/{apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream => apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core}/util/FileUtils.java (96%) diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/AgentStreamModuleProvider.java b/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/AgentStreamModuleProvider.java index 9577b7b14..205f85afb 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/AgentStreamModuleProvider.java +++ b/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/AgentStreamModuleProvider.java @@ -20,7 +20,7 @@ package org.apache.skywalking.apm.collector.agent.stream; import java.util.Properties; -import org.apache.skywalking.apm.collector.agent.stream.buffer.BufferFileConfig; +import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer.BufferFileConfig; import org.apache.skywalking.apm.collector.agent.stream.service.jvm.IGCMetricService; import org.apache.skywalking.apm.collector.agent.stream.service.jvm.IInstanceHeartBeatService; import org.apache.skywalking.apm.collector.agent.stream.service.jvm.IMemoryMetricService; diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/GraphIdDefine.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/GraphIdDefine.java new file mode 100644 index 000000000..3229791b5 --- /dev/null +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/GraphIdDefine.java @@ -0,0 +1,27 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.analysis.segment.parser.define.graph; + +/** + * @author peng-yongsheng + */ +public class GraphIdDefine { + public static final int SEGMENT_PERSISTENCE_GRAPH_ID = 100; + public static final int SEGMENT_STANDARDIZATION_GRAPH_ID = 101; +} diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/WorkerIdDefine.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/WorkerIdDefine.java new file mode 100644 index 000000000..8d082fe31 --- /dev/null +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-define/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/define/graph/WorkerIdDefine.java @@ -0,0 +1,27 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.analysis.segment.parser.define.graph; + +/** + * @author peng-yongsheng + */ +public class WorkerIdDefine { + public static final int SEGMENT_PERSISTENCE_WORKER_ID = 100; + public static final int SEGMENT_STANDARDIZATION_WORKER_ID = 101; +} diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/pom.xml b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/pom.xml index 9eef1371c..aee7752e0 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/pom.xml +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/pom.xml @@ -38,7 +38,7 @@ org.apache.skywalking - layer-register-define + register-define ${project.version} @@ -46,5 +46,15 @@ collector-cache-define ${project.version} + + org.apache.skywalking + collector-storage-define + ${project.version} + + + org.apache.skywalking + analysis-worker-model + ${project.version} + \ No newline at end of file diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/AnalysisSegmentParserModuleProvider.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/AnalysisSegmentParserModuleProvider.java index 1ea256faa..63d0c476d 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/AnalysisSegmentParserModuleProvider.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/AnalysisSegmentParserModuleProvider.java @@ -22,7 +22,10 @@ import java.util.Properties; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.AnalysisSegmentParserModule; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.service.ISegmentParseService; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.service.ISegmentParserListenerRegister; +import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer.SegmentBufferReader; import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser.SegmentParserListenerManager; +import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser.SegmentPersistenceGraph; +import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser.standardization.SegmentStandardizationGraph; import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.service.SegmentParseService; import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.service.SegmentParserListenerRegister; import org.apache.skywalking.apm.collector.core.module.Module; @@ -52,7 +55,12 @@ public class AnalysisSegmentParserModuleProvider extends ModuleProvider { } @Override public void start(Properties config) throws ServiceNotProvidedException { + SegmentPersistenceGraph segmentPersistenceGraph = new SegmentPersistenceGraph(getManager()); + segmentPersistenceGraph.create(); + SegmentStandardizationGraph segmentStandardizationGraph = new SegmentStandardizationGraph(getManager()); + segmentStandardizationGraph.create(); + SegmentBufferReader.INSTANCE.setSegmentParserListenerManager(listenerManager); } @Override public void notifyAfterCompleted() throws ServiceNotProvidedException { diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/BufferFileConfig.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/BufferFileConfig.java similarity index 97% rename from apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/BufferFileConfig.java rename to apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/BufferFileConfig.java index ef84cec62..8ef63a8c3 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/BufferFileConfig.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/BufferFileConfig.java @@ -17,7 +17,7 @@ */ -package org.apache.skywalking.apm.collector.agent.stream.buffer; +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer; import java.util.Properties; diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/Offset.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/Offset.java similarity index 97% rename from apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/Offset.java rename to apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/Offset.java index ec65ad408..7302690a3 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/Offset.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/Offset.java @@ -17,7 +17,7 @@ */ -package org.apache.skywalking.apm.collector.agent.stream.buffer; +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer; /** * @author peng-yongsheng diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/OffsetManager.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/OffsetManager.java similarity index 97% rename from apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/OffsetManager.java rename to apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/OffsetManager.java index e9d8a2ce1..12899710e 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/OffsetManager.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/OffsetManager.java @@ -16,8 +16,7 @@ * */ - -package org.apache.skywalking.apm.collector.agent.stream.buffer; +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer; import java.io.File; import java.io.FilenameFilter; @@ -27,7 +26,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import org.apache.skywalking.apm.collector.core.util.CollectionUtils; import org.apache.skywalking.apm.collector.core.util.Const; -import org.apache.skywalking.apm.collector.agent.stream.util.FileUtils; +import org.apache.skywalking.apm.collector.core.util.FileUtils; import org.apache.skywalking.apm.collector.core.util.TimeBucketUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/SegmentBufferManager.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferManager.java similarity index 97% rename from apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/SegmentBufferManager.java rename to apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferManager.java index 87edff5f6..c146f6653 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/SegmentBufferManager.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferManager.java @@ -17,7 +17,7 @@ */ -package org.apache.skywalking.apm.collector.agent.stream.buffer; +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer; import java.io.File; import java.io.FileOutputStream; diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/SegmentBufferReader.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferReader.java similarity index 89% rename from apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/SegmentBufferReader.java rename to apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferReader.java index 9ec4db053..d2e41f4e5 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/buffer/SegmentBufferReader.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferReader.java @@ -16,8 +16,7 @@ * */ - -package org.apache.skywalking.apm.collector.agent.stream.buffer; +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer; import com.google.protobuf.CodedOutputStream; import java.io.File; @@ -27,7 +26,9 @@ import java.io.IOException; import java.io.InputStream; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; -import org.apache.skywalking.apm.collector.agent.stream.parser.SegmentParse; +import org.apache.skywalking.apm.collector.analysis.segment.parser.define.service.ISegmentParseService; +import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser.SegmentParse; +import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser.SegmentParserListenerManager; import org.apache.skywalking.apm.collector.core.module.ModuleManager; import org.apache.skywalking.apm.collector.core.util.CollectionUtils; import org.apache.skywalking.apm.collector.core.util.Const; @@ -45,12 +46,17 @@ public enum SegmentBufferReader { private final Logger logger = LoggerFactory.getLogger(SegmentBufferReader.class); private InputStream inputStream; private ModuleManager moduleManager; + private SegmentParserListenerManager listenerManager; public void initialize(ModuleManager moduleManager) { this.moduleManager = moduleManager; Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(this::preRead, 3, 3, TimeUnit.SECONDS); } + public void setSegmentParserListenerManager(SegmentParserListenerManager listenerManager) { + this.listenerManager = listenerManager; + } + private void preRead() { String readFileName = OffsetManager.INSTANCE.getReadFileName(); if (StringUtils.isNotEmpty(readFileName)) { @@ -121,8 +127,8 @@ public enum SegmentBufferReader { while (readFile.length() > readFileOffset && readFileOffset < endPoint) { UpstreamSegment upstreamSegment = UpstreamSegment.parser().parseDelimitedFrom(inputStream); - SegmentParse parse = new SegmentParse(moduleManager); - if (!parse.parse(upstreamSegment, SegmentParse.Source.Buffer)) { + SegmentParse parse = new SegmentParse(moduleManager, listenerManager); + if (!parse.parse(upstreamSegment, ISegmentParseService.Source.Buffer)) { return false; } diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentParse.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentParse.java index 60998e3c7..08030f283 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentParse.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentParse.java @@ -20,10 +20,10 @@ package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.par import com.google.protobuf.InvalidProtocolBufferException; import java.util.List; -import javax.swing.text.Segment; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.decorator.ReferenceDecorator; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.decorator.SegmentDecorator; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.decorator.SpanDecorator; +import org.apache.skywalking.apm.collector.analysis.segment.parser.define.graph.GraphIdDefine; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.listener.EntrySpanListener; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.listener.ExitSpanListener; import org.apache.skywalking.apm.collector.analysis.segment.parser.define.listener.FirstSpanListener; @@ -39,6 +39,7 @@ import org.apache.skywalking.apm.collector.core.graph.Graph; import org.apache.skywalking.apm.collector.core.graph.GraphManager; import org.apache.skywalking.apm.collector.core.module.ModuleManager; import org.apache.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.apache.skywalking.apm.collector.storage.table.segment.Segment; import org.apache.skywalking.apm.network.proto.SpanType; import org.apache.skywalking.apm.network.proto.TraceSegmentObject; import org.apache.skywalking.apm.network.proto.UniqueId; @@ -159,7 +160,7 @@ public class SegmentParse { Segment segment = new Segment(id); segment.setDataBinary(dataBinary); segment.setTimeBucket(timeBucket); - Graph graph = GraphManager.INSTANCE.createIfAbsent(TraceStreamGraph.SEGMENT_GRAPH_ID, Segment.class); + Graph graph = GraphManager.INSTANCE.findGraph(GraphIdDefine.SEGMENT_PERSISTENCE_GRAPH_ID, Segment.class); graph.start(segment); } @@ -167,7 +168,7 @@ public class SegmentParse { logger.debug("push to segment buffer write worker, id: {}", id); SegmentStandardization standardization = new SegmentStandardization(id); standardization.setUpstreamSegment(upstreamSegment); - Graph graph = GraphManager.INSTANCE.createIfAbsent(TraceStreamGraph.SEGMENT_STANDARDIZATION_GRAPH_ID, SegmentStandardization.class); + Graph graph = GraphManager.INSTANCE.findGraph(GraphIdDefine.SEGMENT_STANDARDIZATION_GRAPH_ID, SegmentStandardization.class); graph.start(standardization); } diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentPersistenceGraph.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentPersistenceGraph.java new file mode 100644 index 000000000..cc6ec6280 --- /dev/null +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentPersistenceGraph.java @@ -0,0 +1,41 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser; + +import org.apache.skywalking.apm.collector.analysis.segment.parser.define.graph.GraphIdDefine; +import org.apache.skywalking.apm.collector.core.graph.GraphManager; +import org.apache.skywalking.apm.collector.core.module.ModuleManager; +import org.apache.skywalking.apm.collector.storage.table.segment.Segment; + +/** + * @author peng-yongsheng + */ +public class SegmentPersistenceGraph { + + private final ModuleManager moduleManager; + + public SegmentPersistenceGraph(ModuleManager moduleManager) { + this.moduleManager = moduleManager; + } + + public void create() { + GraphManager.INSTANCE.createIfAbsent(GraphIdDefine.SEGMENT_PERSISTENCE_GRAPH_ID, Segment.class) + .addNode(new SegmentPersistenceWorker.Factory(moduleManager).create(null)); + } +} diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/worker/trace/segment/SegmentPersistenceWorker.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentPersistenceWorker.java similarity index 87% rename from apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/worker/trace/segment/SegmentPersistenceWorker.java rename to apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentPersistenceWorker.java index a28ad425a..18949ac9e 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/worker/trace/segment/SegmentPersistenceWorker.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/SegmentPersistenceWorker.java @@ -16,17 +16,16 @@ * */ +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser; -package org.apache.skywalking.apm.collector.agent.stream.worker.trace.segment; - +import org.apache.skywalking.apm.collector.analysis.segment.parser.define.graph.WorkerIdDefine; +import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorkerProvider; +import org.apache.skywalking.apm.collector.analysis.worker.model.impl.PersistenceWorker; import org.apache.skywalking.apm.collector.core.module.ModuleManager; -import org.apache.skywalking.apm.collector.queue.service.QueueCreatorService; import org.apache.skywalking.apm.collector.storage.StorageModule; import org.apache.skywalking.apm.collector.storage.base.dao.IPersistenceDAO; import org.apache.skywalking.apm.collector.storage.dao.ISegmentPersistenceDAO; import org.apache.skywalking.apm.collector.storage.table.segment.Segment; -import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorkerProvider; -import org.apache.skywalking.apm.collector.analysis.worker.model.impl.PersistenceWorker; /** * @author peng-yongsheng @@ -38,7 +37,7 @@ public class SegmentPersistenceWorker extends PersistenceWorker { - public Factory(ModuleManager moduleManager, QueueCreatorService queueCreatorService) { - super(moduleManager, queueCreatorService); + public Factory(ModuleManager moduleManager) { + super(moduleManager); } @Override public SegmentPersistenceWorker workerInstance(ModuleManager moduleManager) { diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationGraph.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationGraph.java new file mode 100644 index 000000000..104c2ee90 --- /dev/null +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationGraph.java @@ -0,0 +1,40 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser.standardization; + +import org.apache.skywalking.apm.collector.analysis.segment.parser.define.graph.GraphIdDefine; +import org.apache.skywalking.apm.collector.core.graph.GraphManager; +import org.apache.skywalking.apm.collector.core.module.ModuleManager; + +/** + * @author peng-yongsheng + */ +public class SegmentStandardizationGraph { + + private final ModuleManager moduleManager; + + public SegmentStandardizationGraph(ModuleManager moduleManager) { + this.moduleManager = moduleManager; + } + + public void create() { + GraphManager.INSTANCE.createIfAbsent(GraphIdDefine.SEGMENT_STANDARDIZATION_GRAPH_ID, SegmentStandardization.class) + .addNode(new SegmentStandardizationWorker.Factory(moduleManager).create(null)); + } +} diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationWorker.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationWorker.java index a5c81bc92..93d0c2ef0 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationWorker.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SegmentStandardizationWorker.java @@ -16,17 +16,16 @@ * */ - package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.parser.standardization; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import org.apache.skywalking.apm.collector.analysis.segment.parser.define.graph.WorkerIdDefine; +import org.apache.skywalking.apm.collector.analysis.segment.parser.provider.buffer.SegmentBufferManager; import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorker; -import org.apache.skywalking.apm.collector.agent.stream.buffer.SegmentBufferManager; -import org.apache.skywalking.apm.collector.core.module.ModuleManager; -import org.apache.skywalking.apm.collector.queue.service.QueueCreatorService; import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorkerProvider; import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException; +import org.apache.skywalking.apm.collector.core.module.ModuleManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -43,7 +42,7 @@ public class SegmentStandardizationWorker extends AbstractLocalAsyncWorker { - public Factory(ModuleManager moduleManager, QueueCreatorService queueCreatorService) { - super(moduleManager, queueCreatorService); + public Factory(ModuleManager moduleManager) { + super(moduleManager); } @Override public SegmentStandardizationWorker workerInstance(ModuleManager moduleManager) { diff --git a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java index 102bed5db..78fd2f1bc 100644 --- a/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java +++ b/apm-collector/apm-collector-analysis/analysis-worker-model/src/main/java/org/apache/skywalking/apm/collector/analysis/worker/model/base/AbstractLocalAsyncWorkerProvider.java @@ -33,7 +33,7 @@ public abstract class AbstractLocalAsyncWorkerProvider create(WorkerCreateListener workerCreateListener) { WORKER_TYPE localAsyncWorker = workerInstance(getModuleManager()); workerCreateListener.addWorker(localAsyncWorker); diff --git a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/graph/GraphManager.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/graph/GraphManager.java index 464cee000..cdd7dac36 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/graph/GraphManager.java +++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/graph/GraphManager.java @@ -16,7 +16,6 @@ * */ - package org.apache.skywalking.apm.collector.core.graph; import java.util.HashMap; @@ -46,7 +45,7 @@ public enum GraphManager { } } - public Graph findGraph(int graphId) { + public Graph findGraph(int graphId, Class input) { Graph graph = allGraphs.get(graphId); if (graph == null) { throw new GraphNotFoundException("Graph id=" + graphId + " not found in this GraphManager"); diff --git a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/util/FileUtils.java b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/util/FileUtils.java similarity index 96% rename from apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/util/FileUtils.java rename to apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/util/FileUtils.java index 9a5aa5027..6036b7fca 100644 --- a/apm-collector/apm-collector-agent-stream/collector-agent-stream-provider/src/main/java/org/apache/skywalking/apm/collector/agent/stream/util/FileUtils.java +++ b/apm-collector/apm-collector-core/src/main/java/org/apache/skywalking/apm/collector/core/util/FileUtils.java @@ -17,13 +17,12 @@ */ -package org.apache.skywalking.apm.collector.agent.stream.util; +package org.apache.skywalking.apm.collector.core.util; import java.io.File; import java.io.FileNotFoundException; import java.io.IOException; import java.io.RandomAccessFile; -import org.apache.skywalking.apm.collector.core.util.Const; import org.slf4j.Logger; import org.slf4j.LoggerFactory;