From d9556c38fca934d0de967b9f1c29de831fe6b986 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=90=B4=E6=99=9F=20Wu=20Sheng?= Date: Fri, 26 Oct 2018 15:17:33 +0800 Subject: [PATCH] Make Mesh and Istio receivers ready (#1821) * Make receiver more effective * Add Istio test case and all source dispatch in * Refactor mock data * FIx rat. * Support call component. * Fix ThermodynamicIndicator bug. * Fix test cases. * Fix missing calculate in db merging. * Add codes for debug. * 1. Fixed elasticsearch bulk process not fresh bug. (#1819) 2. Fixed the bug of source register no queue but wait the end for batch tag. * Remove debug log, and restore TTL timer for real scenarios. --- .../apm/commons/datacarrier/DataCarrier.java | 10 +- .../datacarrier/consumer/ConsumerPool.java | 9 +- .../consumer/ConsumerPoolTest.java | 4 +- .../indicator/ThermodynamicIndicator.java | 18 +- .../worker/IndicatorAggregateWorker.java | 16 +- .../worker/IndicatorPersistentWorker.java | 19 +- .../analysis/worker/IndicatorProcess.java | 4 +- .../worker/IndicatorRemoteWorker.java | 8 +- .../analysis/worker/PersistenceWorker.java | 2 +- .../service/ServiceInventoryRegister.java | 1 + .../worker/RegisterDistinctWorker.java | 6 +- .../worker/RegisterPersistentWorker.java | 49 ++++- .../core/remote/client/GRPCRemoteClient.java | 2 +- .../core/storage/ttl/DataTTLKeeperTimer.java | 21 +- .../indicator/ThermodynamicIndicatorTest.java | 10 +- .../IstioTelemetryHandlerMainTest.java | 86 ++++++++ .../fixture/01-ingress-reviewsv1-policy.msg | 193 ++++++++++++++++++ .../fixture/02-productpage-details-policy.msg | 193 ++++++++++++++++++ .../fixture/03-productpage-details.msg | 193 ++++++++++++++++++ .../fixture/04-productpage-details.msg | 105 ++++++++++ .../fixture/05-productpage-reviews.msg | 105 ++++++++++ .../fixture/06-productpage-reviews.msg | 105 ++++++++++ .../fixture/07-ingress-productpage.msg | 105 ++++++++++ .../fixture/08-ingress-productpage.msg | 105 ++++++++++ .../resources/fixture/09-policy-telemetry.msg | 105 ++++++++++ .../resources/fixture/10-policy-telemetry.msg | 105 ++++++++++ .../fixture/11-details-telemetry.msg | 105 ++++++++++ .../12-ingress-productpage-telemetry.msg | 193 ++++++++++++++++++ .../fixture/13-productpage-telemetry.msg | 105 ++++++++++ .../14-reviews-productpage-telemetry.msg | 193 ++++++++++++++++++ .../mesh/MeshDataBufferFileCache.java | 2 +- .../mesh/TelemetryDataDispatcher.java | 47 ++++- .../main/resources/component-libraries.yml | 6 + .../elasticsearch/base/BatchProcessEsDAO.java | 2 + .../elasticsearch/query/MetricQueryEsDAO.java | 2 +- 35 files changed, 2185 insertions(+), 49 deletions(-) create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/handler/IstioTelemetryHandlerMainTest.java create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/01-ingress-reviewsv1-policy.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/02-productpage-details-policy.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/03-productpage-details.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/04-productpage-details.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/05-productpage-reviews.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/06-productpage-reviews.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/07-ingress-productpage.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/08-ingress-productpage.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/09-policy-telemetry.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/10-policy-telemetry.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/11-details-telemetry.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/12-ingress-productpage-telemetry.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/13-productpage-telemetry.msg create mode 100644 oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/14-reviews-productpage-telemetry.msg diff --git a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/DataCarrier.java b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/DataCarrier.java index d53fda539..76ad609bf 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/DataCarrier.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/DataCarrier.java @@ -33,8 +33,14 @@ public class DataCarrier { private final int channelSize; private Channels channels; private ConsumerPool consumerPool; + private String name; public DataCarrier(int channelSize, int bufferSize) { + this("default", channelSize, bufferSize); + } + + public DataCarrier(String name, int channelSize, int bufferSize) { + this.name = name; this.bufferSize = bufferSize; this.channelSize = channelSize; channels = new Channels(channelSize, bufferSize, new SimpleRollingPartitioner(), BufferStrategy.BLOCKING); @@ -92,7 +98,7 @@ public class DataCarrier { if (consumerPool != null) { consumerPool.close(); } - consumerPool = new ConsumerPool(this.channels, consumerClass, num, consumeCycle); + consumerPool = new ConsumerPool(this.name, this.channels, consumerClass, num, consumeCycle); consumerPool.begin(); return this; } @@ -119,7 +125,7 @@ public class DataCarrier { if (consumerPool != null) { consumerPool.close(); } - consumerPool = new ConsumerPool(this.channels, consumer, num, consumeCycle); + consumerPool = new ConsumerPool(this.name, this.channels, consumer, num, consumeCycle); consumerPool.begin(); return this; } diff --git a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java index 814fdf8f8..d3ba23eba 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java @@ -32,19 +32,20 @@ public class ConsumerPool { private Channels channels; private ReentrantLock lock; - public ConsumerPool(Channels channels, Class> consumerClass, int num, long consumeCycle) { + public ConsumerPool(String name, Channels channels, Class> consumerClass, int num, + long consumeCycle) { this(channels, num); for (int i = 0; i < num; i++) { - consumerThreads[i] = new ConsumerThread("DataCarrier.Consumser." + i + ".Thread", getNewConsumerInstance(consumerClass), consumeCycle); + consumerThreads[i] = new ConsumerThread("DataCarrier." + name + ".Consumser." + i + ".Thread", getNewConsumerInstance(consumerClass), consumeCycle); consumerThreads[i].setDaemon(true); } } - public ConsumerPool(Channels channels, IConsumer prototype, int num, long consumeCycle) { + public ConsumerPool(String name, Channels channels, IConsumer prototype, int num, long consumeCycle) { this(channels, num); prototype.init(); for (int i = 0; i < num; i++) { - consumerThreads[i] = new ConsumerThread("DataCarrier.Consumser." + i + ".Thread", prototype, consumeCycle); + consumerThreads[i] = new ConsumerThread("DataCarrier." + name + ".Consumser." + i + ".Thread", prototype, consumeCycle); consumerThreads[i].setDaemon(true); } diff --git a/apm-commons/apm-datacarrier/src/test/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPoolTest.java b/apm-commons/apm-datacarrier/src/test/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPoolTest.java index 18eab0eac..885a5bf4d 100644 --- a/apm-commons/apm-datacarrier/src/test/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPoolTest.java +++ b/apm-commons/apm-datacarrier/src/test/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerPoolTest.java @@ -34,7 +34,7 @@ public class ConsumerPoolTest { @Test public void testBeginConsumerPool() throws IllegalAccessException { Channels channels = new Channels(2, 100, new SimpleRollingPartitioner(), BufferStrategy.BLOCKING); - ConsumerPool pool = new ConsumerPool(channels, new SampleConsumer(), 2, 20); + ConsumerPool pool = new ConsumerPool("default", channels, new SampleConsumer(), 2, 20); pool.begin(); ConsumerThread[] threads = (ConsumerThread[])MemberModifier.field(ConsumerPool.class, "consumerThreads").get(pool); @@ -46,7 +46,7 @@ public class ConsumerPoolTest { @Test public void testCloseConsumerPool() throws InterruptedException, IllegalAccessException { Channels channels = new Channels(2, 100, new SimpleRollingPartitioner(), BufferStrategy.BLOCKING); - ConsumerPool pool = new ConsumerPool(channels, new SampleConsumer(), 2, 20); + ConsumerPool pool = new ConsumerPool("default", channels, new SampleConsumer(), 2, 20); pool.begin(); Thread.sleep(5000); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicator.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicator.java index 69ba7bae5..830683cc6 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicator.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicator.java @@ -18,9 +18,14 @@ package org.apache.skywalking.oap.server.core.analysis.indicator; -import java.util.*; -import lombok.*; -import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.*; +import java.util.HashMap; +import java.util.Map; +import lombok.Getter; +import lombok.Setter; +import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.Arg; +import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.Entrance; +import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.IndicatorOperator; +import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.SourceFrom; import org.apache.skywalking.oap.server.core.storage.annotation.Column; /** @@ -61,7 +66,7 @@ public abstract class ThermodynamicIndicator extends Indicator { this.step = step; } if (this.numOfSteps == 0) { - this.numOfSteps = maxNumOfSteps + 1; + this.numOfSteps = maxNumOfSteps; } indexCheckAndInit(); @@ -86,14 +91,15 @@ public abstract class ThermodynamicIndicator extends Indicator { ThermodynamicIndicator thermodynamicIndicator = (ThermodynamicIndicator)indicator; this.indexCheckAndInit(); thermodynamicIndicator.indexCheckAndInit(); + final ThermodynamicIndicator self = this; thermodynamicIndicator.detailIndex.forEach((key, element) -> { - IntKeyLongValue existingElement = this.detailIndex.get(key); + IntKeyLongValue existingElement = self.detailIndex.get(key); if (existingElement == null) { existingElement = new IntKeyLongValue(); existingElement.setKey(key); existingElement.setValue(element.getValue()); - addElement(element); + self.addElement(element); } else { existingElement.addValue(element.getValue()); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorAggregateWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorAggregateWorker.java index 15ff40b4f..0f026519f 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorAggregateWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorAggregateWorker.java @@ -18,13 +18,16 @@ package org.apache.skywalking.oap.server.core.analysis.worker; -import java.util.*; +import java.util.Iterator; +import java.util.List; import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer; -import org.apache.skywalking.oap.server.core.analysis.data.*; +import org.apache.skywalking.oap.server.core.analysis.data.EndOfBatchContext; +import org.apache.skywalking.oap.server.core.analysis.data.MergeDataCache; import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; import org.apache.skywalking.oap.server.core.worker.AbstractWorker; -import org.slf4j.*; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng @@ -37,12 +40,14 @@ public class IndicatorAggregateWorker extends AbstractWorker { private final DataCarrier dataCarrier; private final MergeDataCache mergeDataCache; private int messageNum; + private final String modelName; - IndicatorAggregateWorker(int workerId, AbstractWorker nextWorker) { + IndicatorAggregateWorker(int workerId, AbstractWorker nextWorker, String modelName) { super(workerId); + this.modelName = modelName; this.nextWorker = nextWorker; this.mergeDataCache = new MergeDataCache<>(); - this.dataCarrier = new DataCarrier<>(1, 10000); + this.dataCarrier = new DataCarrier<>("IndicatorAggregateWorker." + modelName, 1, 10000); this.dataCarrier.consume(new AggregatorConsumer(this), 1); } @@ -88,6 +93,7 @@ public class IndicatorAggregateWorker extends AbstractWorker { } else { mergeDataCache.put(indicator); } + mergeDataCache.finishWriting(); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorPersistentWorker.java index efecef967..8763d8b12 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorPersistentWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorPersistentWorker.java @@ -18,15 +18,20 @@ package org.apache.skywalking.oap.server.core.analysis.worker; -import java.util.*; +import java.util.Iterator; +import java.util.LinkedList; +import java.util.List; +import java.util.Objects; import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer; -import org.apache.skywalking.oap.server.core.analysis.data.*; +import org.apache.skywalking.oap.server.core.analysis.data.EndOfBatchContext; +import org.apache.skywalking.oap.server.core.analysis.data.MergeDataCache; import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; import org.apache.skywalking.oap.server.core.storage.IIndicatorDAO; import org.apache.skywalking.oap.server.core.worker.AbstractWorker; import org.apache.skywalking.oap.server.library.module.ModuleManager; -import org.slf4j.*; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import static java.util.Objects.nonNull; @@ -50,10 +55,14 @@ public class IndicatorPersistentWorker extends PersistenceWorker(); this.indicatorDAO = indicatorDAO; this.nextWorker = nextWorker; - this.dataCarrier = new DataCarrier<>(1, 10000); + this.dataCarrier = new DataCarrier<>("IndicatorPersistentWorker." + modelName, 1, 10000); this.dataCarrier.consume(new IndicatorPersistentWorker.PersistentConsumer(this), 1); } + @Override void onWork(Indicator indicator) { + super.onWork(indicator); + } + @Override public void in(Indicator indicator) { indicator.setEndOfBatchContext(new EndOfBatchContext(false)); dataCarrier.produce(indicator); @@ -87,6 +96,8 @@ public class IndicatorPersistentWorker extends PersistenceWorker { private final AbstractWorker nextWorker; private final RemoteSenderService remoteSender; + private final String modelName; - IndicatorRemoteWorker(int workerId, ModuleManager moduleManager, AbstractWorker nextWorker) { + IndicatorRemoteWorker(int workerId, ModuleManager moduleManager, AbstractWorker nextWorker, + String modelName) { super(workerId); this.remoteSender = moduleManager.find(CoreModule.NAME).getService(RemoteSenderService.class); this.nextWorker = nextWorker; + this.modelName = modelName; } @Override public final void in(Indicator indicator) { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/PersistenceWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/PersistenceWorker.java index d559ef82c..b9e67f05c 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/PersistenceWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/PersistenceWorker.java @@ -41,7 +41,7 @@ public abstract class PersistenceWorker= batchSize) { try { if (getCache().trySwitchPointer()) { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/ServiceInventoryRegister.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/ServiceInventoryRegister.java index f9a0b78b7..531d3f15c 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/ServiceInventoryRegister.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/ServiceInventoryRegister.java @@ -62,6 +62,7 @@ public class ServiceInventoryRegister implements IServiceInventoryRegister { long now = System.currentTimeMillis(); serviceInventory.setRegisterTime(now); serviceInventory.setHeartbeatTime(now); + serviceInventory.setMappingServiceId(Const.NONE); serviceInventory.setMappingLastUpdateTime(now); InventoryProcess.INSTANCE.in(serviceInventory); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java index 4a677cecf..c02f150b1 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java @@ -22,7 +22,7 @@ import java.util.*; import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer; import org.apache.skywalking.oap.server.core.analysis.data.EndOfBatchContext; -import org.apache.skywalking.oap.server.core.register.RegisterSource; +import org.apache.skywalking.oap.server.core.register.*; import org.apache.skywalking.oap.server.core.worker.AbstractWorker; import org.slf4j.*; @@ -61,7 +61,9 @@ public class RegisterDistinctWorker extends AbstractWorker { } if (messageNum >= 1000 || source.getEndOfBatchContext().isEndOfBatch()) { - sources.values().forEach(nextWorker::in); + sources.values().forEach(source1 -> { + nextWorker.in(source1); + }); messageNum = 0; } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java index cfed8c306..0be299c3d 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java @@ -19,7 +19,10 @@ package org.apache.skywalking.oap.server.core.register.worker; import java.util.*; -import org.apache.skywalking.oap.server.core.register.RegisterSource; +import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; +import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer; +import org.apache.skywalking.oap.server.core.analysis.data.EndOfBatchContext; +import org.apache.skywalking.oap.server.core.register.*; import org.apache.skywalking.oap.server.core.source.Scope; import org.apache.skywalking.oap.server.core.storage.*; import org.apache.skywalking.oap.server.core.worker.AbstractWorker; @@ -38,6 +41,7 @@ public class RegisterPersistentWorker extends AbstractWorker { private final Map sources; private final IRegisterLockDAO registerLockDAO; private final IRegisterDAO registerDAO; + private final DataCarrier dataCarrier; RegisterPersistentWorker(int workerId, String modelName, ModuleManager moduleManager, IRegisterDAO registerDAO, Scope scope) { @@ -47,9 +51,16 @@ public class RegisterPersistentWorker extends AbstractWorker { this.registerDAO = registerDAO; this.registerLockDAO = moduleManager.find(StorageModule.NAME).getService(IRegisterLockDAO.class); this.scope = scope; + this.dataCarrier = new DataCarrier<>("IndicatorPersistentWorker." + modelName, 1, 10000); + this.dataCarrier.consume(new RegisterPersistentWorker.PersistentConsumer(this), 1); } @Override public final void in(RegisterSource registerSource) { + registerSource.setEndOfBatchContext(new EndOfBatchContext(false)); + dataCarrier.produce(registerSource); + } + + private void onWork(RegisterSource registerSource) { if (!sources.containsKey(registerSource)) { sources.put(registerSource, registerSource); } @@ -76,7 +87,43 @@ public class RegisterPersistentWorker extends AbstractWorker { } finally { registerLockDAO.releaseLock(scope); } + } else { + logger.info("Inventory register try lock failure."); } } } + + private class PersistentConsumer implements IConsumer { + + private final RegisterPersistentWorker persistent; + + private PersistentConsumer(RegisterPersistentWorker persistent) { + this.persistent = persistent; + } + + @Override public void init() { + + } + + @Override public void consume(List data) { + Iterator sourceIterator = data.iterator(); + + int i = 0; + while (sourceIterator.hasNext()) { + RegisterSource indicator = sourceIterator.next(); + i++; + if (i == data.size()) { + indicator.getEndOfBatchContext().setEndOfBatch(true); + } + persistent.onWork(indicator); + } + } + + @Override public void onError(List data, Throwable t) { + logger.error(t.getMessage(), t); + } + + @Override public void onExit() { + } + } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java index c2cb31992..48106ee25 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java @@ -45,7 +45,7 @@ public class GRPCRemoteClient implements RemoteClient, Comparable(channelSize, bufferSize); + this.carrier = new DataCarrier<>("GRPCRemoteClient", channelSize, bufferSize); this.carrier.setBufferStrategy(BufferStrategy.BLOCKING); this.carrier.consume(new RemoteMessageConsumer(), 1); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java index 98a08b3eb..2d920ad81 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java @@ -20,20 +20,29 @@ package org.apache.skywalking.oap.server.core.storage.ttl; import java.io.IOException; import java.util.List; -import java.util.concurrent.*; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import lombok.Setter; import org.apache.skywalking.apm.util.RunnableWithExceptionProtection; -import org.apache.skywalking.oap.server.core.*; +import org.apache.skywalking.oap.server.core.Const; +import org.apache.skywalking.oap.server.core.CoreModule; +import org.apache.skywalking.oap.server.core.DataTTL; import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; import org.apache.skywalking.oap.server.core.analysis.record.Record; -import org.apache.skywalking.oap.server.core.cluster.*; +import org.apache.skywalking.oap.server.core.cluster.ClusterModule; +import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery; +import org.apache.skywalking.oap.server.core.cluster.RemoteInstance; import org.apache.skywalking.oap.server.core.config.DownsamplingConfigService; -import org.apache.skywalking.oap.server.core.storage.*; -import org.apache.skywalking.oap.server.core.storage.model.*; +import org.apache.skywalking.oap.server.core.storage.Downsampling; +import org.apache.skywalking.oap.server.core.storage.IHistoryDeleteDAO; +import org.apache.skywalking.oap.server.core.storage.StorageModule; +import org.apache.skywalking.oap.server.core.storage.model.IModelGetter; +import org.apache.skywalking.oap.server.core.storage.model.Model; import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.apache.skywalking.oap.server.library.util.CollectionUtils; import org.joda.time.DateTime; -import org.slf4j.*; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng diff --git a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicatorTest.java b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicatorTest.java index 07d48e3d7..4433692ff 100644 --- a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicatorTest.java +++ b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/ThermodynamicIndicatorTest.java @@ -50,13 +50,12 @@ public class ThermodynamicIndicatorTest { indicatorMocker.combine(100, step, maxNumOfSteps); Map index = Whitebox.getInternalState(indicatorMocker, "detailIndex"); - Assert.assertEquals(5, index.size()); + Assert.assertEquals(4, index.size()); Assert.assertEquals(1, index.get(2).getValue()); Assert.assertEquals(3, index.get(5).getValue()); Assert.assertEquals(1, index.get(6).getValue()); - Assert.assertEquals(6, index.get(10).getValue()); - Assert.assertEquals(2, index.get(11).getValue()); + Assert.assertEquals(8, index.get(10).getValue()); } @Test @@ -83,13 +82,12 @@ public class ThermodynamicIndicatorTest { indicatorMocker.combine(indicatorMocker2); Map index = Whitebox.getInternalState(indicatorMocker, "detailIndex"); - Assert.assertEquals(5, index.size()); + Assert.assertEquals(4, index.size()); Assert.assertEquals(1, index.get(2).getValue()); Assert.assertEquals(3, index.get(5).getValue()); Assert.assertEquals(1, index.get(6).getValue()); - Assert.assertEquals(6, index.get(10).getValue()); - Assert.assertEquals(2, index.get(11).getValue()); + Assert.assertEquals(8, index.get(10).getValue()); } public class ThermodynamicIndicatorMocker extends ThermodynamicIndicator { diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/handler/IstioTelemetryHandlerMainTest.java b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/handler/IstioTelemetryHandlerMainTest.java new file mode 100644 index 000000000..79e501f48 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/handler/IstioTelemetryHandlerMainTest.java @@ -0,0 +1,86 @@ +/* + * 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.oap.server.receiver.istio.telemetry.handler; + +import com.google.protobuf.TextFormat; +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import io.istio.HandleMetricServiceGrpc; +import io.istio.IstioMetricProto; +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStream; +import java.io.InputStreamReader; +import java.util.LinkedList; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + +public class IstioTelemetryHandlerMainTest { + + public static void main(String[] args) throws InterruptedException { + ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).usePlaintext(true).build(); + + HandleMetricServiceGrpc.HandleMetricServiceBlockingStub stub = HandleMetricServiceGrpc.newBlockingStub(channel); + + ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); + executor.schedule(() -> { + try { + send(stub); + } catch (IOException e) { + e.printStackTrace(); + } + }, 1, TimeUnit.SECONDS); + Thread.sleep(5000L); + executor.shutdown(); + } + + private static void send(final HandleMetricServiceGrpc.HandleMetricServiceBlockingStub stub) throws IOException { + for (String s : readData()) { + IstioMetricProto.HandleMetricRequest.Builder requestBuilder = IstioMetricProto.HandleMetricRequest.newBuilder(); + try (InputStreamReader isr = new InputStreamReader(getResourceAsStream(String.format("fixture/%s", s)))) { + TextFormat.getParser().merge(isr, requestBuilder); + } + stub.handleMetric(requestBuilder.build()); + } + } + + private static Iterable readData() throws IOException { + Iterable result = new LinkedList<>(); + try ( + InputStream in = getResourceAsStream("fixture"); + BufferedReader br = new BufferedReader(new InputStreamReader(in))) { + String resource; + + while ((resource = br.readLine()) != null) { + ((LinkedList)result).add(resource); + } + } + return result; + } + + private static InputStream getResourceAsStream(final String resource) { + final InputStream in = getContextClassLoader().getResourceAsStream(resource); + return in == null ? IstioTelemetryHandlerMainTest.class.getResourceAsStream(resource) : in; + } + + private static ClassLoader getContextClassLoader() { + return Thread.currentThread().getContextClassLoader(); + } +} diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/01-ingress-reviewsv1-policy.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/01-ingress-reviewsv1-policy.msg new file mode 100644 index 000000000..9d0ea7fb6 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/01-ingress-reviewsv1-policy.msg @@ -0,0 +1,193 @@ +# 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. + +instances { + value { + int64_value: 948 + } + dimensions { + key: "sourceService" + value { + string_value: "istio-ingressgateway" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://istio-ingressgateway-6f58fdc8d7-m29dq.istio-system" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Check" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 883239533 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 882149494 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-8p5sr.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +instances { + value { + int64_value: 688 + } + dimensions { + key: "sourceService" + value { + string_value: "reviews-v1" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://reviews-v1-59cbdd7959-g69ll.default" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Check" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 905792473 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 904701150 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-8p5sr.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296801" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/02-productpage-details-policy.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/02-productpage-details-policy.msg new file mode 100644 index 000000000..c2ff197c0 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/02-productpage-details-policy.msg @@ -0,0 +1,193 @@ +# 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. + +instances { + value { + int64_value: 940 + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Check" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 887512227 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 886307886 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-dt82j.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +instances { + value { + int64_value: 689 + } + dimensions { + key: "sourceService" + value { + string_value: "details-v1" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://details-v1-7bcdcc4fd6-qb42z.default" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Check" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 895223970 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 894249545 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-dt82j.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296802" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/03-productpage-details.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/03-productpage-details.msg new file mode 100644 index 000000000..c2ff197c0 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/03-productpage-details.msg @@ -0,0 +1,193 @@ +# 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. + +instances { + value { + int64_value: 940 + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Check" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 887512227 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 886307886 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-dt82j.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +instances { + value { + int64_value: 689 + } + dimensions { + key: "sourceService" + value { + string_value: "details-v1" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://details-v1-7bcdcc4fd6-qb42z.default" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Check" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 895223970 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 894249545 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-dt82j.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296802" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/04-productpage-details.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/04-productpage-details.msg new file mode 100644 index 000000000..57de927a8 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/04-productpage-details.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 0 + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/details/0" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "GET" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 899188567 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 892490629 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "details-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://details-v1-7bcdcc4fd6-qb42z.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "source" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296803" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/05-productpage-reviews.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/05-productpage-reviews.msg new file mode 100644 index 000000000..bacdebe9f --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/05-productpage-reviews.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 0 + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 911565364 + } + } + } + } + dimensions { + key: "requestPath" + value { + string_value: "/reviews/0" + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://reviews-v1-59cbdd7959-g69ll.default" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 904205514 + } + } + } + } + dimensions { + key: "destinationService" + value { + string_value: "reviews-v1" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "GET" + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9219286855927306074" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/06-productpage-reviews.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/06-productpage-reviews.msg new file mode 100644 index 000000000..5c4b80565 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/06-productpage-reviews.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 0 + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 912063918 + } + } + } + } + dimensions { + key: "requestPath" + value { + string_value: "/reviews/0" + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://reviews-v1-59cbdd7959-g69ll.default" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 903831844 + } + } + } + } + dimensions { + key: "destinationService" + value { + string_value: "reviews-v1" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "GET" + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "source" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9219286855927306075" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/07-ingress-productpage.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/07-ingress-productpage.msg new file mode 100644 index 000000000..dbf1f7ab0 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/07-ingress-productpage.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 0 + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 915154754 + } + } + } + } + dimensions { + key: "requestPath" + value { + string_value: "/productpage" + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "sourceService" + value { + string_value: "istio-ingressgateway" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 884254913 + } + } + } + } + dimensions { + key: "destinationService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "GET" + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://istio-ingressgateway-6f58fdc8d7-m29dq.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9219286855927306076" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/08-ingress-productpage.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/08-ingress-productpage.msg new file mode 100644 index 000000000..de185c2dd --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/08-ingress-productpage.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 0 + } + dimensions { + key: "sourceService" + value { + string_value: "istio-ingressgateway" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://istio-ingressgateway-6f58fdc8d7-m29dq.istio-system" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/productpage" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "GET" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 915765768 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415799 + nanos: 881503925 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "source" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296804" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/09-policy-telemetry.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/09-policy-telemetry.msg new file mode 100644 index 000000000..5eecca790 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/09-policy-telemetry.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 1346 + } + dimensions { + key: "sourceService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-8p5sr.istio-system" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 887532113 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 882648986 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-bs9v4.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296805" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/10-policy-telemetry.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/10-policy-telemetry.msg new file mode 100644 index 000000000..3a2adf0a9 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/10-policy-telemetry.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 1321 + } + dimensions { + key: "sourceService" + value { + string_value: "istio-policy" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://istio-policy-7fbd997765-dt82j.istio-system" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 892667623 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 888995363 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-bs9v4.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296806" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/11-details-telemetry.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/11-details-telemetry.msg new file mode 100644 index 000000000..d5d6255be --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/11-details-telemetry.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 814 + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 903133331 + } + } + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "sourceService" + value { + string_value: "details-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-vx5dl.istio-system" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 899051593 + } + } + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://details-v1-7bcdcc4fd6-qb42z.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9219286855927306077" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/12-ingress-productpage-telemetry.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/12-ingress-productpage-telemetry.msg new file mode 100644 index 000000000..28a53fc08 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/12-ingress-productpage-telemetry.msg @@ -0,0 +1,193 @@ +# 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. + +instances { + value { + int64_value: 858 + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 903322857 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 900194688 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-bs9v4.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +instances { + value { + int64_value: 1154 + } + dimensions { + key: "sourceService" + value { + string_value: "istio-ingressgateway" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://istio-ingressgateway-6f58fdc8d7-m29dq.istio-system" + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 919832087 + } + } + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 916615524 + } + } + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-bs9v4.istio-system" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9307733128061296807" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/13-productpage-telemetry.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/13-productpage-telemetry.msg new file mode 100644 index 000000000..0d151c4ef --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/13-productpage-telemetry.msg @@ -0,0 +1,105 @@ +# 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. + +instances { + value { + int64_value: 904 + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 917907619 + } + } + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-vx5dl.istio-system" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 914305682 + } + } + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9219286855927306078" diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/14-reviews-productpage-telemetry.msg b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/14-reviews-productpage-telemetry.msg new file mode 100644 index 000000000..72b7ad007 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/test/resources/fixture/14-reviews-productpage-telemetry.msg @@ -0,0 +1,193 @@ +# 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. + +instances { + value { + int64_value: 860 + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 917979066 + } + } + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "sourceService" + value { + string_value: "reviews-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-vx5dl.istio-system" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 913319273 + } + } + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://reviews-v1-59cbdd7959-g69ll.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + name: "swmetric.instance.istio-system" +} +instances { + value { + int64_value: 1068 + } + dimensions { + key: "responseTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 918964543 + } + } + } + } + dimensions { + key: "requestPath" + value { + string_value: "/istio.mixer.v1.Mixer/Report" + } + } + dimensions { + key: "responseCode" + value { + int64_value: 200 + } + } + dimensions { + key: "sourceService" + value { + string_value: "productpage-v1" + } + } + dimensions { + key: "destinationUID" + value { + string_value: "kubernetes://istio-telemetry-796dbc5d46-vx5dl.istio-system" + } + } + dimensions { + key: "requestTime" + value { + timestamp_value { + value { + seconds: 1537415800 + nanos: 916174445 + } + } + } + } + dimensions { + key: "destinationService" + value { + string_value: "istio-telemetry" + } + } + dimensions { + key: "requestMethod" + value { + string_value: "POST" + } + } + dimensions { + key: "apiProtocol" + value { + string_value: "" + } + } + dimensions { + key: "sourceUID" + value { + string_value: "kubernetes://productpage-v1-8584c875d8-fq8wh.default" + } + } + dimensions { + key: "reporter" + value { + string_value: "destination" + } + } + dimensions { + key: "requestScheme" + value { + string_value: "http" + } + } + name: "swmetric.instance.istio-system" +} +dedup_id: "9219286855927306079" diff --git a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshDataBufferFileCache.java b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshDataBufferFileCache.java index 2ea9237ca..8089fe23a 100644 --- a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshDataBufferFileCache.java +++ b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshDataBufferFileCache.java @@ -33,7 +33,7 @@ public class MeshDataBufferFileCache implements IConsumer(3, 1024); + dataCarrier = new DataCarrier<>("MeshDataBufferFileCache", 3, 1024); } void start() throws IOException { diff --git a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java index 4dfcc4e19..70b755a9f 100644 --- a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java +++ b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java @@ -24,6 +24,7 @@ import org.apache.skywalking.apm.network.servicemesh.ServiceMeshMetric; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.cache.ServiceInstanceInventoryCache; import org.apache.skywalking.oap.server.core.cache.ServiceInventoryCache; +import org.apache.skywalking.oap.server.core.source.All; import org.apache.skywalking.oap.server.core.source.DetectPoint; import org.apache.skywalking.oap.server.core.source.Endpoint; import org.apache.skywalking.oap.server.core.source.RequestType; @@ -59,7 +60,12 @@ public class TelemetryDataDispatcher { } public static void preProcess(ServiceMeshMetric data) { - CACHE.in(data); + ServiceMeshMetricDataDecorator decorator = new ServiceMeshMetricDataDecorator(data); + if (decorator.tryMetaDataRegister()) { + TelemetryDataDispatcher.doDispatch(decorator); + } else { + CACHE.in(data); + } } /** @@ -70,11 +76,29 @@ public class TelemetryDataDispatcher { static void doDispatch(ServiceMeshMetricDataDecorator decorator) { ServiceMeshMetric metric = decorator.getMetric(); long minuteTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(metric.getStartTime()); - toService(decorator, minuteTimeBucket); + + if (org.apache.skywalking.apm.network.common.DetectPoint.server.equals(metric.getDetectPoint())) { + toAll(decorator, minuteTimeBucket); + toService(decorator, minuteTimeBucket); + toEndpoint(decorator, minuteTimeBucket); + } toServiceRelation(decorator, minuteTimeBucket); toServiceInstance(decorator, minuteTimeBucket); toServiceInstanceRelation(decorator, minuteTimeBucket); - toEndpoint(decorator, minuteTimeBucket); + } + + private static void toAll(ServiceMeshMetricDataDecorator decorator, long minuteTimeBucket) { + ServiceMeshMetric metric = decorator.getMetric(); + All all = new All(); + all.setTimeBucket(minuteTimeBucket); + all.setName(getServiceName(metric.getDestServiceId(), metric.getDestServiceName())); + all.setServiceInstanceName(getServiceInstanceName(metric.getDestServiceInstanceId(), metric.getDestServiceInstance())); + all.setEndpointName(metric.getEndpoint()); + all.setLatency(metric.getLatency()); + all.setStatus(metric.getStatus()); + all.setType(protocol2Type(metric.getProtocol())); + + SOURCE_RECEIVER.receive(all); } private static void toService(ServiceMeshMetricDataDecorator decorator, long minuteTimeBucket) { @@ -110,6 +134,7 @@ public class TelemetryDataDispatcher { serviceRelation.setType(protocol2Type(metric.getProtocol())); serviceRelation.setResponseCode(metric.getResponseCode()); serviceRelation.setDetectPoint(detectPointMapping(metric.getDetectPoint())); + serviceRelation.setComponentId(protocol2Component(metric.getProtocol())); SOURCE_RECEIVER.receive(serviceRelation); } @@ -150,6 +175,7 @@ public class TelemetryDataDispatcher { serviceRelation.setType(protocol2Type(metric.getProtocol())); serviceRelation.setResponseCode(metric.getResponseCode()); serviceRelation.setDetectPoint(detectPointMapping(metric.getDetectPoint())); + serviceRelation.setComponentId(protocol2Component(metric.getProtocol())); SOURCE_RECEIVER.receive(serviceRelation); } @@ -184,6 +210,21 @@ public class TelemetryDataDispatcher { } } + private static int protocol2Component(Protocol protocol) { + switch (protocol) { + case gRPC: + // GRPC in component-libraries.yml + return 23; + case HTTP: + // HTTP in component-libraries.yml + return 49; + case UNRECOGNIZED: + default: + // RPC in component-libraries.yml + return 50; + } + } + private static DetectPoint detectPointMapping(org.apache.skywalking.apm.network.common.DetectPoint detectPoint) { switch (detectPoint) { case client: diff --git a/oap-server/server-starter/src/main/resources/component-libraries.yml b/oap-server/server-starter/src/main/resources/component-libraries.yml index 7291c7f7e..0878d4db9 100644 --- a/oap-server/server-starter/src/main/resources/component-libraries.yml +++ b/oap-server/server-starter/src/main/resources/component-libraries.yml @@ -171,6 +171,12 @@ Elasticsearch: transport-client: id: 48 languages: Java +http: + id: 49 + languages: Java,C#,Node.js +rpc: + id: 50 + languages: Java,C#,Node.js # .NET/.NET Core components # [3000, 4000) for C#/.NET only diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/BatchProcessEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/BatchProcessEsDAO.java index 0e9ff2fe7..c00a7755d 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/BatchProcessEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/BatchProcessEsDAO.java @@ -68,5 +68,7 @@ public class BatchProcessEsDAO extends EsDAO implements IBatchDAO { } }); } + + this.bulkProcessor.flush(); } } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetricQueryEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetricQueryEsDAO.java index d40f533f2..ad897b642 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetricQueryEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetricQueryEsDAO.java @@ -150,7 +150,7 @@ public class MetricQueryEsDAO extends EsDAO implements IMetricQueryDAO { } for (IntKeyLongValue intKeyLongValue : intKeyLongValues) { - axisYValues.set(intKeyLongValue.getKey() - 1, intKeyLongValue.getValue()); + axisYValues.set(intKeyLongValue.getKey(), intKeyLongValue.getValue()); } thermodynamicValueMatrix.add(axisYValues);