parent
ecc050c63d
commit
258a870c89
|
|
@ -13,6 +13,11 @@
|
|||
<packaging>jar</packaging>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.skywalking</groupId>
|
||||
<artifactId>apm-collector-stream</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.skywalking</groupId>
|
||||
<artifactId>apm-collector-cluster</artifactId>
|
||||
|
|
@ -25,7 +30,7 @@
|
|||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.skywalking</groupId>
|
||||
<artifactId>apm-network</artifactId>
|
||||
<artifactId>apm-collector-storage</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
|
|
|
|||
|
|
@ -1,17 +1,133 @@
|
|||
package org.skywalking.apm.collector.agentjvm.grpc.handler;
|
||||
|
||||
import io.grpc.stub.StreamObserver;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.cpu.CpuMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.cpu.define.CpuMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.gc.GCMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.gc.define.GCMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memory.MemoryMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memory.define.MemoryMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memorypool.MemoryPoolMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memorypool.define.MemoryPoolMetricDataDefine;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.server.grpc.GRPCHandler;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
import org.skywalking.apm.collector.stream.worker.WorkerInvokeException;
|
||||
import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.network.proto.CPU;
|
||||
import org.skywalking.apm.network.proto.Downstream;
|
||||
import org.skywalking.apm.network.proto.GC;
|
||||
import org.skywalking.apm.network.proto.JVMMetrics;
|
||||
import org.skywalking.apm.network.proto.JVMMetricsServiceGrpc;
|
||||
import org.skywalking.apm.network.proto.Memory;
|
||||
import org.skywalking.apm.network.proto.MemoryPool;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsServiceImplBase implements GRPCHandler {
|
||||
|
||||
private final Logger logger = LoggerFactory.getLogger(JVMMetricsServiceHandler.class);
|
||||
|
||||
@Override public void collect(JVMMetrics request, StreamObserver<Downstream> responseObserver) {
|
||||
super.collect(request, responseObserver);
|
||||
int applicationInstanceId = request.getApplicationInstanceId();
|
||||
logger.debug("receive the jvm metric from application instance, id: {}", applicationInstanceId);
|
||||
|
||||
StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME);
|
||||
request.getMetricsList().forEach(metric -> {
|
||||
long time = TimeBucketUtils.INSTANCE.getSecondTimeBucket(metric.getTime());
|
||||
sendToCpuMetricPersistenceWorker(context, applicationInstanceId, time, metric.getCpu());
|
||||
sendToMemoryMetricPersistenceWorker(context, applicationInstanceId, time, metric.getMemoryList());
|
||||
sendToMemoryPoolMetricPersistenceWorker(context, applicationInstanceId, time, metric.getMemoryPoolList());
|
||||
sendToGCMetricPersistenceWorker(context, applicationInstanceId, time, metric.getGcList());
|
||||
});
|
||||
|
||||
responseObserver.onNext(Downstream.newBuilder().build());
|
||||
responseObserver.onCompleted();
|
||||
}
|
||||
|
||||
private void sendToCpuMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, CPU cpu) {
|
||||
CpuMetricDataDefine.CpuMetric cpuMetric = new CpuMetricDataDefine.CpuMetric();
|
||||
cpuMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId);
|
||||
cpuMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
cpuMetric.setUsagePercent(cpu.getUsagePercent());
|
||||
cpuMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to cpu metric persistence worker, id: {}", cpuMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(CpuMetricPersistenceWorker.WorkerRole.INSTANCE).tell(cpuMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
|
||||
private void sendToMemoryMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, List<Memory> memories) {
|
||||
|
||||
memories.forEach(memory -> {
|
||||
MemoryMetricDataDefine.MemoryMetric memoryMetric = new MemoryMetricDataDefine.MemoryMetric();
|
||||
memoryMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + String.valueOf(memory.getIsHeap()));
|
||||
memoryMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
memoryMetric.setHeap(memory.getIsHeap());
|
||||
memoryMetric.setInit(memory.getInit());
|
||||
memoryMetric.setMax(memory.getMax());
|
||||
memoryMetric.setUsed(memory.getUsed());
|
||||
memoryMetric.setCommitted(memory.getCommitted());
|
||||
memoryMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to memory metric persistence worker, id: {}", memoryMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(MemoryMetricPersistenceWorker.WorkerRole.INSTANCE).tell(memoryMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private void sendToMemoryPoolMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, List<MemoryPool> memoryPools) {
|
||||
|
||||
memoryPools.forEach(memoryPool -> {
|
||||
MemoryPoolMetricDataDefine.MemoryPoolMetric memoryPoolMetric = new MemoryPoolMetricDataDefine.MemoryPoolMetric();
|
||||
memoryPoolMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + String.valueOf(memoryPool.getType().getNumber()));
|
||||
memoryPoolMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
memoryPoolMetric.setPoolType(memoryPool.getType().getNumber());
|
||||
memoryPoolMetric.setHeap(memoryPool.getIsHeap());
|
||||
memoryPoolMetric.setInit(memoryPool.getInit());
|
||||
memoryPoolMetric.setMax(memoryPool.getMax());
|
||||
memoryPoolMetric.setUsed(memoryPool.getUsed());
|
||||
memoryPoolMetric.setCommitted(memoryPool.getCommited());
|
||||
memoryPoolMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to memory pool metric persistence worker, id: {}", memoryPoolMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(MemoryPoolMetricPersistenceWorker.WorkerRole.INSTANCE).tell(memoryPoolMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private void sendToGCMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, List<GC> gcs) {
|
||||
gcs.forEach(gc -> {
|
||||
GCMetricDataDefine.GCMetric gcMetric = new GCMetricDataDefine.GCMetric();
|
||||
gcMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + String.valueOf(gc.getPhraseValue()));
|
||||
gcMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
gcMetric.setPhrase(gc.getPhraseValue());
|
||||
gcMetric.setCount(gc.getCount());
|
||||
gcMetric.setTime(gc.getTime());
|
||||
gcMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to gc metric persistence worker, id: {}", gcMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(GCMetricPersistenceWorker.WorkerRole.INSTANCE).tell(gcMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,7 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.cpu;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.dao.ICpuMetricDAO;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define.CpuMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.cpu.dao.ICpuMetricDAO;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.cpu.define.CpuMetricDataDefine;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider;
|
||||
import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext;
|
||||
|
|
@ -1,10 +1,10 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.cpu.dao;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.elasticsearch.action.index.IndexRequestBuilder;
|
||||
import org.elasticsearch.action.update.UpdateRequestBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define.CpuMetricTable;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.cpu.define.CpuMetricTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Data;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.cpu.dao;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
|
||||
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.cpu.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface ICpuMetricDAO {
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.cpu.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Attribute;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.cpu.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.cpu.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.cpu.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,7 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.gc;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.dao.IGCMetricDAO;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define.GCMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.gc.dao.IGCMetricDAO;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.gc.define.GCMetricDataDefine;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider;
|
||||
import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext;
|
||||
|
|
@ -1,10 +1,10 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.gc.dao;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.elasticsearch.action.index.IndexRequestBuilder;
|
||||
import org.elasticsearch.action.update.UpdateRequestBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define.GCMetricTable;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.gc.define.GCMetricTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Data;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.gc.dao;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
|
||||
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.gc.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IGCMetricDAO {
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.gc.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Attribute;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.gc.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.gc.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.gc.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,7 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memory;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.dao.IMemoryMetricDAO;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define.MemoryMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memory.dao.IMemoryMetricDAO;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memory.define.MemoryMetricDataDefine;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider;
|
||||
import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext;
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.memory.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IMemoryMetricDAO {
|
||||
}
|
||||
|
|
@ -1,10 +1,10 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memory.dao;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.elasticsearch.action.index.IndexRequestBuilder;
|
||||
import org.elasticsearch.action.update.UpdateRequestBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define.MemoryMetricTable;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memory.define.MemoryMetricTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Data;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memory.dao;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
|
||||
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memory.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Attribute;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memory.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memory.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memory.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,7 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memorypool;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.dao.IMemoryPoolMetricDAO;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define.MemoryPoolMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memorypool.dao.IMemoryPoolMetricDAO;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memorypool.define.MemoryPoolMetricDataDefine;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider;
|
||||
import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext;
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IMemoryPoolMetricDAO {
|
||||
}
|
||||
|
|
@ -1,10 +1,10 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.dao;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.elasticsearch.action.index.IndexRequestBuilder;
|
||||
import org.elasticsearch.action.update.UpdateRequestBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define.MemoryPoolMetricTable;
|
||||
import org.skywalking.apm.collector.agentjvm.worker.memorypool.define.MemoryPoolMetricTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Data;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.dao;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.dao;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
|
||||
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Attribute;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.define;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define;
|
||||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -0,0 +1,4 @@
|
|||
org.skywalking.apm.collector.agentjvm.worker.cpu.dao.CpuMetricEsDAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.memory.dao.MemoryMetricEsDAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.memorypool.dao.MemoryPoolMetricEsDAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.gc.dao.GCMetricEsDAO
|
||||
|
|
@ -0,0 +1,4 @@
|
|||
org.skywalking.apm.collector.agentjvm.worker.cpu.dao.CpuMetricH2DAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.memory.dao.MemoryMetricH2DAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.memorypool.dao.MemoryPoolMetricH2DAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.gc.dao.GCMetricH2DAO
|
||||
|
|
@ -0,0 +1,4 @@
|
|||
org.skywalking.apm.collector.agentjvm.worker.cpu.CpuMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentjvm.worker.memory.MemoryMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentjvm.worker.memorypool.MemoryPoolMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentjvm.worker.gc.GCMetricPersistenceWorker$Factory
|
||||
|
|
@ -0,0 +1,11 @@
|
|||
org.skywalking.apm.collector.agentjvm.worker.cpu.define.CpuMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentjvm.worker.cpu.define.CpuMetricH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentjvm.worker.memory.define.MemoryMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentjvm.worker.memory.define.MemoryMetricH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentjvm.worker.memorypool.define.MemoryPoolMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentjvm.worker.memorypool.define.MemoryPoolMetricH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentjvm.worker.gc.define.GCMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentjvm.worker.gc.define.GCMetricH2TableDefine
|
||||
|
|
@ -0,0 +1,88 @@
|
|||
package org.skywalking.apm.collector.agentjvm.grpc.handler;
|
||||
|
||||
import io.grpc.ManagedChannel;
|
||||
import io.grpc.ManagedChannelBuilder;
|
||||
import org.junit.Test;
|
||||
import org.skywalking.apm.network.proto.CPU;
|
||||
import org.skywalking.apm.network.proto.GC;
|
||||
import org.skywalking.apm.network.proto.GCPhrase;
|
||||
import org.skywalking.apm.network.proto.JVMMetric;
|
||||
import org.skywalking.apm.network.proto.JVMMetrics;
|
||||
import org.skywalking.apm.network.proto.JVMMetricsServiceGrpc;
|
||||
import org.skywalking.apm.network.proto.Memory;
|
||||
import org.skywalking.apm.network.proto.MemoryPool;
|
||||
import org.skywalking.apm.network.proto.PoolType;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class JVMMetricsServiceHandlerTestCase {
|
||||
|
||||
private final Logger logger = LoggerFactory.getLogger(JVMMetricsServiceHandlerTestCase.class);
|
||||
|
||||
private JVMMetricsServiceGrpc.JVMMetricsServiceBlockingStub stub;
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).usePlaintext(true).build();
|
||||
stub = JVMMetricsServiceGrpc.newBlockingStub(channel);
|
||||
|
||||
JVMMetrics.Builder jvmMetricsBuilder = JVMMetrics.newBuilder();
|
||||
jvmMetricsBuilder.setApplicationInstanceId(1);
|
||||
|
||||
JVMMetric.Builder jvmMetric = JVMMetric.newBuilder();
|
||||
jvmMetric.setTime(System.currentTimeMillis());
|
||||
buildCpuMetric(jvmMetric);
|
||||
buildMemoryMetric(jvmMetric);
|
||||
buildMemoryPoolMetric(jvmMetric);
|
||||
buildGcMetric(jvmMetric);
|
||||
|
||||
jvmMetricsBuilder.addMetrics(jvmMetric.build());
|
||||
stub.collect(jvmMetricsBuilder.build());
|
||||
}
|
||||
|
||||
private void buildCpuMetric(JVMMetric.Builder jvmMetric) {
|
||||
CPU.Builder cpuBuilder = CPU.newBuilder();
|
||||
cpuBuilder.setUsagePercent(70);
|
||||
jvmMetric.setCpu(cpuBuilder);
|
||||
}
|
||||
|
||||
private void buildMemoryMetric(JVMMetric.Builder jvmMetric) {
|
||||
Memory.Builder builder_1 = Memory.newBuilder();
|
||||
builder_1.setIsHeap(true);
|
||||
builder_1.setInit(20);
|
||||
builder_1.setMax(100);
|
||||
builder_1.setUsed(50);
|
||||
builder_1.setCommitted(30);
|
||||
jvmMetric.addMemory(builder_1.build());
|
||||
|
||||
Memory.Builder builder_2 = Memory.newBuilder();
|
||||
builder_2.setIsHeap(false);
|
||||
builder_2.setInit(200);
|
||||
builder_2.setMax(1000);
|
||||
builder_2.setUsed(500);
|
||||
builder_2.setCommitted(300);
|
||||
jvmMetric.addMemory(builder_2.build());
|
||||
}
|
||||
|
||||
private void buildMemoryPoolMetric(JVMMetric.Builder jvmMetric) {
|
||||
MemoryPool.Builder builder_1 = MemoryPool.newBuilder();
|
||||
builder_1.setType(PoolType.NEWGEN_USAGE);
|
||||
builder_1.setIsHeap(true);
|
||||
builder_1.setInit(20);
|
||||
builder_1.setMax(100);
|
||||
builder_1.setUsed(50);
|
||||
builder_1.setCommited(30);
|
||||
jvmMetric.addMemoryPool(builder_1.build());
|
||||
}
|
||||
|
||||
private void buildGcMetric(JVMMetric.Builder jvmMetric) {
|
||||
GC.Builder gcBuilder = GC.newBuilder();
|
||||
gcBuilder.setPhrase(GCPhrase.NEW);
|
||||
gcBuilder.setCount(2);
|
||||
gcBuilder.setTime(100);
|
||||
jvmMetric.addGc(gcBuilder.build());
|
||||
}
|
||||
}
|
||||
|
|
@ -4,7 +4,6 @@ import java.util.LinkedList;
|
|||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine;
|
||||
import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine;
|
||||
import org.skywalking.apm.collector.agentstream.grpc.handler.JVMMetricsServiceHandler;
|
||||
import org.skywalking.apm.collector.agentstream.grpc.handler.TraceSegmentServiceHandler;
|
||||
import org.skywalking.apm.collector.core.cluster.ClusterDataListener;
|
||||
import org.skywalking.apm.collector.core.framework.Handler;
|
||||
|
|
@ -47,7 +46,6 @@ public class AgentStreamGRPCModuleDefine extends AgentStreamModuleDefine {
|
|||
@Override public List<Handler> handlerList() {
|
||||
List<Handler> handlers = new LinkedList<>();
|
||||
handlers.add(new TraceSegmentServiceHandler());
|
||||
handlers.add(new JVMMetricsServiceHandler());
|
||||
return handlers;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,133 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.grpc.handler;
|
||||
|
||||
import io.grpc.stub.StreamObserver;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.CpuMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define.CpuMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.GCMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define.GCMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.MemoryMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define.MemoryMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.MemoryPoolMetricPersistenceWorker;
|
||||
import org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define.MemoryPoolMetricDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.server.grpc.GRPCHandler;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
import org.skywalking.apm.collector.stream.worker.WorkerInvokeException;
|
||||
import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException;
|
||||
import org.skywalking.apm.network.proto.CPU;
|
||||
import org.skywalking.apm.network.proto.Downstream;
|
||||
import org.skywalking.apm.network.proto.GC;
|
||||
import org.skywalking.apm.network.proto.JVMMetrics;
|
||||
import org.skywalking.apm.network.proto.JVMMetricsServiceGrpc;
|
||||
import org.skywalking.apm.network.proto.Memory;
|
||||
import org.skywalking.apm.network.proto.MemoryPool;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsServiceImplBase implements GRPCHandler {
|
||||
|
||||
private final Logger logger = LoggerFactory.getLogger(JVMMetricsServiceHandler.class);
|
||||
|
||||
@Override public void collect(JVMMetrics request, StreamObserver<Downstream> responseObserver) {
|
||||
int applicationInstanceId = request.getApplicationInstanceId();
|
||||
logger.debug("receive the jvm metric from application instance, id: {}", applicationInstanceId);
|
||||
|
||||
StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME);
|
||||
request.getMetricsList().forEach(metric -> {
|
||||
long time = TimeBucketUtils.INSTANCE.getSecondTimeBucket(metric.getTime());
|
||||
sendToCpuMetricPersistenceWorker(context, applicationInstanceId, time, metric.getCpu());
|
||||
sendToMemoryMetricPersistenceWorker(context, applicationInstanceId, time, metric.getMemoryList());
|
||||
sendToMemoryPoolMetricPersistenceWorker(context, applicationInstanceId, time, metric.getMemoryPoolList());
|
||||
sendToGCMetricPersistenceWorker(context, applicationInstanceId, time, metric.getGcList());
|
||||
});
|
||||
|
||||
responseObserver.onNext(Downstream.newBuilder().build());
|
||||
responseObserver.onCompleted();
|
||||
}
|
||||
|
||||
private void sendToCpuMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, CPU cpu) {
|
||||
CpuMetricDataDefine.CpuMetric cpuMetric = new CpuMetricDataDefine.CpuMetric();
|
||||
cpuMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId);
|
||||
cpuMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
cpuMetric.setUsagePercent(cpu.getUsagePercent());
|
||||
cpuMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to cpu metric persistence worker, id: {}", cpuMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(CpuMetricPersistenceWorker.WorkerRole.INSTANCE).tell(cpuMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
|
||||
private void sendToMemoryMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, List<Memory> memories) {
|
||||
|
||||
memories.forEach(memory -> {
|
||||
MemoryMetricDataDefine.MemoryMetric memoryMetric = new MemoryMetricDataDefine.MemoryMetric();
|
||||
memoryMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + String.valueOf(memory.getIsHeap()));
|
||||
memoryMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
memoryMetric.setHeap(memory.getIsHeap());
|
||||
memoryMetric.setInit(memory.getInit());
|
||||
memoryMetric.setMax(memory.getMax());
|
||||
memoryMetric.setUsed(memory.getUsed());
|
||||
memoryMetric.setCommitted(memory.getCommitted());
|
||||
memoryMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to memory metric persistence worker, id: {}", memoryMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(MemoryMetricPersistenceWorker.WorkerRole.INSTANCE).tell(memoryMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private void sendToMemoryPoolMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, List<MemoryPool> memoryPools) {
|
||||
|
||||
memoryPools.forEach(memoryPool -> {
|
||||
MemoryPoolMetricDataDefine.MemoryPoolMetric memoryPoolMetric = new MemoryPoolMetricDataDefine.MemoryPoolMetric();
|
||||
memoryPoolMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + String.valueOf(memoryPool.getType().getNumber()));
|
||||
memoryPoolMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
memoryPoolMetric.setPoolType(memoryPool.getType().getNumber());
|
||||
memoryPoolMetric.setHeap(memoryPool.getIsHeap());
|
||||
memoryPoolMetric.setInit(memoryPool.getInit());
|
||||
memoryPoolMetric.setMax(memoryPool.getMax());
|
||||
memoryPoolMetric.setUsed(memoryPool.getUsed());
|
||||
memoryPoolMetric.setCommitted(memoryPool.getCommited());
|
||||
memoryPoolMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to memory pool metric persistence worker, id: {}", memoryPoolMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(MemoryPoolMetricPersistenceWorker.WorkerRole.INSTANCE).tell(memoryPoolMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private void sendToGCMetricPersistenceWorker(StreamModuleContext context, int applicationInstanceId,
|
||||
long timeBucket, List<GC> gcs) {
|
||||
gcs.forEach(gc -> {
|
||||
GCMetricDataDefine.GCMetric gcMetric = new GCMetricDataDefine.GCMetric();
|
||||
gcMetric.setId(timeBucket + Const.ID_SPLIT + applicationInstanceId + Const.ID_SPLIT + String.valueOf(gc.getPhraseValue()));
|
||||
gcMetric.setApplicationInstanceId(applicationInstanceId);
|
||||
gcMetric.setPhrase(gc.getPhraseValue());
|
||||
gcMetric.setCount(gc.getCount());
|
||||
gcMetric.setTime(gc.getTime());
|
||||
gcMetric.setTimeBucket(timeBucket);
|
||||
try {
|
||||
logger.debug("send to gc metric persistence worker, id: {}", gcMetric.getId());
|
||||
context.getClusterWorkerContext().lookup(GCMetricPersistenceWorker.WorkerRole.INSTANCE).tell(gcMetric.toData());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.cache;
|
|||
|
||||
import com.google.common.cache.Cache;
|
||||
import com.google.common.cache.CacheBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ import java.util.List;
|
|||
import org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.GlobalTraceIdsListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.global.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,7 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface ICpuMetricDAO {
|
||||
}
|
||||
|
|
@ -1,7 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IGCMetricDAO {
|
||||
}
|
||||
|
|
@ -1,7 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IMemoryMetricDAO {
|
||||
}
|
||||
|
|
@ -1,7 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IMemoryPoolMetricDAO {
|
||||
}
|
||||
|
|
@ -2,14 +2,14 @@ package org.skywalking.apm.collector.agentstream.worker.node.component;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.LocalSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.component.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,12 +2,12 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
|
|
@ -10,7 +10,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener;
|
|||
import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
|
|
@ -10,7 +10,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener;
|
|||
import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.application;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.IdAutoIncrement;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.application.dao.IApplicationDAO;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.application;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.instance;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.servicename;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener
|
|||
import org.skywalking.apm.collector.agentstream.worker.segment.LocalSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.segment.cost.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.segment.origin.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,12 +1,12 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.service.entry;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.service.entry.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.serviceref.reference;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener;
|
||||
|
|
@ -10,7 +10,7 @@ import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener
|
|||
import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleContext;
|
||||
import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.CommonTable;
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -9,8 +9,4 @@ org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.NodeRefSumEs
|
|||
org.skywalking.apm.collector.agentstream.worker.segment.cost.dao.SegmentCostEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.service.entry.dao.ServiceEntryEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.serviceref.reference.dao.ServiceRefEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.dao.CpuMetricEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.dao.MemoryMetricEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.dao.MemoryPoolMetricEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.dao.GCMetricEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.serviceref.reference.dao.ServiceRefEsDAO
|
||||
|
|
@ -10,7 +10,7 @@ org.skywalking.apm.collector.agentstream.worker.segment.cost.dao.SegmentCostH2DA
|
|||
org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.service.entry.dao.ServiceEntryH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.serviceref.reference.dao.ServiceRefH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.dao.CpuMetricH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.dao.MemoryMetricH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.dao.MemoryPoolMetricH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.dao.GCMetricH2DAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.cpu.dao.CpuMetricH2DAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.memory.dao.MemoryMetricH2DAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.memorypool.dao.MemoryPoolMetricH2DAO
|
||||
org.skywalking.apm.collector.agentjvm.worker.gc.dao.GCMetricH2DAO
|
||||
|
|
@ -20,10 +20,10 @@ org.skywalking.apm.collector.agentstream.worker.segment.origin.SegmentPersistenc
|
|||
org.skywalking.apm.collector.agentstream.worker.segment.cost.SegmentCostPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.global.GlobalTracePersistenceWorker$Factory
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.CpuMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.MemoryMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.MemoryPoolMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.GCMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentjvm.worker.cpu.CpuMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentjvm.worker.memory.MemoryMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentjvm.worker.memorypool.MemoryPoolMetricPersistenceWorker$Factory
|
||||
org.skywalking.apm.collector.agentjvm.worker.gc.GCMetricPersistenceWorker$Factory
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterSerialWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterSerialWorker$Factory
|
||||
|
|
|
|||
|
|
@ -32,16 +32,4 @@ org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntr
|
|||
org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefEsTableDefine
|
||||
org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define.CpuMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.cpu.define.CpuMetricH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define.MemoryMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memory.define.MemoryMetricH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define.MemoryPoolMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.memorypool.define.MemoryPoolMetricH2TableDefine
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define.GCMetricEsTableDefine
|
||||
org.skywalking.apm.collector.agentstream.worker.jvmmetric.gc.define.GCMetricH2TableDefine
|
||||
org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefH2TableDefine
|
||||
|
|
@ -4,7 +4,7 @@ import java.util.Calendar;
|
|||
import java.util.TimeZone;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker;
|
||||
package org.skywalking.apm.collector.stream.worker.storage;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker;
|
||||
package org.skywalking.apm.collector.stream.worker.util;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -6,8 +6,6 @@ package org.skywalking.apm.collector.agentstream.worker;
|
|||
public class Const {
|
||||
public static final String ID_SPLIT = "..-..";
|
||||
public static final String IDS_SPLIT = "\\.\\.-\\.\\.";
|
||||
public static final String PEERS_FRONT_SPLIT = "[";
|
||||
public static final String PEERS_BEHIND_SPLIT = "]";
|
||||
public static final int USER_ID = 1;
|
||||
public static final String USER_CODE = "User";
|
||||
public static final String SEGMENT_SPAN_SPLIT = "S";
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.util;
|
||||
package org.skywalking.apm.collector.stream.worker.util;
|
||||
|
||||
import java.text.SimpleDateFormat;
|
||||
import java.util.Calendar;
|
||||
|
|
@ -8,7 +8,7 @@ import org.elasticsearch.action.search.SearchType;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.aggregations.AggregationBuilders;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.slf4j.Logger;
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ import org.elasticsearch.action.search.SearchType;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.aggregations.AggregationBuilders;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.slf4j.Logger;
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ import org.elasticsearch.search.aggregations.AggregationBuilders;
|
|||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.TermsAggregationBuilder;
|
||||
import org.elasticsearch.search.aggregations.metrics.sum.Sum;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.slf4j.Logger;
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ import org.elasticsearch.action.search.SearchType;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.aggregations.AggregationBuilders;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.slf4j.Logger;
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ import com.google.gson.JsonArray;
|
|||
import com.google.gson.JsonObject;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ import com.google.gson.JsonArray;
|
|||
import com.google.gson.JsonObject;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.agentstream.worker.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.CollectionUtils;
|
||||
import org.skywalking.apm.collector.core.util.ObjectUtils;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
|
|
|
|||
Loading…
Reference in New Issue