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 4f5eabd60..daff12275 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 @@ -31,6 +31,10 @@ import org.apache.skywalking.apm.commons.datacarrier.partition.SimpleRollingPart /** * DataCarrier main class. use this instance to set Producer/Consumer Model. */ + +/** + * 数据缓冲层 + */ public class DataCarrier { private Channels channels; private IDriver driver; @@ -50,8 +54,11 @@ public class DataCarrier { public DataCarrier(String name, String envPrefix, int channelSize, int bufferSize, BufferStrategy strategy) { this.name = name; + // 每个分区的内存大小 bufferSize = EnvUtil.getInt(envPrefix + "_BUFFER_SIZE", bufferSize); + // 缓存的数据分区,提高并发度 channelSize = EnvUtil.getInt(envPrefix + "_CHANNEL_SIZE", channelSize); + // 数据分区 channels = new Channels<>(channelSize, bufferSize, new SimpleRollingPartitioner(), strategy); } @@ -126,6 +133,7 @@ public class DataCarrier { if (driver != null) { driver.close(channels); } + // 创建数据采集线程 driver = new ConsumeDriver(this.name, this.channels, consumer, num, consumeCycle); driver.begin(channels); return this; diff --git a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/buffer/Channels.java b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/buffer/Channels.java index b9ad4fd77..5f7137e7e 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/buffer/Channels.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/buffer/Channels.java @@ -33,6 +33,7 @@ public class Channels { public Channels(int channelSize, int bufferSize, IDataPartitioner partitioner, BufferStrategy strategy) { this.dataPartitioner = partitioner; this.strategy = strategy; + // 数据分区 bufferChannels = new QueueBuffer[channelSize]; for (int i = 0; i < channelSize; i++) { if (BufferStrategy.BLOCKING.equals(strategy)) { diff --git a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerThread.java b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerThread.java index 6c8d5ed80..2986205f6 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerThread.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/apache/skywalking/apm/commons/datacarrier/consumer/ConsumerThread.java @@ -50,6 +50,7 @@ public class ConsumerThread extends Thread { final List consumeList = new ArrayList(1500); while (running) { + // 消费缓冲数据 if (!consume(consumeList)) { try { Thread.sleep(consumeCycle); @@ -67,11 +68,13 @@ public class ConsumerThread extends Thread { private boolean consume(List consumeList) { for (DataSource dataSource : dataSources) { + // 将数据读取到consumeList dataSource.obtain(consumeList); } if (!consumeList.isEmpty()) { try { + // 数据消费 consumer.consume(consumeList); } catch (Throwable t) { consumer.onError(consumeList, t); @@ -91,6 +94,9 @@ public class ConsumerThread extends Thread { /** * DataSource is a reference to {@link Buffer}. */ + /** + * 存储数据 + */ class DataSource { private QueueBuffer sourceBuffer; diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java index 015596ad2..a14dc1f56 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java @@ -62,6 +62,7 @@ public enum ServiceManager { private Map loadAllServices() { Map bootedServices = new LinkedHashMap<>(); List allServices = new LinkedList<>(); + // 使用spi方式加载BootService load(allServices); for (final BootService bootService : allServices) { Class bootServiceClass = bootService.getClass(); diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java index cc1faa9af..cbfc687ba 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java @@ -467,6 +467,7 @@ public class TracingContext implements AbstractTracerContext { } AgentSo11y.measureTracingContextCompletion(false); TraceSegment finishedSegment = segment.finish(limitMechanismWorking); + // 通知TraceSegmentServiceClient上报数据 TracingContext.ListenerManager.notifyFinish(finishedSegment); running = false; } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java index 6c682ca41..8879f24ee 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java @@ -331,6 +331,9 @@ public abstract class AbstractTracingSpan implements AbstractSpan { return this; } + /** + * 异步存储span数据 + */ @Override public AbstractSpan asyncFinish() { if (!isInAsyncMode) { diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/jvm/JVMService.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/jvm/JVMService.java index ff6f6234e..8c4763bd2 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/jvm/JVMService.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/jvm/JVMService.java @@ -42,6 +42,10 @@ import org.apache.skywalking.apm.util.RunnableWithExceptionProtection; * The JVMService represents a timer, which collectors JVM cpu, memory, memorypool, gc, thread and class info, * and send the collected info to Collector through the channel provided by {@link GRPCChannelManager} */ + +/** + * 开始 JVM 指标收集 + */ @DefaultImplementor public class JVMService implements BootService, Runnable { private static final ILog LOGGER = LogManager.getLogger(JVMService.class); diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java index 9744398fd..5cd4e8882 100755 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java @@ -45,6 +45,9 @@ import org.apache.skywalking.apm.util.StringUtil; import static org.apache.skywalking.apm.agent.core.conf.Config.Collector.IS_RESOLVE_DNS_PERIODICALLY; +/** + * 建立OAP的gRPC连接 + */ @DefaultImplementor public class GRPCChannelManager implements BootService, Runnable { private static final ILog LOGGER = LogManager.getLogger(GRPCChannelManager.class); @@ -74,6 +77,7 @@ public class GRPCChannelManager implements BootService, Runnable { connectCheckFuture = Executors.newSingleThreadScheduledExecutor( new DefaultNamedThreadFactory("GRPCChannelManager") ).scheduleAtFixedRate( + // 心跳检查,skywalking通过随机选择一个oap节点并保持长连接的方式来代替负载均衡,如果连接不通重新选择一个节点 new RunnableWithExceptionProtection( this, t -> LOGGER.error("unexpected exception.", t) diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/ServiceManagementClient.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/ServiceManagementClient.java index f1299559b..7ef4bb291 100755 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/ServiceManagementClient.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/ServiceManagementClient.java @@ -44,6 +44,9 @@ import org.apache.skywalking.apm.util.RunnableWithExceptionProtection; import static org.apache.skywalking.apm.agent.core.conf.Config.Collector.GRPC_UPSTREAM_TIMEOUT; +/** + * 注册服务实例 + */ @DefaultImplementor public class ServiceManagementClient implements BootService, Runnable, GRPCChannelListener { private static final ILog LOGGER = LogManager.getLogger(ServiceManagementClient.class); diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java index a0b86957f..3f5b60275 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java @@ -44,6 +44,9 @@ import static org.apache.skywalking.apm.agent.core.conf.Config.Buffer.BUFFER_SIZ import static org.apache.skywalking.apm.agent.core.conf.Config.Buffer.CHANNEL_SIZE; import static org.apache.skywalking.apm.agent.core.remote.GRPCChannelStatus.CONNECTED; +/** + * 启动链路数据上报 + */ @DefaultImplementor public class TraceSegmentServiceClient implements BootService, IConsumer, TracingContextListener, GRPCChannelListener { private static final ILog LOGGER = LogManager.getLogger(TraceSegmentServiceClient.class); @@ -179,6 +182,7 @@ public class TraceSegmentServiceClient implements BootService, IConsumer