From a9ff88fe958900e555b75591f06b9dc5af4f6cc1 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 3 Sep 2017 16:19:44 +0800 Subject: [PATCH 1/2] Performance evaluation test result: Because of laptop just have 100M network card and sata disk, so 100 million segment process by collector used 255 second, but cpu just used 28% --- .../instance/InstanceIDService.java | 4 +- .../worker/cache/ApplicationCache.java | 16 +- .../worker/cache/InstanceCache.java | 2 +- .../worker/cache/ServiceCache.java | 2 +- .../define/GlobalTraceEsTableDefine.java | 2 +- .../ServiceNameRegisterSerialWorker.java | 4 +- .../servicename/dao/ServiceNameEsDAO.java | 2 +- .../cost/define/SegmentCostEsTableDefine.java | 2 +- .../mock/grpc/GrpcSegmentPost.java | 473 ++++++++++++++++++ .../src/test/resources/logback.xml | 16 + .../src/main/resources/application.yml | 2 +- .../collector/core/server/ServerHolder.java | 9 +- .../src/main/proto/RemoteCommonService.proto | 2 +- .../handler/RemoteCommonServiceHandler.java | 34 +- .../stream/worker/RemoteWorkerRef.java | 74 ++- .../stream/worker/impl/AggregationWorker.java | 8 +- .../stream/worker/impl/PersistenceWorker.java | 51 +- .../stream/worker/impl/data/DataCache.java | 10 +- .../worker/impl/data/DataCollection.java | 34 +- .../stream/worker/impl/data/Window.java | 22 +- 20 files changed, 696 insertions(+), 73 deletions(-) create mode 100644 apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/resources/logback.xml diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java index fa5f5bb22..40743ab21 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java @@ -1,6 +1,6 @@ package org.skywalking.apm.collector.agentregister.instance; -import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterRemoteWorker; +import org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterRemoteWorker; import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.IInstanceDAO; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -28,7 +28,7 @@ public class InstanceIDService { StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance("0", applicationId, agentUUID, registerTime, 0, registerTime, osInfo); try { - context.getClusterWorkerContext().lookup(ApplicationRegisterRemoteWorker.WorkerRole.INSTANCE).tell(instance); + context.getClusterWorkerContext().lookup(InstanceRegisterRemoteWorker.WorkerRole.INSTANCE).tell(instance); } catch (WorkerNotFoundException | WorkerInvokeException e) { logger.error(e.getMessage(), e); } 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 index ff5574029..a38cb3976 100644 --- 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 @@ -10,16 +10,26 @@ import org.skywalking.apm.collector.storage.dao.DAOContainer; */ public class ApplicationCache { - private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(1000).build(); + private static Cache CACHE = CacheBuilder.newBuilder().initialCapacity(100).maximumSize(1000).build(); public static int get(String applicationCode) { + int applicationId = 0; try { - return CACHE.get(applicationCode, () -> { + applicationId = CACHE.get(applicationCode, () -> { IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName()); return dao.getApplicationId(applicationCode); }); } catch (Throwable e) { - return 0; + return applicationId; } + + if (applicationId == 0) { + IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName()); + applicationId = dao.getApplicationId(applicationCode); + if (applicationId != 0) { + CACHE.put(applicationCode, applicationId); + } + } + return applicationId; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java index 54ac23020..e2808ef6f 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java @@ -10,7 +10,7 @@ import org.skywalking.apm.collector.storage.dao.DAOContainer; */ public class InstanceCache { - private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(1000).build(); + private static Cache CACHE = CacheBuilder.newBuilder().initialCapacity(100).maximumSize(5000).build(); public static int get(int applicationInstanceId) { try { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java index 6d0294bd7..1dcc065b0 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java @@ -11,7 +11,7 @@ import org.skywalking.apm.collector.storage.dao.DAOContainer; */ public class ServiceCache { - private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(10000).build(); + private static Cache CACHE = CacheBuilder.newBuilder().initialCapacity(1000).maximumSize(20000).build(); public static String getServiceName(int serviceId) { try { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java index 657606ca8..8a2b244a2 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java @@ -14,7 +14,7 @@ public class GlobalTraceEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 2; + return 5; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java index 7b98a3556..5d6fd9efd 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.servicename; import org.skywalking.apm.collector.agentstream.worker.register.IdAutoIncrement; import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO; import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.register.ServiceNameDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorker; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; @@ -10,7 +11,6 @@ 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.storage.define.DataDefine; import org.skywalking.apm.collector.stream.worker.selector.ForeverFirstSelector; import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; import org.slf4j.Logger; @@ -47,8 +47,8 @@ public class ServiceNameRegisterSerialWorker extends AbstractLocalAsyncWorker { } else { int max = dao.getMaxServiceId(); serviceId = IdAutoIncrement.INSTANCE.increment(min, max); - serviceName.setApplicationId(serviceId); serviceName.setId(String.valueOf(serviceId)); + serviceName.setServiceId(serviceId); } dao.save(serviceName); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java index 150a0b601..a2e3c63d8 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java @@ -30,7 +30,7 @@ public class ServiceNameEsDAO extends EsDAO implements IServiceNameDAO { ElasticSearchClient client = getClient(); SearchRequestBuilder searchRequestBuilder = client.prepareSearch(ServiceNameTable.TABLE); - searchRequestBuilder.setTypes("type"); + searchRequestBuilder.setTypes(ServiceNameTable.TABLE_TYPE); searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH); BoolQueryBuilder builder = QueryBuilders.boolQuery(); builder.must().add(QueryBuilders.termQuery(ServiceNameTable.COLUMN_APPLICATION_ID, applicationId)); 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 a25bf0fba..82f05c550 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 @@ -14,7 +14,7 @@ public class SegmentCostEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 2; + return 5; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java new file mode 100644 index 000000000..3672c633b --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/grpc/GrpcSegmentPost.java @@ -0,0 +1,473 @@ +package org.skywalking.apm.collector.agentstream.mock.grpc; + +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import io.grpc.stub.StreamObserver; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import org.junit.Test; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.network.proto.Application; +import org.skywalking.apm.network.proto.ApplicationInstance; +import org.skywalking.apm.network.proto.ApplicationInstanceMapping; +import org.skywalking.apm.network.proto.ApplicationMapping; +import org.skywalking.apm.network.proto.ApplicationRegisterServiceGrpc; +import org.skywalking.apm.network.proto.Downstream; +import org.skywalking.apm.network.proto.InstanceDiscoveryServiceGrpc; +import org.skywalking.apm.network.proto.KeyWithStringValue; +import org.skywalking.apm.network.proto.LogMessage; +import org.skywalking.apm.network.proto.OSInfo; +import org.skywalking.apm.network.proto.RefType; +import org.skywalking.apm.network.proto.ServiceNameCollection; +import org.skywalking.apm.network.proto.ServiceNameDiscoveryServiceGrpc; +import org.skywalking.apm.network.proto.ServiceNameElement; +import org.skywalking.apm.network.proto.ServiceNameMappingCollection; +import org.skywalking.apm.network.proto.SpanLayer; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.SpanType; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.skywalking.apm.network.proto.TraceSegmentServiceGrpc; +import org.skywalking.apm.network.proto.UniqueId; +import org.skywalking.apm.network.proto.UpstreamSegment; +import org.skywalking.apm.network.trace.component.ComponentsDefine; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class GrpcSegmentPost { + + private final Logger logger = LoggerFactory.getLogger(GrpcSegmentPost.class); + + private AtomicLong sequence = new AtomicLong(1); + + @Test + public void init() { + ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).maxInboundMessageSize(1024 * 1024 * 50).usePlaintext(true).build(); + + int consumerApplicationId = 0; + int providerApplicationId = 0; + int consumerInstanceId = 0; + int providerInstanceId = 0; + int consumerEntryServiceId = 0; + int consumerExitServiceId = 0; + int consumerExitApplicationId = 0; + int providerEntryServiceId = 0; + + while (consumerApplicationId == 0) { + consumerApplicationId = registerApplication(channel, "consumer"); + } + while (consumerExitApplicationId == 0) { + consumerExitApplicationId = registerApplication(channel, "172.25.0.4:20880"); + } + while (providerApplicationId == 0) { + providerApplicationId = registerApplication(channel, "provider"); + } + while (consumerInstanceId == 0) { + consumerInstanceId = registerInstanceId(channel, "ConsumerUUID", consumerApplicationId, "consumer_host_name", 1); + } + while (providerInstanceId == 0) { + providerInstanceId = registerInstanceId(channel, "ProviderUUID", providerApplicationId, "provider_host_name", 2); + } + while (consumerEntryServiceId == 0) { + consumerEntryServiceId = registerServiceId(channel, consumerApplicationId, "/dubbox-case/case/dubbox-rest"); + } + while (consumerExitServiceId == 0) { + consumerExitServiceId = registerServiceId(channel, consumerApplicationId, "org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + } + while (providerEntryServiceId == 0) { + providerEntryServiceId = registerServiceId(channel, providerApplicationId, "org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + } + + Ids ids = new Ids(); + ids.setConsumerApplicationId(consumerApplicationId); + ids.setProviderApplicationId(providerApplicationId); + ids.setConsumerInstanceId(consumerInstanceId); + ids.setProviderInstanceId(providerInstanceId); + ids.setConsumerEntryServiceId(consumerEntryServiceId); + ids.setConsumerExitServiceId(consumerExitServiceId); + ids.setConsumerExitApplicationId(consumerExitApplicationId); + + long startTime = TimeBucketUtils.INSTANCE.getSecondTimeBucket(System.currentTimeMillis()); + logger.info("start time: {}", startTime); + + int count = 10; + ThreadCount threadCount = new ThreadCount(count); + for (int i = 0; i < count; i++) { + Status status = new Status(); + BuildNewSegment buildNewSegment = new BuildNewSegment(channel, ids, threadCount, i, status); + Executors.newSingleThreadExecutor().execute(buildNewSegment); + } + + while (threadCount.getCount() != 0) { + try { + Thread.sleep(100); + } catch (InterruptedException e) { + } + } + long endTime = TimeBucketUtils.INSTANCE.getSecondTimeBucket(System.currentTimeMillis()); + logger.info("end time: {}", endTime); + + channel.shutdownNow(); + while (!channel.isTerminated()) { + try { + channel.awaitTermination(100, TimeUnit.SECONDS); + } catch (InterruptedException e) { + logger.error(e.getMessage(), e); + } + } + } + + private int registerApplication(ManagedChannel channel, String applicationCode) { + ApplicationRegisterServiceGrpc.ApplicationRegisterServiceBlockingStub stub = ApplicationRegisterServiceGrpc.newBlockingStub(channel); + Application application = Application.newBuilder().addApplicationCode(applicationCode).build(); + ApplicationMapping mapping = stub.register(application); + int applicationId = mapping.getApplication(0).getValue(); + + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + return applicationId; + } + + private int registerInstanceId(ManagedChannel channel, String agentUUId, Integer applicationId, + String hostName, int processNo) { + InstanceDiscoveryServiceGrpc.InstanceDiscoveryServiceBlockingStub stub = InstanceDiscoveryServiceGrpc.newBlockingStub(channel); + ApplicationInstance.Builder instance = ApplicationInstance.newBuilder(); + instance.setApplicationId(applicationId); + instance.setRegisterTime(System.currentTimeMillis()); + instance.setAgentUUID(agentUUId); + + OSInfo.Builder osInfo = OSInfo.newBuilder(); + osInfo.setHostname(hostName); + osInfo.setOsName("Linux"); + osInfo.setProcessNo(processNo); + osInfo.addIpv4S("10.0.0.1"); + osInfo.addIpv4S("10.0.0.2"); + instance.setOsinfo(osInfo.build()); + + ApplicationInstanceMapping mapping = stub.register(instance.build()); + int instanceId = mapping.getApplicationInstanceId(); + + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + return instanceId; + } + + private int registerServiceId(ManagedChannel channel, int applicationId, String serviceName) { + ServiceNameDiscoveryServiceGrpc.ServiceNameDiscoveryServiceBlockingStub stub = ServiceNameDiscoveryServiceGrpc.newBlockingStub(channel); + ServiceNameCollection.Builder collection = ServiceNameCollection.newBuilder(); + + ServiceNameElement.Builder element = ServiceNameElement.newBuilder(); + element.setApplicationId(applicationId); + element.setServiceName(serviceName); + collection.addElements(element); + + ServiceNameMappingCollection mappingCollection = stub.discovery(collection.build()); + int serviceId = mappingCollection.getElements(0).getServiceId(); + + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + return serviceId; + } + + class BuildNewSegment implements Runnable { + private final ManagedChannel segmentChannel; + private final Ids ids; + private final ThreadCount threadCount; + private final int procNo; + private final Status status; + private StreamObserver streamObserver; + + public BuildNewSegment(ManagedChannel segmentChannel, + Ids ids, ThreadCount threadCount, int procNo, + Status status) { + this.segmentChannel = segmentChannel; + this.ids = ids; + this.threadCount = threadCount; + this.procNo = procNo; + this.status = status; + } + + @Override public void run() { + statusChange(); + int i = 0; + while (i < 50000) { + send(streamObserver, ids); + + i++; + if (i % 10000 == 0) { + logger.info("process no: {}, send segment count: {}", procNo, i); + streamObserver.onCompleted(); + while (!status.isFinish) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + } + } + status.setFinish(false); + statusChange(); + } + } + this.threadCount.finishOne(); + } + + private void statusChange() { + TraceSegmentServiceGrpc.TraceSegmentServiceStub stub = TraceSegmentServiceGrpc.newStub(segmentChannel); + streamObserver = stub.collect(new StreamObserver() { + @Override public void onNext(Downstream downstream) { + } + + @Override public void onError(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + + @Override public void onCompleted() { + status.setFinish(true); + logger.info("process no: {}, server completed", procNo); + } + }); + } + } + + public void send(StreamObserver streamObserver, Ids ids) { + long now = System.currentTimeMillis(); + UniqueId consumerSegmentId = createSegmentId(); + UniqueId providerSegmentId = createSegmentId(); + + streamObserver.onNext(createConsumerSegment(consumerSegmentId, ids, now)); + streamObserver.onNext(createProviderSegment(consumerSegmentId, providerSegmentId, ids, now)); + } + + private UpstreamSegment createConsumerSegment(UniqueId segmentId, Ids ids, long timestamp) { + UpstreamSegment.Builder upstream = UpstreamSegment.newBuilder(); + upstream.addGlobalTraceIds(segmentId); + + TraceSegmentObject.Builder segmentBuilder = TraceSegmentObject.newBuilder(); + segmentBuilder.setApplicationId(ids.consumerApplicationId); + segmentBuilder.setApplicationInstanceId(ids.consumerInstanceId); + segmentBuilder.setTraceSegmentId(segmentId); + + SpanObject.Builder entrySpan = SpanObject.newBuilder(); + entrySpan.setSpanId(0); + entrySpan.setSpanType(SpanType.Entry); + entrySpan.setSpanLayer(SpanLayer.Http); + entrySpan.setParentSpanId(-1); + entrySpan.setStartTime(timestamp); + entrySpan.setEndTime(timestamp + 3000); + entrySpan.setComponentId(ComponentsDefine.TOMCAT.getId()); + entrySpan.setOperationNameId(ids.getConsumerEntryServiceId()); + entrySpan.setIsError(false); + + LogMessage.Builder entryLogMessage = LogMessage.newBuilder(); + entryLogMessage.setTime(timestamp); + + KeyWithStringValue.Builder data_1 = KeyWithStringValue.newBuilder(); + data_1.setKey("url"); + data_1.setValue("http://localhost:18080/dubbox-case/case/dubbox-rest"); + entryLogMessage.addData(data_1); + + KeyWithStringValue.Builder data_2 = KeyWithStringValue.newBuilder(); + data_2.setKey("http.method"); + data_2.setValue("GET"); + entryLogMessage.addData(data_2); + entrySpan.addLogs(entryLogMessage); + segmentBuilder.addSpans(entrySpan); + + SpanObject.Builder exitSpan = SpanObject.newBuilder(); + exitSpan.setSpanId(1); + exitSpan.setSpanType(SpanType.Exit); + exitSpan.setSpanLayer(SpanLayer.RPCFramework); + exitSpan.setParentSpanId(0); + exitSpan.setStartTime(timestamp + 500); + exitSpan.setEndTime(timestamp + 2500); + exitSpan.setComponentId(ComponentsDefine.TOMCAT.getId()); + exitSpan.setOperationNameId(ids.getConsumerExitServiceId()); + exitSpan.setPeerId(ids.consumerExitApplicationId); + exitSpan.setIsError(false); + + LogMessage.Builder exitLogMessage = LogMessage.newBuilder(); + exitLogMessage.setTime(timestamp); + + KeyWithStringValue.Builder data = KeyWithStringValue.newBuilder(); + data.setKey("url"); + data.setValue("rest://172.25.0.4:20880/org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + exitLogMessage.addData(data); + exitSpan.addLogs(exitLogMessage); + segmentBuilder.addSpans(exitSpan); + + upstream.setSegment(segmentBuilder.build().toByteString()); + return upstream.build(); + } + + private UpstreamSegment createProviderSegment(UniqueId consumerSegmentId, UniqueId providerSegmentId, Ids ids, + long timestamp) { + UpstreamSegment.Builder upstream = UpstreamSegment.newBuilder(); + upstream.addGlobalTraceIds(consumerSegmentId); + + TraceSegmentObject.Builder segmentBuilder = TraceSegmentObject.newBuilder(); + segmentBuilder.setApplicationId(ids.providerApplicationId); + segmentBuilder.setApplicationInstanceId(ids.providerInstanceId); + segmentBuilder.setTraceSegmentId(providerSegmentId); + + TraceSegmentReference.Builder referenceBuilder = TraceSegmentReference.newBuilder(); + referenceBuilder.setParentTraceSegmentId(consumerSegmentId); + referenceBuilder.setParentApplicationInstanceId(ids.getConsumerInstanceId()); + referenceBuilder.setParentSpanId(1); + referenceBuilder.setParentServiceId(ids.getConsumerExitServiceId()); + referenceBuilder.setEntryApplicationInstanceId(ids.getConsumerInstanceId()); + referenceBuilder.setEntryServiceId(ids.getConsumerEntryServiceId()); + referenceBuilder.setNetworkAddressId(ids.consumerExitApplicationId); + referenceBuilder.setRefType(RefType.CrossProcess); + segmentBuilder.addRefs(referenceBuilder); + + SpanObject.Builder entrySpan = SpanObject.newBuilder(); + entrySpan.setSpanId(0); + entrySpan.setSpanType(SpanType.Entry); + entrySpan.setSpanLayer(SpanLayer.RPCFramework); + entrySpan.setParentSpanId(-1); + entrySpan.setStartTime(timestamp + 1000); + entrySpan.setEndTime(timestamp + 2000); + entrySpan.setComponentId(ComponentsDefine.TOMCAT.getId()); + entrySpan.setOperationNameId(ids.getProviderEntryServiceId()); + entrySpan.setIsError(false); + + LogMessage.Builder entryLogMessage = LogMessage.newBuilder(); + entryLogMessage.setTime(timestamp); + + KeyWithStringValue.Builder data_1 = KeyWithStringValue.newBuilder(); + data_1.setKey("url"); + data_1.setValue("rest://172.25.0.4:20880/org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()"); + entryLogMessage.addData(data_1); + + KeyWithStringValue.Builder data_2 = KeyWithStringValue.newBuilder(); + data_2.setKey("http.method"); + data_2.setValue("GET"); + entryLogMessage.addData(data_2); + entrySpan.addLogs(entryLogMessage); + segmentBuilder.addSpans(entrySpan); + + upstream.setSegment(segmentBuilder.build().toByteString()); + return upstream.build(); + } + + private UniqueId createSegmentId() { + long id = sequence.getAndIncrement(); + UniqueId.Builder builder = UniqueId.newBuilder(); + builder.addIdParts(id); + builder.addIdParts(id); + builder.addIdParts(id); + return builder.build(); + } + + class Ids { + private int consumerApplicationId = 0; + private int providerApplicationId = 0; + private int consumerInstanceId = 0; + private int providerInstanceId = 0; + private int consumerEntryServiceId = 0; + private int consumerExitServiceId = 0; + private int consumerExitApplicationId = 0; + private int providerEntryServiceId = 0; + + public int getConsumerApplicationId() { + return consumerApplicationId; + } + + public void setConsumerApplicationId(int consumerApplicationId) { + this.consumerApplicationId = consumerApplicationId; + } + + public int getProviderApplicationId() { + return providerApplicationId; + } + + public void setProviderApplicationId(int providerApplicationId) { + this.providerApplicationId = providerApplicationId; + } + + public int getConsumerInstanceId() { + return consumerInstanceId; + } + + public void setConsumerInstanceId(int consumerInstanceId) { + this.consumerInstanceId = consumerInstanceId; + } + + public int getProviderInstanceId() { + return providerInstanceId; + } + + public void setProviderInstanceId(int providerInstanceId) { + this.providerInstanceId = providerInstanceId; + } + + public int getConsumerEntryServiceId() { + return consumerEntryServiceId; + } + + public void setConsumerEntryServiceId(int consumerEntryServiceId) { + this.consumerEntryServiceId = consumerEntryServiceId; + } + + public int getConsumerExitServiceId() { + return consumerExitServiceId; + } + + public void setConsumerExitServiceId(int consumerExitServiceId) { + this.consumerExitServiceId = consumerExitServiceId; + } + + public int getConsumerExitApplicationId() { + return consumerExitApplicationId; + } + + public void setConsumerExitApplicationId(int consumerExitApplicationId) { + this.consumerExitApplicationId = consumerExitApplicationId; + } + + public int getProviderEntryServiceId() { + return providerEntryServiceId; + } + + public void setProviderEntryServiceId(int providerEntryServiceId) { + this.providerEntryServiceId = providerEntryServiceId; + } + } + + class ThreadCount { + private int count; + + public ThreadCount(int count) { + this.count = count; + } + + public void finishOne() { + count--; + } + + public int getCount() { + return count; + } + } + + class Status { + private boolean isFinish = false; + + public boolean isFinish() { + return isFinish; + } + + public void setFinish(boolean finish) { + isFinish = finish; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/test/resources/logback.xml b/apm-collector/apm-collector-agentstream/src/test/resources/logback.xml new file mode 100644 index 000000000..46eba2b93 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/resources/logback.xml @@ -0,0 +1,16 @@ + + + + + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + + + + + + + \ No newline at end of file diff --git a/apm-collector/apm-collector-boot/src/main/resources/application.yml b/apm-collector/apm-collector-boot/src/main/resources/application.yml index 1a84f5740..b1d440ec5 100644 --- a/apm-collector/apm-collector-boot/src/main/resources/application.yml +++ b/apm-collector/apm-collector-boot/src/main/resources/application.yml @@ -25,6 +25,6 @@ storage: elasticsearch: cluster_name: CollectorDBCluster cluster_transport_sniffer: true - cluster_nodes: 127.0.0.1:9300 + cluster_nodes: 10.0.0.19:9300,10.0.0.6:9300 diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java index 96edf8f02..2653ab530 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java @@ -4,12 +4,16 @@ import java.util.LinkedList; import java.util.List; import org.skywalking.apm.collector.core.framework.Handler; import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class ServerHolder { + private final Logger logger = LoggerFactory.getLogger(ServerHolder.class); + private List servers; public ServerHolder() { @@ -33,7 +37,10 @@ public class ServerHolder { private void addHandler(List handlers, Server server) { if (CollectionUtils.isNotEmpty(handlers)) { - handlers.forEach(handler -> server.addHandler(handler)); + handlers.forEach(handler -> { + server.addHandler(handler); + logger.debug("add handler into server: {}, handler name: {}", server.hostPort(), handler.getClass().getName()); + }); } } diff --git a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto index 159bfe9bc..b2100424f 100644 --- a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto +++ b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto @@ -4,7 +4,7 @@ option java_multiple_files = true; option java_package = "org.skywalking.apm.collector.remote.grpc.proto"; service RemoteCommonService { - rpc call (RemoteMessage) returns (Empty) { + rpc call (stream RemoteMessage) returns (Empty) { } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java index 86b5dc2b0..d1d92ebd3 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java @@ -22,17 +22,29 @@ public class RemoteCommonServiceHandler extends RemoteCommonServiceGrpc.RemoteCo private final Logger logger = LoggerFactory.getLogger(RemoteCommonServiceHandler.class); - @Override public void call(RemoteMessage request, StreamObserver responseObserver) { - String roleName = request.getWorkerRole(); - RemoteData remoteData = request.getRemoteData(); + @Override public StreamObserver call(StreamObserver responseObserver) { + return new StreamObserver() { + @Override public void onNext(RemoteMessage message) { + String roleName = message.getWorkerRole(); + RemoteData remoteData = message.getRemoteData(); - StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); - Role role = context.getClusterWorkerContext().getRole(roleName); - Object object = role.dataDefine().deserialize(remoteData); - try { - context.getClusterWorkerContext().lookupInSide(roleName).tell(object); - } catch (WorkerNotFoundException | WorkerInvokeException e) { - logger.error(e.getMessage(), e); - } + StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); + Role role = context.getClusterWorkerContext().getRole(roleName); + Object object = role.dataDefine().deserialize(remoteData); + try { + context.getClusterWorkerContext().lookupInSide(roleName).tell(object); + } catch (WorkerNotFoundException | WorkerInvokeException e) { + logger.error(e.getMessage(), e); + } + } + + @Override public void onError(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + + @Override public void onCompleted() { + responseObserver.onCompleted(); + } + }; } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java index 5e99ab101..707230f6c 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/RemoteWorkerRef.java @@ -1,17 +1,24 @@ package org.skywalking.apm.collector.stream.worker; +import io.grpc.stub.StreamObserver; import org.skywalking.apm.collector.client.grpc.GRPCClient; +import org.skywalking.apm.collector.remote.grpc.proto.Empty; import org.skywalking.apm.collector.remote.grpc.proto.RemoteCommonServiceGrpc; import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; import org.skywalking.apm.collector.remote.grpc.proto.RemoteMessage; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class RemoteWorkerRef extends WorkerRef { + private final Logger logger = LoggerFactory.getLogger(RemoteWorkerRef.class); + private final Boolean acrossJVM; - private final RemoteCommonServiceGrpc.RemoteCommonServiceBlockingStub stub; + private final RemoteCommonServiceGrpc.RemoteCommonServiceStub stub; + private StreamObserver streamObserver; private final AbstractRemoteWorker remoteWorker; public RemoteWorkerRef(Role role, AbstractRemoteWorker remoteWorker) { @@ -25,7 +32,8 @@ public class RemoteWorkerRef extends WorkerRef { super(role); this.remoteWorker = null; this.acrossJVM = true; - this.stub = RemoteCommonServiceGrpc.newBlockingStub(client.getChannel()); + this.stub = RemoteCommonServiceGrpc.newStub(client.getChannel()); + createStreamObserver(); } @Override @@ -36,7 +44,8 @@ public class RemoteWorkerRef extends WorkerRef { RemoteMessage.Builder builder = RemoteMessage.newBuilder(); builder.setWorkerRole(getRole().roleName()); builder.setRemoteData(remoteData); - stub.call(builder.build()); + + streamObserver.onNext(builder.build()); } else { remoteWorker.allocateJob(message); } @@ -45,4 +54,63 @@ public class RemoteWorkerRef extends WorkerRef { public Boolean isAcrossJVM() { return acrossJVM; } + + private void createStreamObserver() { + StreamStatus status = new StreamStatus(false); + streamObserver = stub.call(new StreamObserver() { + @Override public void onNext(Empty empty) { + } + + @Override public void onError(Throwable throwable) { + logger.error(throwable.getMessage(), throwable); + } + + @Override public void onCompleted() { + status.finished(); + } + }); + } + + class StreamStatus { + private volatile boolean status; + + public StreamStatus(boolean status) { + this.status = status; + } + + public boolean isFinish() { + return status; + } + + public void finished() { + this.status = true; + } + + /** + * @param maxTimeout max wait time, milliseconds. + */ + public void wait4Finish(long maxTimeout) { + long time = 0; + while (!status) { + if (time > maxTimeout) { + break; + } + try2Sleep(5); + time += 5; + } + } + + /** + * Try to sleep, and ignore the {@link InterruptedException} + * + * @param millis the length of time to sleep in milliseconds + */ + private void try2Sleep(long millis) { + try { + Thread.sleep(millis); + } catch (InterruptedException e) { + + } + } + } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java index 4cc4676d6..255d07016 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java @@ -51,7 +51,7 @@ public abstract class AggregationWorker extends AbstractLocalAsyncWorker { private void sendToNext() throws WorkerException { dataCache.switchPointer(); - while (dataCache.getLast().isHolding()) { + while (dataCache.getLast().isWriting()) { try { Thread.sleep(10); } catch (InterruptedException e) { @@ -66,17 +66,17 @@ public abstract class AggregationWorker extends AbstractLocalAsyncWorker { logger.error(e.getMessage(), e); } }); - dataCache.releaseLast(); + dataCache.finishReadingLast(); } protected final void aggregate(Object message) { Data data = (Data)message; - dataCache.hold(); + dataCache.writing(); if (dataCache.containsKey(data.id())) { getRole().dataDefine().mergeData(data, dataCache.get(data.id())); } else { dataCache.put(data.id(), data); } - dataCache.release(); + dataCache.finishWriting(); } } 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 875cba06b..c1e509c20 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,11 +1,13 @@ package org.skywalking.apm.collector.stream.worker.impl; -import java.util.ArrayList; +import java.util.LinkedList; import java.util.List; import java.util.Map; import org.skywalking.apm.collector.core.queue.EndOfBatchCommand; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.core.util.ObjectUtils; +import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.storage.dao.IBatchDAO; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorker; import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; @@ -35,24 +37,37 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { } @Override protected final void onWork(Object message) throws WorkerException { - if (message instanceof EndOfBatchCommand || message instanceof FlushAndSwitch) { - if (dataCache.trySwitchPointer()) { - dataCache.switchPointer(); - } - } else { - if (dataCache.currentCollectionSize() >= 1000) { + if (message instanceof FlushAndSwitch) { + try { if (dataCache.trySwitchPointer()) { dataCache.switchPointer(); } + } finally { + dataCache.trySwitchPointerFinally(); + } + } else if (message instanceof EndOfBatchCommand) { + } else { + if (dataCache.currentCollectionSize() >= 5000) { + try { + if (dataCache.trySwitchPointer()) { + dataCache.switchPointer(); + + List collection = buildBatchCollection(); + IBatchDAO dao = (IBatchDAO)DAOContainer.INSTANCE.get(IBatchDAO.class.getName()); + dao.batchPersistence(collection); + } + } finally { + dataCache.trySwitchPointerFinally(); + } } aggregate(message); } } public final List buildBatchCollection() throws WorkerException { - List batchCollection; + List batchCollection = new LinkedList<>(); try { - while (dataCache.getLast().isHolding()) { + while (dataCache.getLast().isWriting()) { try { Thread.sleep(10); } catch (InterruptedException e) { @@ -60,16 +75,18 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { } } - batchCollection = prepareBatch(dataCache.getLast().asMap()); + if (dataCache.getLast().asMap() != null) { + batchCollection = prepareBatch(dataCache.getLast().asMap()); + } } finally { - dataCache.releaseLast(); + dataCache.finishReadingLast(); } return batchCollection; } protected final List prepareBatch(Map dataMap) { - List insertBatchCollection = new ArrayList<>(); - List updateBatchCollection = new ArrayList<>(); + List insertBatchCollection = new LinkedList<>(); + List updateBatchCollection = new LinkedList<>(); dataMap.forEach((id, data) -> { if (needMergeDBData()) { Data dbData = persistenceDAO().get(id, getRole().dataDefine()); @@ -101,18 +118,16 @@ public abstract class PersistenceWorker extends AbstractLocalAsyncWorker { } private void aggregate(Object message) { - dataCache.hold(); + dataCache.writing(); Data data = (Data)message; if (dataCache.containsKey(data.id())) { getRole().dataDefine().mergeData(data, dataCache.get(data.id())); } else { - if (dataCache.currentCollectionSize() < 1000) { - dataCache.put(data.id(), data); - } + dataCache.put(data.id(), data); } - dataCache.release(); + dataCache.finishWriting(); } protected abstract IPersistenceDAO persistenceDAO(); diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java index ff81c3ea8..459e7062f 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCache.java @@ -21,16 +21,16 @@ public class DataCache extends Window { lockedDataCollection.put(id, data); } - public void hold() { - lockedDataCollection = getCurrentAndHold(); + public void writing() { + lockedDataCollection = getCurrentAndWriting(); } public int currentCollectionSize() { - return getCurrentAndHold().size(); + return getCurrent().size(); } - public void release() { - lockedDataCollection.release(); + public void finishWriting() { + lockedDataCollection.finishWriting(); lockedDataCollection = null; } } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java index bc280d8d4..ee599d983 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataCollection.java @@ -1,7 +1,7 @@ package org.skywalking.apm.collector.stream.worker.impl.data; -import java.util.HashMap; import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import org.skywalking.apm.collector.core.stream.Data; /** @@ -9,23 +9,37 @@ import org.skywalking.apm.collector.core.stream.Data; */ public class DataCollection { private Map data; - private volatile boolean isHold; + private volatile boolean writing; + private volatile boolean reading; public DataCollection() { - this.data = new HashMap<>(); - this.isHold = false; + this.data = new ConcurrentHashMap<>(); + this.writing = false; + this.reading = false; } - public void release() { - isHold = false; + public void finishWriting() { + writing = false; } - public void hold() { - isHold = true; + public void writing() { + writing = true; } - public boolean isHolding() { - return isHold; + public boolean isWriting() { + return writing; + } + + public void finishReading() { + reading = false; + } + + public void reading() { + reading = true; + } + + public boolean isReading() { + return reading; } public boolean containsKey(String key) { diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java index 887f24809..f2f729f7d 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java @@ -21,32 +21,40 @@ public abstract class Window { } public boolean trySwitchPointer() { - if (windowSwitch.incrementAndGet() == 1) { + if (windowSwitch.incrementAndGet() == 1 && !getLast().isReading()) { return true; } else { - windowSwitch.addAndGet(-1); return false; } } + public void trySwitchPointerFinally() { + windowSwitch.addAndGet(-1); + } + public void switchPointer() { if (pointer == windowDataA) { pointer = windowDataB; } else { pointer = windowDataA; } + getLast().reading(); } - protected DataCollection getCurrentAndHold() { + protected DataCollection getCurrentAndWriting() { if (pointer == windowDataA) { - windowDataA.hold(); + windowDataA.writing(); return windowDataA; } else { - windowDataB.hold(); + windowDataB.writing(); return windowDataB; } } + protected DataCollection getCurrent() { + return pointer; + } + public DataCollection getLast() { if (pointer == windowDataA) { return windowDataB; @@ -55,8 +63,8 @@ public abstract class Window { } } - public void releaseLast() { + public void finishReadingLast() { getLast().clear(); - windowSwitch.addAndGet(-1); + getLast().finishReading(); } } From 2178a736a91f5c8593013248c5521c701d658d17 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 3 Sep 2017 16:39:26 +0800 Subject: [PATCH 2/2] fixed check style error --- .../apm/collector/stream/worker/impl/data/Window.java | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java index f2f729f7d..851285f01 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Window.java @@ -21,11 +21,7 @@ public abstract class Window { } public boolean trySwitchPointer() { - if (windowSwitch.incrementAndGet() == 1 && !getLast().isReading()) { - return true; - } else { - return false; - } + return windowSwitch.incrementAndGet() == 1 && !getLast().isReading(); } public void trySwitchPointerFinally() {