[Collector] Topology query tuning, Batch Process instead of bulk. (#1384)

* #1202

1. Determine the log is enabled for the DEBUG level before printing message.
2. Make the columns initialize to be static attribute.
3. Topology build optimize: Cache query result to avoid repeating queries.

* #1202

Add elasticsearch batch process setting into application.yml.

* #1202

Fixed check style error.

* #1202

Use XContentFactory to build source to insert into elasticsearch.
This commit is contained in:
彭勇升 pengys 2018-06-25 21:19:54 +08:00 committed by GitHub
parent 48dacdb3f8
commit 5e03ec8845
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
132 changed files with 1069 additions and 801 deletions

View File

@ -45,7 +45,9 @@ public class ApplicationRegisterServiceHandler extends ApplicationRegisterServic
@Override
public void applicationCodeRegister(Application request, StreamObserver<ApplicationMapping> responseObserver) {
logger.debug("register application");
if (logger.isDebugEnabled()) {
logger.debug("register application");
}
ApplicationMapping.Builder builder = ApplicationMapping.newBuilder();
String applicationCode = request.getApplicationCode();

View File

@ -21,24 +21,14 @@ package org.apache.skywalking.apm.collector.agent.grpc.provider.handler;
import io.grpc.stub.StreamObserver;
import java.util.List;
import org.apache.skywalking.apm.collector.analysis.jvm.define.AnalysisJVMModule;
import org.apache.skywalking.apm.collector.analysis.jvm.define.service.ICpuMetricService;
import org.apache.skywalking.apm.collector.analysis.jvm.define.service.IGCMetricService;
import org.apache.skywalking.apm.collector.analysis.jvm.define.service.IMemoryMetricService;
import org.apache.skywalking.apm.collector.analysis.jvm.define.service.IMemoryPoolMetricService;
import org.apache.skywalking.apm.collector.analysis.jvm.define.service.*;
import org.apache.skywalking.apm.collector.analysis.metric.define.AnalysisMetricModule;
import org.apache.skywalking.apm.collector.analysis.metric.define.service.IInstanceHeartBeatService;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.core.util.TimeBucketUtils;
import org.apache.skywalking.apm.collector.server.grpc.GRPCHandler;
import org.apache.skywalking.apm.network.proto.CPU;
import org.apache.skywalking.apm.network.proto.Downstream;
import org.apache.skywalking.apm.network.proto.GC;
import org.apache.skywalking.apm.network.proto.JVMMetrics;
import org.apache.skywalking.apm.network.proto.JVMMetricsServiceGrpc;
import org.apache.skywalking.apm.network.proto.Memory;
import org.apache.skywalking.apm.network.proto.MemoryPool;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.skywalking.apm.network.proto.*;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -63,7 +53,10 @@ public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsSe
@Override public void collect(JVMMetrics request, StreamObserver<Downstream> responseObserver) {
int instanceId = request.getApplicationInstanceId();
logger.debug("receive the jvm metric from application instance, id: {}", instanceId);
if (logger.isDebugEnabled()) {
logger.debug("receive the jvm metric from application instance, id: {}", instanceId);
}
request.getMetricsList().forEach(metric -> {
long minuteTimeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(metric.getTime());

View File

@ -46,7 +46,10 @@ public class NetworkAddressRegisterServiceHandler extends NetworkAddressRegister
@Override
public void batchRegister(NetworkAddresses request, StreamObserver<NetworkAddressMappings> responseObserver) {
logger.debug("register application");
if (logger.isDebugEnabled()) {
logger.debug("register application");
}
ProtocolStringList addressesList = request.getAddressesList();
NetworkAddressMappings.Builder builder = NetworkAddressMappings.newBuilder();

View File

@ -44,7 +44,10 @@ public class TraceSegmentServiceHandler extends TraceSegmentServiceGrpc.TraceSeg
@Override public StreamObserver<UpstreamSegment> collect(StreamObserver<Downstream> responseObserver) {
return new StreamObserver<UpstreamSegment>() {
@Override public void onNext(UpstreamSegment segment) {
logger.debug("receive segment");
if (logger.isDebugEnabled()) {
logger.debug("receive segment");
}
segmentParseService.parse(segment, ISegmentParseService.Source.Agent);
if (debug) {

View File

@ -18,19 +18,14 @@
package org.apache.skywalking.apm.collector.agent.jetty.provider.handler;
import com.google.gson.Gson;
import com.google.gson.JsonArray;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import com.google.gson.*;
import java.io.IOException;
import javax.servlet.http.HttpServletRequest;
import org.apache.skywalking.apm.collector.analysis.register.define.AnalysisRegisterModule;
import org.apache.skywalking.apm.collector.analysis.register.define.service.INetworkAddressIDService;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.server.jetty.ArgumentsParseException;
import org.apache.skywalking.apm.collector.server.jetty.JettyJsonHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -52,17 +47,21 @@ public class NetworkAddressRegisterServletHandler extends JettyJsonHandler {
return "/networkAddress/register";
}
@Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException {
@Override protected JsonElement doGet(HttpServletRequest req) {
throw new UnsupportedOperationException();
}
@Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException {
@Override protected JsonElement doPost(HttpServletRequest req) {
JsonArray responseArray = new JsonArray();
try {
JsonArray networkAddresses = gson.fromJson(req.getReader(), JsonArray.class);
for (int i = 0; i < networkAddresses.size(); i++) {
String networkAddress = networkAddresses.get(i).getAsString();
logger.debug("network address register, network address: {}", networkAddress);
if (logger.isDebugEnabled()) {
logger.debug("network address register, network address: {}", networkAddress);
}
int addressId = networkAddressIDService.get(networkAddress);
JsonObject mapping = new JsonObject();
mapping.addProperty(ADDRESS_ID, addressId);

View File

@ -20,18 +20,14 @@ package org.apache.skywalking.apm.collector.agent.jetty.provider.handler;
import com.google.gson.JsonElement;
import com.google.gson.stream.JsonReader;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.*;
import javax.servlet.http.HttpServletRequest;
import org.apache.skywalking.apm.collector.agent.jetty.provider.handler.reader.TraceSegment;
import org.apache.skywalking.apm.collector.agent.jetty.provider.handler.reader.TraceSegmentJsonReader;
import org.apache.skywalking.apm.collector.agent.jetty.provider.handler.reader.*;
import org.apache.skywalking.apm.collector.analysis.segment.parser.define.AnalysisSegmentParserModule;
import org.apache.skywalking.apm.collector.analysis.segment.parser.define.service.ISegmentParseService;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.server.jetty.ArgumentsParseException;
import org.apache.skywalking.apm.collector.server.jetty.JettyJsonHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -50,18 +46,22 @@ public class TraceSegmentServletHandler extends JettyJsonHandler {
return "/segments";
}
@Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException {
@Override protected JsonElement doGet(HttpServletRequest req) {
throw new UnsupportedOperationException();
}
@Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException {
logger.debug("receive stream segment");
@Override protected JsonElement doPost(HttpServletRequest req) {
if (logger.isDebugEnabled()) {
logger.debug("receive stream segment");
}
try {
BufferedReader bufferedReader = req.getReader();
read(bufferedReader);
} catch (IOException e) {
logger.error(e.getMessage(), e);
}
return null;
}

View File

@ -34,7 +34,7 @@ import static java.util.Objects.isNull;
*/
public class GCMetricService implements IGCMetricService {
private final Logger logger = LoggerFactory.getLogger(GCMetricService.class);
private static final Logger logger = LoggerFactory.getLogger(GCMetricService.class);
private Graph<GCMetric> gcMetricGraph;
@ -59,7 +59,9 @@ public class GCMetricService implements IGCMetricService {
gcMetric.setTimes(1L);
gcMetric.setTimeBucket(timeBucket);
logger.debug("push to gc metric graph, id: {}", gcMetric.getId());
if (logger.isDebugEnabled()) {
logger.debug("push to gc metric graph, id: {}", gcMetric.getId());
}
getGcMetricGraph().start(gcMetric);
}
}

View File

@ -100,14 +100,14 @@ public class AnalysisMetricModuleProvider extends ModuleProvider {
private void segmentParserListenerRegister() {
ISegmentParserListenerRegister segmentParserListenerRegister = getManager().find(AnalysisSegmentParserModule.NAME).getService(ISegmentParserListenerRegister.class);
segmentParserListenerRegister.register(new ServiceReferenceMetricSpanListener.Factory());
segmentParserListenerRegister.register(new ApplicationComponentSpanListener.Factory());
segmentParserListenerRegister.register(new ApplicationMappingSpanListener.Factory());
segmentParserListenerRegister.register(new InstanceMappingSpanListener.Factory());
segmentParserListenerRegister.register(new GlobalTraceSpanListener.Factory());
segmentParserListenerRegister.register(new SegmentDurationSpanListener.Factory());
segmentParserListenerRegister.register(new ResponseTimeDistributionSpanListener.Factory());
segmentParserListenerRegister.register(new ServiceNameSpanListener.Factory());
segmentParserListenerRegister.register(new ServiceReferenceMetricSpanListener.Factory()); //11000TPS
segmentParserListenerRegister.register(new ApplicationComponentSpanListener.Factory()); //17000TPS
segmentParserListenerRegister.register(new ApplicationMappingSpanListener.Factory()); //22000TPS
segmentParserListenerRegister.register(new InstanceMappingSpanListener.Factory()); //22000TPS
segmentParserListenerRegister.register(new GlobalTraceSpanListener.Factory()); //15000TPS
segmentParserListenerRegister.register(new SegmentDurationSpanListener.Factory()); //13000TPS
segmentParserListenerRegister.register(new ResponseTimeDistributionSpanListener.Factory()); //25000TPS
segmentParserListenerRegister.register(new ServiceNameSpanListener.Factory()); //20000TPS
}
private void graphCreate(WorkerCreateListener workerCreateListener) {

View File

@ -51,7 +51,9 @@ public class InstanceHeartBeatService implements IInstanceHeartBeatService {
instance.setHeartBeatTime(TimeBucketUtils.INSTANCE.getSecondTimeBucket(heartBeatTime));
instance.setInstanceId(instanceId);
logger.debug("push to instance heart beat persistence worker, id: {}", instance.getId());
if (logger.isDebugEnabled()) {
logger.debug("push to instance heart beat persistence worker, id: {}", instance.getId());
}
getHeartBeatGraph().start(instance);
}
}

View File

@ -36,7 +36,7 @@ import org.apache.skywalking.apm.collector.storage.table.application.Application
public class ApplicationComponentSpanListener implements EntrySpanListener, ExitSpanListener {
private final ApplicationCacheService applicationCacheService;
private List<ApplicationComponent> applicationComponents = new LinkedList<>();
private final List<ApplicationComponent> applicationComponents = new LinkedList<>();
private ApplicationComponentSpanListener(ModuleManager moduleManager) {
this.applicationCacheService = moduleManager.find(CacheModule.NAME).getService(ApplicationCacheService.class);

View File

@ -51,7 +51,10 @@ public class ApplicationMappingSpanListener implements EntrySpanListener {
}
@Override public void parseEntry(SpanDecorator spanDecorator, SegmentCoreInfo segmentCoreInfo) {
logger.debug("application mapping listener parse reference");
if (logger.isDebugEnabled()) {
logger.debug("application mapping listener parse reference");
}
if (!spanDecorator.getSpanLayer().equals(SpanLayer.MQ)) {
if (spanDecorator.getRefsCount() > 0) {
for (int i = 0; i < spanDecorator.getRefsCount(); i++) {
@ -73,10 +76,16 @@ public class ApplicationMappingSpanListener implements EntrySpanListener {
}
@Override public void build() {
logger.debug("application mapping listener build");
if (logger.isDebugEnabled()) {
logger.debug("application mapping listener build");
}
Graph<ApplicationMapping> graph = GraphManager.INSTANCE.findGraph(MetricGraphIdDefine.APPLICATION_MAPPING_GRAPH_ID, ApplicationMapping.class);
applicationMappings.forEach(applicationMapping -> {
logger.debug("push to application mapping aggregation worker, id: {}", applicationMapping.getId());
if (logger.isDebugEnabled()) {
logger.debug("push to application mapping aggregation worker, id: {}", applicationMapping.getId());
}
graph.start(applicationMapping);
});
}

View File

@ -37,7 +37,7 @@ public class GlobalTraceSpanListener implements GlobalTraceIdsListener {
private static final Logger logger = LoggerFactory.getLogger(GlobalTraceSpanListener.class);
private List<String> globalTraceIds = new LinkedList<>();
private final List<String> globalTraceIds = new LinkedList<>();
private SegmentCoreInfo segmentCoreInfo;
@Override public boolean containsPoint(Point point) {
@ -58,17 +58,19 @@ public class GlobalTraceSpanListener implements GlobalTraceIdsListener {
}
@Override public void build() {
logger.debug("global trace listener build");
if (logger.isDebugEnabled()) {
logger.debug("global trace listener build");
}
Graph<GlobalTrace> graph = GraphManager.INSTANCE.findGraph(MetricGraphIdDefine.GLOBAL_TRACE_GRAPH_ID, GlobalTrace.class);
for (String globalTraceId : globalTraceIds) {
globalTraceIds.forEach(globalTraceId -> {
GlobalTrace globalTrace = new GlobalTrace();
globalTrace.setId(segmentCoreInfo.getSegmentId() + Const.ID_SPLIT + globalTraceId);
globalTrace.setTraceId(globalTraceId);
globalTrace.setSegmentId(segmentCoreInfo.getSegmentId());
globalTrace.setTimeBucket(segmentCoreInfo.getMinuteTimeBucket());
graph.start(globalTrace);
}
});
}
public static class Factory implements SpanListenerFactory {

View File

@ -43,7 +43,10 @@ public class InstanceMappingSpanListener implements EntrySpanListener {
}
@Override public void parseEntry(SpanDecorator spanDecorator, SegmentCoreInfo segmentCoreInfo) {
logger.debug("instance mapping listener parse reference");
if (logger.isDebugEnabled()) {
logger.debug("instance mapping listener parse reference");
}
if (spanDecorator.getRefsCount() > 0) {
for (int i = 0; i < spanDecorator.getRefsCount(); i++) {
InstanceMapping instanceMapping = new InstanceMapping();
@ -60,10 +63,16 @@ public class InstanceMappingSpanListener implements EntrySpanListener {
}
@Override public void build() {
logger.debug("instance mapping listener build");
if (logger.isDebugEnabled()) {
logger.debug("instance mapping listener build");
}
Graph<InstanceMapping> graph = GraphManager.INSTANCE.findGraph(MetricGraphIdDefine.INSTANCE_MAPPING_GRAPH_ID, InstanceMapping.class);
instanceMappings.forEach(instanceMapping -> {
logger.debug("push to instance mapping aggregation worker, id: {}", instanceMapping.getId());
if (logger.isDebugEnabled()) {
logger.debug("push to instance mapping aggregation worker, id: {}", instanceMapping.getId());
}
graph.start(instanceMapping);
});
}

View File

@ -76,7 +76,11 @@ public class SegmentDurationSpanListener implements FirstSpanListener, EntrySpan
@Override public void build() {
Graph<SegmentDuration> graph = GraphManager.INSTANCE.findGraph(MetricGraphIdDefine.SEGMENT_DURATION_GRAPH_ID, SegmentDuration.class);
logger.debug("segment duration listener build");
if (logger.isDebugEnabled()) {
logger.debug("segment duration listener build");
}
if (entryOperationNameIds.size() == 0) {
segmentDuration.getServiceName().add(serviceNameCacheService.get(firstOperationNameId).getServiceName());
} else {

View File

@ -49,7 +49,7 @@ public class ServiceNameAggregationWorker extends AggregationWorker<ServiceName,
}
@Override public int queueSize() {
return 256;
return 4096;
}
}

View File

@ -61,7 +61,7 @@ public class ServiceNameHeartBeatPersistenceWorker extends MergePersistenceWorke
@Override
public int queueSize() {
return 1024;
return 4096;
}
}

View File

@ -32,7 +32,7 @@ import org.apache.skywalking.apm.collector.storage.table.register.ServiceName;
*/
public class ServiceNameSpanListener implements EntrySpanListener, ExitSpanListener, LocalSpanListener {
private List<ServiceName> serviceNames;
private final List<ServiceName> serviceNames;
private ServiceNameSpanListener() {
this.serviceNames = new LinkedList<>();

View File

@ -148,7 +148,10 @@ public class ServiceReferenceMetricSpanListener implements EntrySpanListener, Ex
}
@Override public void build() {
logger.debug("service reference listener build");
if (logger.isDebugEnabled()) {
logger.debug("service reference listener build");
}
Graph<ServiceReferenceMetric> graph = GraphManager.INSTANCE.findGraph(MetricGraphIdDefine.SERVICE_REFERENCE_METRIC_GRAPH_ID, ServiceReferenceMetric.class);
entryReferenceMetric.forEach(serviceReferenceMetric -> {
String metricId = serviceReferenceMetric.getFrontServiceId() + Const.ID_SPLIT + serviceReferenceMetric.getBehindServiceId() + Const.ID_SPLIT + serviceReferenceMetric.getSourceValue();
@ -157,7 +160,10 @@ public class ServiceReferenceMetricSpanListener implements EntrySpanListener, Ex
serviceReferenceMetric.setId(id);
serviceReferenceMetric.setMetricId(metricId);
serviceReferenceMetric.setTimeBucket(minuteTimeBucket);
logger.debug("push to service reference aggregation worker, id: {}", serviceReferenceMetric.getId());
if (logger.isDebugEnabled()) {
logger.debug("push to service reference aggregation worker, id: {}", serviceReferenceMetric.getId());
}
graph.start(serviceReferenceMetric);
});

View File

@ -19,15 +19,11 @@
package org.apache.skywalking.apm.collector.analysis.register.provider.register;
import org.apache.skywalking.apm.collector.analysis.register.define.graph.WorkerIdDefine;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractRemoteWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractRemoteWorkerProvider;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.*;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.remote.service.RemoteSenderService;
import org.apache.skywalking.apm.collector.remote.service.Selector;
import org.apache.skywalking.apm.collector.remote.service.*;
import org.apache.skywalking.apm.collector.storage.table.register.Application;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -44,8 +40,11 @@ public class ApplicationRegisterRemoteWorker extends AbstractRemoteWorker<Applic
return WorkerIdDefine.APPLICATION_REGISTER_REMOTE_WORKER;
}
@Override protected void onWork(Application message) throws WorkerException {
logger.debug("application code: {}", message.getApplicationCode());
@Override protected void onWork(Application message) {
if (logger.isDebugEnabled()) {
logger.debug("application code: {}", message.getApplicationCode());
}
onNext(message);
}

View File

@ -19,19 +19,15 @@
package org.apache.skywalking.apm.collector.analysis.register.provider.register;
import org.apache.skywalking.apm.collector.analysis.register.define.graph.WorkerIdDefine;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorkerProvider;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.*;
import org.apache.skywalking.apm.collector.cache.CacheModule;
import org.apache.skywalking.apm.collector.cache.service.ApplicationCacheService;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.core.util.BooleanUtils;
import org.apache.skywalking.apm.collector.core.util.Const;
import org.apache.skywalking.apm.collector.core.util.*;
import org.apache.skywalking.apm.collector.storage.StorageModule;
import org.apache.skywalking.apm.collector.storage.dao.register.IApplicationRegisterDAO;
import org.apache.skywalking.apm.collector.storage.table.register.Application;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -53,8 +49,10 @@ public class ApplicationRegisterSerialWorker extends AbstractLocalAsyncWorker<Ap
return WorkerIdDefine.APPLICATION_REGISTER_SERIAL_WORKER;
}
@Override protected void onWork(Application application) throws WorkerException {
logger.debug("register application, application code: {}", application.getApplicationCode());
@Override protected void onWork(Application application) {
if (logger.isDebugEnabled()) {
logger.debug("register application, application code: {}", application.getApplicationCode());
}
int applicationId;

View File

@ -19,15 +19,11 @@
package org.apache.skywalking.apm.collector.analysis.register.provider.register;
import org.apache.skywalking.apm.collector.analysis.register.define.graph.WorkerIdDefine;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractRemoteWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractRemoteWorkerProvider;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.*;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.remote.service.RemoteSenderService;
import org.apache.skywalking.apm.collector.remote.service.Selector;
import org.apache.skywalking.apm.collector.remote.service.*;
import org.apache.skywalking.apm.collector.storage.table.register.Instance;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -44,8 +40,10 @@ public class InstanceRegisterRemoteWorker extends AbstractRemoteWorker<Instance,
super(moduleManager);
}
@Override protected void onWork(Instance instance) throws WorkerException {
logger.debug("application id: {}, agentUUID: {}, register time: {}", instance.getApplicationId(), instance.getAgentUUID(), instance.getRegisterTime());
@Override protected void onWork(Instance instance) {
if (logger.isDebugEnabled()) {
logger.debug("application id: {}, agentUUID: {}, register time: {}", instance.getApplicationId(), instance.getAgentUUID(), instance.getRegisterTime());
}
onNext(instance);
}

View File

@ -19,19 +19,15 @@
package org.apache.skywalking.apm.collector.analysis.register.provider.register;
import org.apache.skywalking.apm.collector.analysis.register.define.graph.WorkerIdDefine;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorkerProvider;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.*;
import org.apache.skywalking.apm.collector.cache.CacheModule;
import org.apache.skywalking.apm.collector.cache.service.InstanceCacheService;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.core.util.BooleanUtils;
import org.apache.skywalking.apm.collector.core.util.Const;
import org.apache.skywalking.apm.collector.core.util.*;
import org.apache.skywalking.apm.collector.storage.StorageModule;
import org.apache.skywalking.apm.collector.storage.dao.register.IInstanceRegisterDAO;
import org.apache.skywalking.apm.collector.storage.table.register.Instance;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -53,8 +49,10 @@ public class InstanceRegisterSerialWorker extends AbstractLocalAsyncWorker<Insta
return WorkerIdDefine.INSTANCE_REGISTER_SERIAL_WORKER;
}
@Override protected void onWork(Instance instance) throws WorkerException {
logger.debug("register instance, application id: {}, agentUUID: {}", instance.getApplicationId(), instance.getAgentUUID());
@Override protected void onWork(Instance instance) {
if (logger.isDebugEnabled()) {
logger.debug("register instance, application id: {}, agentUUID: {}", instance.getApplicationId(), instance.getAgentUUID());
}
int instanceId;
if (BooleanUtils.valueToBoolean(instance.getIsAddress())) {

View File

@ -19,15 +19,11 @@
package org.apache.skywalking.apm.collector.analysis.register.provider.register;
import org.apache.skywalking.apm.collector.analysis.register.define.graph.WorkerIdDefine;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractRemoteWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractRemoteWorkerProvider;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.*;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.remote.service.RemoteSenderService;
import org.apache.skywalking.apm.collector.remote.service.Selector;
import org.apache.skywalking.apm.collector.remote.service.*;
import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -44,8 +40,11 @@ public class NetworkAddressRegisterRemoteWorker extends AbstractRemoteWorker<Net
return WorkerIdDefine.NETWORK_ADDRESS_REGISTER_REMOTE_WORKER;
}
@Override protected void onWork(NetworkAddress message) throws WorkerException {
logger.debug("network address: {}", message.getNetworkAddress());
@Override protected void onWork(NetworkAddress message) {
if (logger.isDebugEnabled()) {
logger.debug("network address: {}", message.getNetworkAddress());
}
onNext(message);
}

View File

@ -19,17 +19,14 @@
package org.apache.skywalking.apm.collector.analysis.register.provider.register;
import org.apache.skywalking.apm.collector.analysis.register.define.graph.WorkerIdDefine;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorkerProvider;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.*;
import org.apache.skywalking.apm.collector.cache.CacheModule;
import org.apache.skywalking.apm.collector.cache.service.NetworkAddressCacheService;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.storage.StorageModule;
import org.apache.skywalking.apm.collector.storage.dao.register.INetworkAddressRegisterDAO;
import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -51,8 +48,11 @@ public class NetworkAddressRegisterSerialWorker extends AbstractLocalAsyncWorker
return WorkerIdDefine.NETWORK_ADDRESS_REGISTER_SERIAL_WORKER;
}
@Override protected void onWork(NetworkAddress networkAddress) throws WorkerException {
logger.debug("register network address, address: {}", networkAddress.getNetworkAddress());
@Override protected void onWork(NetworkAddress networkAddress) {
if (logger.isDebugEnabled()) {
logger.debug("register network address, address: {}", networkAddress.getNetworkAddress());
}
if (networkAddress.getAddressId() == 0) {
int addressId = networkAddressCacheService.getAddressId(networkAddress.getNetworkAddress());

View File

@ -50,7 +50,10 @@ public class ServiceNameRegisterSerialWorker extends AbstractLocalAsyncWorker<Se
}
@Override protected void onWork(ServiceName serviceName) {
logger.debug("register service name: {}, application id: {}", serviceName.getServiceName(), serviceName.getApplicationId());
if (logger.isDebugEnabled()) {
logger.debug("register service name: {}, application id: {}", serviceName.getServiceName(), serviceName.getApplicationId());
}
int serviceId = serviceIdCacheService.get(serviceName.getApplicationId(), serviceName.getSrcSpanType(), serviceName.getServiceName());
if (serviceId == 0) {
long now = System.currentTimeMillis();

View File

@ -73,7 +73,10 @@ public class InstanceIDService implements IInstanceIDService {
}
@Override public int getOrCreateByAgentUUID(int applicationId, String agentUUID, long registerTime, AgentOsInfo osInfo) {
logger.debug("get or getOrCreate instance id by agent UUID, application id: {}, agentUUID: {}, registerTime: {}, osInfo: {}", applicationId, agentUUID, registerTime, osInfo);
if (logger.isDebugEnabled()) {
logger.debug("get or getOrCreate instance id by agent UUID, application id: {}, agentUUID: {}, registerTime: {}, osInfo: {}", applicationId, agentUUID, registerTime, osInfo);
}
int instanceId = getInstanceCacheService().getInstanceIdByAgentUUID(applicationId, agentUUID);
if (instanceId == 0) {
@ -95,7 +98,10 @@ public class InstanceIDService implements IInstanceIDService {
}
@Override public int getOrCreateByAddressId(int applicationId, int addressId, long registerTime) {
logger.debug("get or getOrCreate instance id by address id, application id: {}, address id: {}, registerTime: {}", applicationId, addressId, registerTime);
if (logger.isDebugEnabled()) {
logger.debug("get or getOrCreate instance id by address id, application id: {}, address id: {}, registerTime: {}", applicationId, addressId, registerTime);
}
int instanceId = getInstanceCacheService().getInstanceIdByAddressId(applicationId, addressId);
if (instanceId == 0) {

View File

@ -97,7 +97,10 @@ public enum SegmentBufferManager {
}
private void newDataFile() throws IOException {
logger.debug("getOrCreate new segment buffer file");
if (logger.isDebugEnabled()) {
logger.debug("getOrCreate new segment buffer file");
}
String timeBucket = String.valueOf(TimeBucketUtils.INSTANCE.getSecondTimeBucket(System.currentTimeMillis()));
String writeFileName = DATA_FILE_PREFIX + "_" + timeBucket + "." + Const.FILE_SUFFIX;
File dataFile = new File(BufferFileConfig.BUFFER_PATH + writeFileName);

View File

@ -109,7 +109,10 @@ public enum SegmentBufferReader {
}
for (File dataFile : dataFiles) {
logger.debug("Reading segment buffer data file, file name: {}", dataFile.getAbsolutePath());
if (logger.isDebugEnabled()) {
logger.debug("Reading segment buffer data file, file name: {}", dataFile.getAbsolutePath());
}
OffsetManager.INSTANCE.setReadOffset(dataFile.getName(), 0);
if (!read(dataFile, 0)) {
break;
@ -137,7 +140,11 @@ public enum SegmentBufferReader {
final int serialized = upstreamSegment.getSerializedSize();
readFileOffset = readFileOffset + CodedOutputStream.computeUInt32SizeNoTag(serialized) + serialized;
logger.debug("read segment buffer from file: {}, offset: {}, file length: {}", readFile.getName(), readFileOffset, readFile.length());
if (logger.isDebugEnabled()) {
logger.debug("read segment buffer from file: {}, offset: {}, file length: {}", readFile.getName(), readFileOffset, readFile.length());
}
OffsetManager.INSTANCE.setReadOffset(readFileOffset);
}

View File

@ -41,7 +41,7 @@ public class SegmentParse {
private static final Logger logger = LoggerFactory.getLogger(SegmentParse.class);
private final ModuleManager moduleManager;
private List<SpanListener> spanListeners;
private final List<SpanListener> spanListeners;
private final SegmentParserListenerManager listenerManager;
private final SegmentCoreInfo segmentCoreInfo;
@ -65,14 +65,19 @@ public class SegmentParse {
SegmentDecorator segmentDecorator = new SegmentDecorator(segmentObject);
if (!preBuild(traceIds, segmentDecorator)) {
logger.debug("This segment id exchange not success, write to buffer file, id: {}", segmentCoreInfo.getSegmentId());
if (logger.isDebugEnabled()) {
logger.debug("This segment id exchange not success, write to buffer file, id: {}", segmentCoreInfo.getSegmentId());
}
if (source.equals(ISegmentParseService.Source.Agent)) {
writeToBufferFile(segmentCoreInfo.getSegmentId(), segment);
}
return false;
} else {
logger.debug("This segment id exchange success, id: {}", segmentCoreInfo.getSegmentId());
if (logger.isDebugEnabled()) {
logger.debug("This segment id exchange success, id: {}", segmentCoreInfo.getSegmentId());
}
notifyListenerToBuild();
buildSegment(segmentCoreInfo.getSegmentId(), segmentDecorator.toByteArray());
return true;
@ -167,7 +172,10 @@ public class SegmentParse {
@GraphComputingMetric(name = "/segment/parse/bufferFile/write")
private void writeToBufferFile(String id, UpstreamSegment upstreamSegment) {
logger.debug("push to segment buffer write worker, id: {}", id);
if (logger.isDebugEnabled()) {
logger.debug("push to segment buffer write worker, id: {}", id);
}
SegmentStandardization standardization = new SegmentStandardization(id);
standardization.setUpstreamSegment(upstreamSegment);
Graph<SegmentStandardization> graph = GraphManager.INSTANCE.findGraph(GraphIdDefine.SEGMENT_STANDARDIZATION_GRAPH_ID, SegmentStandardization.class);

View File

@ -59,7 +59,7 @@ public class SegmentPersistenceWorker extends NonMergePersistenceWorker<Segment>
@Override
public int queueSize() {
return 1024;
return 2048;
}
}
}

View File

@ -62,7 +62,9 @@ public class SpanIdExchanger implements IdExchanger<SpanDecorator> {
int componentId = componentLibraryCatalogService.getComponentId(standardBuilder.getComponent());
if (componentId == 0) {
logger.debug("component: {} in application: {} exchange failed", standardBuilder.getComponent(), applicationId);
if (logger.isDebugEnabled()) {
logger.debug("component: {} in application: {} exchange failed", standardBuilder.getComponent(), applicationId);
}
return false;
} else {
standardBuilder.toBuilder();
@ -75,7 +77,9 @@ public class SpanIdExchanger implements IdExchanger<SpanDecorator> {
int peerId = networkAddressIDService.getOrCreate(standardBuilder.getPeer());
if (peerId == 0) {
logger.debug("peer: {} in application: {} exchange failed", standardBuilder.getPeer(), applicationId);
if (logger.isDebugEnabled()) {
logger.debug("peer: {} in application: {} exchange failed", standardBuilder.getPeer(), applicationId);
}
return false;
} else {
standardBuilder.toBuilder();
@ -93,7 +97,9 @@ public class SpanIdExchanger implements IdExchanger<SpanDecorator> {
int operationNameId = serviceNameService.getOrCreate(applicationId, standardBuilder.getSpanTypeValue(), operationName);
if (operationNameId == 0) {
logger.debug("service name: {} from application id: {} exchange failed", operationName, applicationId);
if (logger.isDebugEnabled()) {
logger.debug("service name: {} from application id: {} exchange failed", operationName, applicationId);
}
return false;
} else {
standardBuilder.toBuilder();

View File

@ -35,7 +35,7 @@ public class SegmentBase64Printer {
private static final Logger LOGGER = LoggerFactory.getLogger(SegmentBase64Printer.class);
public static void main(String[] args) throws InvalidProtocolBufferException {
String segmentBase64 = "CgwKCgMXjPKUga3WgBsSvAEIARiF7Jq1nywgp+yatZ8sKlASDAoKAnPAqKD5rNaAGxgBIAIqDjEyNy4wLjAuMTo5MDkyOAJCFC9zZW5kTWVzc2FnZS97Y291bnR9UhQvc2VuZE1lc3NhZ2Uve2NvdW50fTocS2Fma2EvVHJhY2UtdG9waWMtMS9Db25zdW1lclgEYBt6GwoJbXEuYnJva2VyEg4xMjcuMC4wLjE6OTA5MnoZCghtcS50b3BpYxINVHJhY2UtdG9waWMtMRImEP///////////wEY/+uatZ8sILTsmrWfLDD///////////8BUAIYAiAD";
String segmentBase64 = "CgoKCJbf2NPCLBAQEiAQ////////////ARiV39jTwiwg2+7Y08IsMNQPWANgARIlCAEYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAIQARif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIAxACGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgEEAMYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAUQBBif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIBhAFGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgHEAYYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAgQBxif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicICRAIGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgKEAkYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAsQChif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIDBALGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgNEAwYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCA4QDRif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEvMCCA8QDhif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADggHIAhKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXIS8wIIEBAPGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAOCAcgCEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchLzAggREBAYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgA4IByAISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEvMCCBIQERif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADggHIAhKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXIS8wIIExASGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAOCAcgCEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchLzAggUEBMYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgA4IByAISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyGP7//////////wEgCA==";
byte[] binarySegment = Base64.getDecoder().decode(segmentBase64);
TraceSegmentObject segmentObject = TraceSegmentObject.parseFrom(binarySegment);

View File

@ -39,7 +39,7 @@ public abstract class AbstractLocalAsyncWorkerProvider<INPUT extends QueueData,
workerCreateListener.addWorker(localAsyncWorker);
LocalAsyncWorkerRef<INPUT, OUTPUT> localAsyncWorkerRef = new LocalAsyncWorkerRef<>(localAsyncWorker);
DataCarrier<INPUT> dataCarrier = new DataCarrier<>(1, queueSize());
DataCarrier<INPUT> dataCarrier = new DataCarrier<>(1, 10000);
localAsyncWorkerRef.setQueueEventHandler(dataCarrier);
dataCarrier.consume(localAsyncWorkerRef, 1);
return localAsyncWorkerRef;

View File

@ -29,7 +29,7 @@ import org.slf4j.LoggerFactory;
*/
public abstract class AbstractWorker<INPUT, OUTPUT> implements NodeProcessor<INPUT, OUTPUT> {
private final Logger logger = LoggerFactory.getLogger(AbstractWorker.class);
private static final Logger logger = LoggerFactory.getLogger(AbstractWorker.class);
private final ModuleManager moduleManager;

View File

@ -18,23 +18,21 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.base;
import java.util.Iterator;
import java.util.List;
import java.util.*;
import org.apache.skywalking.apm.collector.core.annotations.trace.BatchParameter;
import org.apache.skywalking.apm.collector.core.data.QueueData;
import org.apache.skywalking.apm.collector.core.graph.NodeProcessor;
import org.apache.skywalking.apm.collector.core.queue.EndOfBatchContext;
import org.apache.skywalking.apm.commons.datacarrier.DataCarrier;
import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class LocalAsyncWorkerRef<INPUT extends QueueData, OUTPUT extends QueueData> extends WorkerRef<INPUT, OUTPUT> implements IConsumer<INPUT> {
private final Logger logger = LoggerFactory.getLogger(LocalAsyncWorkerRef.class);
private static final Logger logger = LoggerFactory.getLogger(LocalAsyncWorkerRef.class);
private DataCarrier<INPUT> dataCarrier;

View File

@ -28,7 +28,7 @@ import org.slf4j.LoggerFactory;
*/
public class RemoteWorkerRef<INPUT extends RemoteData, OUTPUT extends RemoteData> extends WorkerRef<INPUT, OUTPUT> {
private final Logger logger = LoggerFactory.getLogger(RemoteWorkerRef.class);
private static final Logger logger = LoggerFactory.getLogger(RemoteWorkerRef.class);
private final AbstractRemoteWorker<INPUT, OUTPUT> remoteWorker;
private final RemoteSenderService remoteSenderService;

View File

@ -18,22 +18,20 @@
package org.apache.skywalking.apm.collector.analysis.worker.model.impl;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.AbstractLocalAsyncWorker;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.WorkerException;
import org.apache.skywalking.apm.collector.analysis.worker.model.base.*;
import org.apache.skywalking.apm.collector.analysis.worker.model.impl.data.MergeDataCache;
import org.apache.skywalking.apm.collector.core.data.StreamData;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public abstract class AggregationWorker<INPUT extends StreamData, OUTPUT extends StreamData> extends AbstractLocalAsyncWorker<INPUT, OUTPUT> {
private final Logger logger = LoggerFactory.getLogger(AggregationWorker.class);
private static final Logger logger = LoggerFactory.getLogger(AggregationWorker.class);
private MergeDataCache<OUTPUT> mergeDataCache;
private final MergeDataCache<OUTPUT> mergeDataCache;
private int messageNum;
public AggregationWorker(ModuleManager moduleManager) {
@ -52,13 +50,10 @@ public abstract class AggregationWorker<INPUT extends StreamData, OUTPUT extends
messageNum++;
aggregate(output);
if (messageNum >= 100) {
if (messageNum >= 1000 || message.getEndOfBatchContext().isEndOfBatch()) {
sendToNext();
messageNum = 0;
}
if (message.getEndOfBatchContext().isEndOfBatch()) {
sendToNext();
}
}
private void sendToNext() throws WorkerException {
@ -70,8 +65,12 @@ public abstract class AggregationWorker<INPUT extends StreamData, OUTPUT extends
throw new WorkerException(e.getMessage(), e);
}
}
mergeDataCache.getLast().collection().forEach((String id, OUTPUT data) -> {
logger.debug(data.toString());
if (logger.isDebugEnabled()) {
logger.debug(data.toString());
}
onNext(data);
});
mergeDataCache.finishReadingLast();

View File

@ -32,7 +32,7 @@ import static java.util.Objects.nonNull;
*/
public abstract class MergePersistenceWorker<INPUT_AND_OUTPUT extends StreamData> extends PersistenceWorker<INPUT_AND_OUTPUT, MergeDataCollection<INPUT_AND_OUTPUT>> {
private final Logger logger = LoggerFactory.getLogger(MergePersistenceWorker.class);
private static final Logger logger = LoggerFactory.getLogger(MergePersistenceWorker.class);
private final MergeDataCache<INPUT_AND_OUTPUT> mergeDataCache;
@ -46,22 +46,21 @@ public abstract class MergePersistenceWorker<INPUT_AND_OUTPUT extends StreamData
}
@Override protected List<Object> prepareBatch(MergeDataCollection<INPUT_AND_OUTPUT> collection) {
List<Object> insertBatchCollection = new LinkedList<>();
List<Object> updateBatchCollection = new LinkedList<>();
List<Object> batchCollection = new LinkedList<>();
collection.collection().forEach((id, data) -> {
if (needMergeDBData()) {
INPUT_AND_OUTPUT dbData = persistenceDAO().get(id);
if (nonNull(dbData)) {
dbData.mergeAndFormulaCalculateData(data);
try {
updateBatchCollection.add(persistenceDAO().prepareBatchUpdate(dbData));
batchCollection.add(persistenceDAO().prepareBatchUpdate(dbData));
onNext(dbData);
} catch (Throwable t) {
logger.error(t.getMessage(), t);
}
} else {
try {
insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data));
batchCollection.add(persistenceDAO().prepareBatchInsert(data));
onNext(data);
} catch (Throwable t) {
logger.error(t.getMessage(), t);
@ -69,7 +68,7 @@ public abstract class MergePersistenceWorker<INPUT_AND_OUTPUT extends StreamData
}
} else {
try {
insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data));
batchCollection.add(persistenceDAO().prepareBatchInsert(data));
onNext(data);
} catch (Throwable t) {
logger.error(t.getMessage(), t);
@ -77,8 +76,7 @@ public abstract class MergePersistenceWorker<INPUT_AND_OUTPUT extends StreamData
}
});
insertBatchCollection.addAll(updateBatchCollection);
return insertBatchCollection;
return batchCollection;
}
@Override protected void cacheData(INPUT_AND_OUTPUT input) {

View File

@ -30,7 +30,7 @@ import org.slf4j.*;
*/
public abstract class NonMergePersistenceWorker<INPUT_AND_OUTPUT extends StreamData> extends PersistenceWorker<INPUT_AND_OUTPUT, NonMergeDataCollection<INPUT_AND_OUTPUT>> {
private final Logger logger = LoggerFactory.getLogger(NonMergePersistenceWorker.class);
private static final Logger logger = LoggerFactory.getLogger(NonMergePersistenceWorker.class);
private final NonMergeDataCache<INPUT_AND_OUTPUT> mergeDataCache;
@ -50,7 +50,7 @@ public abstract class NonMergePersistenceWorker<INPUT_AND_OUTPUT extends StreamD
}
@Override protected List<Object> prepareBatch(NonMergeDataCollection<INPUT_AND_OUTPUT> collection) {
List<Object> insertBatchCollection = new LinkedList<>();
List<Object> insertBatchCollection = new ArrayList<>(collection.collection().size());
collection.collection().forEach(data -> {
try {
insertBatchCollection.add(persistenceDAO().prepareBatchInsert(data));

View File

@ -36,7 +36,7 @@ import org.slf4j.*;
*/
public abstract class PersistenceWorker<INPUT_AND_OUTPUT extends StreamData, COLLECTION extends Collection> extends AbstractLocalAsyncWorker<INPUT_AND_OUTPUT, INPUT_AND_OUTPUT> {
private final Logger logger = LoggerFactory.getLogger(PersistenceWorker.class);
private static final Logger logger = LoggerFactory.getLogger(PersistenceWorker.class);
private final IBatchDAO batchDAO;
private final int blockBatchPersistenceSize;

View File

@ -62,29 +62,42 @@ public enum PersistenceTimer {
@SuppressWarnings("unchecked")
private void extractDataAndSave(IBatchDAO batchDAO, List<PersistenceWorker> persistenceWorkers) {
logger.debug("Extract data and save");
if (logger.isDebugEnabled()) {
logger.debug("Extract data and save");
}
long startTime = System.currentTimeMillis();
try {
List batchAllCollection = new LinkedList();
persistenceWorkers.forEach((PersistenceWorker worker) -> {
logger.debug("extract {} worker data and save", worker.getClass().getName());
if (logger.isDebugEnabled()) {
logger.debug("extract {} worker data and save", worker.getClass().getName());
}
if (worker.flushAndSwitch()) {
List<?> batchCollection = worker.buildBatchCollection();
logger.debug("extract {} worker data size: {}", worker.getClass().getName(), batchCollection.size());
if (logger.isDebugEnabled()) {
logger.debug("extract {} worker data size: {}", worker.getClass().getName(), batchCollection.size());
}
batchAllCollection.addAll(batchCollection);
}
});
if (debug) {
logger.info("build batch persistence duration: {} ms", System.currentTimeMillis() - startTime);
}
batchDAO.batchPersistence(batchAllCollection);
} catch (Throwable e) {
logger.error(e.getMessage(), e);
} finally {
logger.debug("persistence data save finish");
if (logger.isDebugEnabled()) {
logger.debug("persistence data save finish");
}
}
if (debug) {
long endTime = System.currentTimeMillis();
logger.info("batch persistence duration: {} ms", endTime - startTime);
logger.info("batch persistence duration: {} ms", System.currentTimeMillis() - startTime);
}
}
}

View File

@ -76,6 +76,11 @@ storage:
indexShardsNumber: 2
indexReplicasNumber: 0
highPerformanceMode: true
# Batch process setting, refer to https://www.elastic.co/guide/en/elasticsearch/client/java-api/5.5/java-docs-bulk-processor.html
bulkActions: 2000 # Execute the bulk every 2000 requests
bulkSize: 20 # flush the bulk every 20mb
flushInterval: 10 # flush the bulk every 10 seconds whatever the number of requests
concurrentRequests: 2 # the number of concurrent requests
# Set a timeout on metric data. After the timeout has expired, the metric data will automatically be deleted.
traceDataTTL: 90 # Unit is minute
minuteMetricDataTTL: 90 # Unit is minute
@ -106,5 +111,4 @@ configuration:
# default:
# host: localhost
# port: 9411
# contextPath: /
#
# contextPath: /

View File

@ -37,7 +37,7 @@ public class ServiceIdCacheGuavaService implements ServiceIdCacheService {
private final Logger logger = LoggerFactory.getLogger(ServiceIdCacheGuavaService.class);
private final Cache<String, Integer> serviceIdCache = CacheBuilder.newBuilder().maximumSize(10000).build();
private final Cache<String, Integer> serviceIdCache = CacheBuilder.newBuilder().maximumSize(1000000).build();
private final ModuleManager moduleManager;
private IServiceNameCacheDAO serviceNameCacheDAO;

View File

@ -18,18 +18,15 @@
package org.apache.skywalking.apm.collector.cache.guava.service;
import com.google.common.cache.Cache;
import com.google.common.cache.CacheBuilder;
import com.google.common.cache.*;
import org.apache.skywalking.apm.collector.cache.service.ServiceNameCacheService;
import org.apache.skywalking.apm.collector.core.module.ModuleManager;
import org.apache.skywalking.apm.collector.storage.StorageModule;
import org.apache.skywalking.apm.collector.storage.dao.cache.IServiceNameCacheDAO;
import org.apache.skywalking.apm.collector.storage.table.register.ServiceName;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
import static java.util.Objects.isNull;
import static java.util.Objects.nonNull;
import static java.util.Objects.*;
/**
* @author peng-yongsheng
@ -38,7 +35,7 @@ public class ServiceNameCacheGuavaService implements ServiceNameCacheService {
private final Logger logger = LoggerFactory.getLogger(ServiceNameCacheGuavaService.class);
private final Cache<Integer, ServiceName> serviceCache = CacheBuilder.newBuilder().maximumSize(10000).build();
private final Cache<Integer, ServiceName> serviceCache = CacheBuilder.newBuilder().maximumSize(1000000).build();
private final ModuleManager moduleManager;
private IServiceNameCacheDAO serviceNameCacheDAO;
@ -66,6 +63,8 @@ public class ServiceNameCacheGuavaService implements ServiceNameCacheService {
serviceName = getServiceNameCacheDAO().get(serviceId);
if (nonNull(serviceName)) {
serviceCache.put(serviceId, serviceName);
} else {
logger.warn("Service id {} is not in cache and persistent storage.", serviceId);
}
}

View File

@ -18,24 +18,18 @@
package org.apache.skywalking.apm.collector.client.elasticsearch;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.LinkedList;
import java.util.List;
import java.net.*;
import java.util.*;
import java.util.function.Consumer;
import org.apache.skywalking.apm.collector.client.Client;
import org.apache.skywalking.apm.collector.client.ClientException;
import org.apache.skywalking.apm.collector.client.NameSpace;
import org.apache.skywalking.apm.collector.client.*;
import org.apache.skywalking.apm.collector.core.data.CommonTable;
import org.apache.skywalking.apm.collector.core.util.Const;
import org.apache.skywalking.apm.collector.core.util.StringUtils;
import org.apache.skywalking.apm.collector.core.util.*;
import org.elasticsearch.action.admin.indices.create.CreateIndexResponse;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexResponse;
import org.elasticsearch.action.admin.indices.exists.indices.IndicesExistsResponse;
import org.elasticsearch.action.admin.indices.mapping.get.GetFieldMappingsResponse;
import org.elasticsearch.action.bulk.BulkRequestBuilder;
import org.elasticsearch.action.get.GetRequestBuilder;
import org.elasticsearch.action.get.MultiGetRequestBuilder;
import org.elasticsearch.action.get.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.search.SearchRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
@ -45,11 +39,9 @@ import org.elasticsearch.common.settings.Settings;
import org.elasticsearch.common.transport.InetSocketTransportAddress;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.index.reindex.DeleteByQueryAction;
import org.elasticsearch.index.reindex.DeleteByQueryRequestBuilder;
import org.elasticsearch.index.reindex.*;
import org.elasticsearch.transport.client.PreBuiltTransportClient;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author peng-yongsheng

View File

@ -67,7 +67,7 @@ public abstract class Window<WINDOW_COLLECTION extends Collection> {
}
}
protected WINDOW_COLLECTION getCurrent() {
private WINDOW_COLLECTION getCurrent() {
return pointer;
}

View File

@ -26,6 +26,16 @@ import org.apache.skywalking.apm.collector.core.queue.EndOfBatchContext;
*/
public abstract class StreamData extends AbstractData implements QueueData {
private static final ByteColumn[] BYTE_COLUMNS = {};
private static final StringListColumn[] STRING_LIST_COLUMNS = {};
private static final LongListColumn[] LONG_LIST_COLUMNS = {};
private static final IntegerListColumn[] INTEGER_LIST_COLUMNS = {};
private static final DoubleListColumn[] DOUBLE_LIST_COLUMNS = {};
private EndOfBatchContext endOfBatchContext;
@Override public final EndOfBatchContext getEndOfBatchContext() {
@ -41,18 +51,18 @@ public abstract class StreamData extends AbstractData implements QueueData {
DoubleColumn[] doubleColumns, StringListColumn[] stringListColumns,
LongListColumn[] longListColumns,
IntegerListColumn[] integerListColumns, DoubleListColumn[] doubleListColumns) {
super(stringColumns, longColumns, integerColumns, doubleColumns, new ByteColumn[0], stringListColumns, longListColumns, integerListColumns, doubleListColumns);
super(stringColumns, longColumns, integerColumns, doubleColumns, BYTE_COLUMNS, stringListColumns, longListColumns, integerListColumns, doubleListColumns);
}
public StreamData(StringColumn[] stringColumns, LongColumn[] longColumns,
IntegerColumn[] integerColumns, DoubleColumn[] doubleColumns) {
super(stringColumns, longColumns, integerColumns, doubleColumns, new ByteColumn[0], new StringListColumn[0], new LongListColumn[0], new IntegerListColumn[0], new DoubleListColumn[0]);
super(stringColumns, longColumns, integerColumns, doubleColumns, BYTE_COLUMNS, STRING_LIST_COLUMNS, LONG_LIST_COLUMNS, INTEGER_LIST_COLUMNS, DOUBLE_LIST_COLUMNS);
}
public StreamData(StringColumn[] stringColumns, LongColumn[] longColumns,
IntegerColumn[] integerColumns, DoubleColumn[] doubleColumns,
ByteColumn[] byteColumns) {
super(stringColumns, longColumns, integerColumns, doubleColumns, byteColumns, new StringListColumn[0], new LongListColumn[0], new IntegerListColumn[0], new DoubleListColumn[0]);
super(stringColumns, longColumns, integerColumns, doubleColumns, byteColumns, STRING_LIST_COLUMNS, LONG_LIST_COLUMNS, INTEGER_LIST_COLUMNS, DOUBLE_LIST_COLUMNS);
}
@Override public final String selectKey() {

View File

@ -24,6 +24,7 @@ import org.apache.skywalking.apm.collector.core.data.*;
* @author peng-yongsheng
*/
public class NonMergeOperation implements MergeOperation {
@Override public String operate(String newValue, String oldValue) {
return oldValue;
}

View File

@ -20,14 +20,10 @@ package org.apache.skywalking.apm.collector.instrument;
import java.lang.annotation.Annotation;
import java.lang.reflect.Method;
import java.util.LinkedList;
import java.util.List;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.*;
import java.util.concurrent.*;
import org.apache.skywalking.apm.collector.core.annotations.trace.BatchParameter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.*;
/**
* @author wusheng, peng-yongsheng
@ -61,8 +57,8 @@ public enum MetricTree implements Runnable {
logBuffer.append("##################################################################################################################").append(lineSeparator);
logBuffer.append("# Collector Service Report #").append(lineSeparator);
logBuffer.append("##################################################################################################################").append(lineSeparator);
metrics.forEach((MetricNode metric) -> metric.toOutput(new ReportWriter() {
metrics.forEach((MetricNode metric) -> metric.toOutput(new ReportWriter() {
@Override public void writeMetricName(String name) {
logBuffer.append(name).append(lineSeparator);
}

View File

@ -19,19 +19,15 @@
package org.apache.skywalking.apm.collector.instrument;
import java.lang.reflect.Method;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentHashMap;
import net.bytebuddy.implementation.bind.annotation.AllArguments;
import net.bytebuddy.implementation.bind.annotation.Origin;
import net.bytebuddy.implementation.bind.annotation.RuntimeType;
import net.bytebuddy.implementation.bind.annotation.SuperCall;
import net.bytebuddy.implementation.bind.annotation.This;
import java.util.concurrent.*;
import net.bytebuddy.implementation.bind.annotation.*;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
/**
* @author wu-sheng
*/
public class ServiceMetricTracing {
private volatile ConcurrentHashMap<Method, ServiceMetric> metrics = new ConcurrentHashMap<>();
ServiceMetricTracing() {

View File

@ -18,10 +18,8 @@
package org.apache.skywalking.apm.collector.instrument.tools;
import java.util.LinkedHashMap;
import java.util.Map;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.*;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -36,20 +34,13 @@ class ReportFormatter {
logger.info(System.lineSeparator() + "Formatted report: ");
report.getMetrics().forEach(metric -> {
String[] subMetricNames = metric.getMetricName().split("/");
String metricName = "";
for (String subMetricName : subMetricNames) {
if (subMetricName != null && !subMetricName.equals("")) {
metricName = metricName + "/" + subMetricName;
if (!metricMap.containsKey(metricName)) {
Metric newMetric = new Metric();
newMetric.setMetricName(metricName);
metricMap.put(metricName, newMetric);
}
metricMap.get(metricName).merge(metric);
}
if (metricMap.containsKey(metric.getMetricName())) {
Metric existMetric = metricMap.get(metric.getMetricName());
existMetric.setTotal(existMetric.getTotal() + metric.getTotal());
existMetric.setCalls(existMetric.getCalls() + metric.getCalls());
existMetric.setAvg(existMetric.getTotal() / existMetric.getCalls());
} else {
metricMap.put(metric.getMetricName(), metric);
}
});

View File

@ -32,7 +32,9 @@ public class ForeverFirstSelector implements RemoteClientSelector {
private static final Logger logger = LoggerFactory.getLogger(ForeverFirstSelector.class);
@Override public RemoteClient select(List<RemoteClient> clients, RemoteData remoteData) {
logger.debug("clients size: {}", clients.size());
if (logger.isDebugEnabled()) {
logger.debug("clients size: {}", clients.size());
}
return clients.get(0);
}
}

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.apm.collector.storage.base.dao;
import java.io.IOException;
import org.apache.skywalking.apm.collector.core.data.StreamData;
/**
@ -27,9 +28,9 @@ public interface IPersistenceDAO<INSERT, UPDATE, STREAM_DATA extends StreamData>
STREAM_DATA get(String id);
INSERT prepareBatchInsert(STREAM_DATA data);
INSERT prepareBatchInsert(STREAM_DATA data) throws IOException;
UPDATE prepareBatchUpdate(STREAM_DATA data);
UPDATE prepareBatchUpdate(STREAM_DATA data) throws IOException;
void deleteHistory(Long timeBucketBefore);
}

View File

@ -43,8 +43,11 @@ public class ApplicationAlarm extends StreamData implements Alarm {
new IntegerColumn(ApplicationAlarmTable.APPLICATION_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationAlarm() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class ApplicationAlarmList extends StreamData {
new IntegerColumn(ApplicationAlarmListTable.APPLICATION_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationAlarmList() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class ApplicationReferenceAlarm extends StreamData implements Alarm {
new IntegerColumn(ApplicationReferenceAlarmTable.BEHIND_APPLICATION_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationReferenceAlarm() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class ApplicationReferenceAlarmList extends StreamData {
new IntegerColumn(ApplicationReferenceAlarmListTable.BEHIND_APPLICATION_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationReferenceAlarmList() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class InstanceAlarm extends StreamData implements Alarm {
new IntegerColumn(InstanceAlarmTable.INSTANCE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public InstanceAlarm() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class InstanceAlarmList extends StreamData {
new IntegerColumn(InstanceAlarmListTable.INSTANCE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public InstanceAlarmList() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -46,8 +46,11 @@ public class InstanceReferenceAlarm extends StreamData implements Alarm {
new IntegerColumn(InstanceReferenceAlarmTable.BEHIND_INSTANCE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public InstanceReferenceAlarm() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -46,8 +46,11 @@ public class InstanceReferenceAlarmList extends StreamData {
new IntegerColumn(InstanceReferenceAlarmListTable.BEHIND_INSTANCE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public InstanceReferenceAlarmList() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -45,8 +45,11 @@ public class ServiceAlarm extends StreamData implements Alarm {
new IntegerColumn(ServiceAlarmTable.SERVICE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ServiceAlarm() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -45,8 +45,11 @@ public class ServiceAlarmList extends StreamData {
new IntegerColumn(ServiceAlarmListTable.SERVICE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ServiceAlarmList() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -48,8 +48,11 @@ public class ServiceReferenceAlarm extends StreamData implements Alarm {
new IntegerColumn(ServiceReferenceAlarmTable.BEHIND_SERVICE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ServiceReferenceAlarm() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -48,8 +48,11 @@ public class ServiceReferenceAlarmList extends StreamData {
new IntegerColumn(ServiceReferenceAlarmListTable.BEHIND_SERVICE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ServiceReferenceAlarmList() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -42,8 +42,11 @@ public class ApplicationComponent extends StreamData {
new IntegerColumn(ApplicationComponentTable.APPLICATION_ID, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationComponent() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -42,8 +42,11 @@ public class ApplicationMapping extends StreamData {
new IntegerColumn(ApplicationMappingTable.MAPPING_APPLICATION_ID, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationMapping() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -62,8 +62,11 @@ public class ApplicationMetric extends StreamData implements Metric {
new IntegerColumn(ApplicationMetricTable.APPLICATION_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -64,8 +64,11 @@ public class ApplicationReferenceMetric extends StreamData implements Metric {
new IntegerColumn(ApplicationReferenceMetricTable.BEHIND_APPLICATION_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ApplicationReferenceMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -37,8 +37,14 @@ public class GlobalTrace extends StreamData {
new LongColumn(GlobalTraceTable.TIME_BUCKET, new CoverMergeOperation()),
};
private static final IntegerColumn[] INTEGER_COLUMNS = {
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public GlobalTrace() {
super(STRING_COLUMNS, LONG_COLUMNS, new IntegerColumn[0], new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class ResponseTimeDistribution extends StreamData {
new IntegerColumn(ResponseTimeDistributionTable.STEP, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ResponseTimeDistribution() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -43,8 +43,11 @@ public class InstanceMapping extends StreamData {
new IntegerColumn(InstanceMappingTable.ADDRESS_ID, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public InstanceMapping() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -60,8 +60,11 @@ public class InstanceMetric extends StreamData implements Metric {
new IntegerColumn(InstanceMetricTable.INSTANCE_ID, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public InstanceMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -62,8 +62,11 @@ public class InstanceReferenceMetric extends StreamData implements Metric {
new IntegerColumn(InstanceReferenceMetricTable.BEHIND_INSTANCE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public InstanceReferenceMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class GCMetric extends StreamData {
new IntegerColumn(GCMetricTable.PHRASE, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public GCMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -46,8 +46,11 @@ public class MemoryMetric extends StreamData {
new IntegerColumn(MemoryMetricTable.IS_HEAP, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public MemoryMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -46,8 +46,11 @@ public class MemoryPoolMetric extends StreamData {
new IntegerColumn(MemoryPoolMetricTable.POOL_TYPE, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public MemoryPoolMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -39,8 +39,14 @@ public class Application extends StreamData {
new IntegerColumn(ApplicationTable.IS_ADDRESS, new CoverMergeOperation()),
};
private static final LongColumn[] LONG_COLUMNS = {
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public Application() {
super(STRING_COLUMNS, new LongColumn[0], INTEGER_COLUMNS, new DoubleColumn[0], new StringListColumn[0], new LongListColumn[0], new IntegerListColumn[0], new DoubleListColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -47,8 +47,11 @@ public class Instance extends StreamData {
new IntegerColumn(InstanceTable.IS_ADDRESS, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public Instance() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -39,9 +39,14 @@ public class NetworkAddress extends StreamData {
new IntegerColumn(NetworkAddressTable.SERVER_TYPE, new CoverMergeOperation()),
};
public NetworkAddress() {
super(STRING_COLUMNS, new LongColumn[0], INTEGER_COLUMNS, new DoubleColumn[0], new StringListColumn[0], new LongListColumn[0], new IntegerListColumn[0], new DoubleListColumn[0]);
private static final LongColumn[] LONG_COLUMNS = {
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public NetworkAddress() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -44,8 +44,11 @@ public class ServiceName extends StreamData {
new IntegerColumn(ServiceNameTable.SRC_SPAN_TYPE, new CoverMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ServiceName() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -39,8 +39,14 @@ public class Segment extends StreamData {
new ByteColumn(SegmentTable.DATA_BINARY, new CoverMergeOperation()),
};
private static final IntegerColumn[] INTEGER_COLUMNS = {
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public Segment() {
super(STRING_COLUMNS, LONG_COLUMNS, new IntegerColumn[0], new DoubleColumn[0], BYTE_COLUMNS);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS, BYTE_COLUMNS);
}
@Override public String getId() {

View File

@ -61,8 +61,11 @@ public class ServiceMetric extends StreamData implements Metric {
new IntegerColumn(ServiceMetricTable.SERVICE_ID, new NonMergeOperation()),
};
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ServiceMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0]);
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -67,9 +67,11 @@ public class ServiceReferenceMetric extends StreamData implements Metric {
new IntegerColumn(ServiceReferenceMetricTable.BEHIND_APPLICATION_ID, new NonMergeOperation()),
};
public ServiceReferenceMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, new DoubleColumn[0], new StringListColumn[0], new LongListColumn[0], new IntegerListColumn[0], new DoubleListColumn[0]);
private static final DoubleColumn[] DOUBLE_COLUMNS = {
};
public ServiceReferenceMetric() {
super(STRING_COLUMNS, LONG_COLUMNS, INTEGER_COLUMNS, DOUBLE_COLUMNS);
}
@Override public String getId() {

View File

@ -18,9 +18,10 @@
package org.apache.skywalking.apm.collector.storage.es;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.storage.table.Metric;
import org.apache.skywalking.apm.collector.storage.table.MetricColumns;
import org.apache.skywalking.apm.collector.storage.table.*;
import org.elasticsearch.common.xcontent.XContentBuilder;
/**
* @author peng-yongsheng
@ -51,26 +52,26 @@ public enum MetricTransformUtil {
target.setMqTransactionAverageDuration(((Number)source.get(MetricColumns.MQ_TRANSACTION_AVERAGE_DURATION.getName())).longValue());
}
public void esStreamDataToEsData(Metric source, Map<String, Object> target) {
target.put(MetricColumns.TIME_BUCKET.getName(), source.getTimeBucket());
target.put(MetricColumns.SOURCE_VALUE.getName(), source.getSourceValue());
public void esStreamDataToEsData(Metric source, XContentBuilder target) throws IOException {
target.field(MetricColumns.TIME_BUCKET.getName(), source.getTimeBucket());
target.field(MetricColumns.SOURCE_VALUE.getName(), source.getSourceValue());
target.put(MetricColumns.TRANSACTION_CALLS.getName(), source.getTransactionCalls());
target.put(MetricColumns.TRANSACTION_ERROR_CALLS.getName(), source.getTransactionErrorCalls());
target.put(MetricColumns.TRANSACTION_DURATION_SUM.getName(), source.getTransactionDurationSum());
target.put(MetricColumns.TRANSACTION_ERROR_DURATION_SUM.getName(), source.getTransactionErrorDurationSum());
target.put(MetricColumns.TRANSACTION_AVERAGE_DURATION.getName(), source.getTransactionAverageDuration());
target.field(MetricColumns.TRANSACTION_CALLS.getName(), source.getTransactionCalls());
target.field(MetricColumns.TRANSACTION_ERROR_CALLS.getName(), source.getTransactionErrorCalls());
target.field(MetricColumns.TRANSACTION_DURATION_SUM.getName(), source.getTransactionDurationSum());
target.field(MetricColumns.TRANSACTION_ERROR_DURATION_SUM.getName(), source.getTransactionErrorDurationSum());
target.field(MetricColumns.TRANSACTION_AVERAGE_DURATION.getName(), source.getTransactionAverageDuration());
target.put(MetricColumns.BUSINESS_TRANSACTION_CALLS.getName(), source.getBusinessTransactionCalls());
target.put(MetricColumns.BUSINESS_TRANSACTION_ERROR_CALLS.getName(), source.getBusinessTransactionErrorCalls());
target.put(MetricColumns.BUSINESS_TRANSACTION_DURATION_SUM.getName(), source.getBusinessTransactionDurationSum());
target.put(MetricColumns.BUSINESS_TRANSACTION_ERROR_DURATION_SUM.getName(), source.getBusinessTransactionErrorDurationSum());
target.put(MetricColumns.BUSINESS_TRANSACTION_AVERAGE_DURATION.getName(), source.getBusinessTransactionAverageDuration());
target.field(MetricColumns.BUSINESS_TRANSACTION_CALLS.getName(), source.getBusinessTransactionCalls());
target.field(MetricColumns.BUSINESS_TRANSACTION_ERROR_CALLS.getName(), source.getBusinessTransactionErrorCalls());
target.field(MetricColumns.BUSINESS_TRANSACTION_DURATION_SUM.getName(), source.getBusinessTransactionDurationSum());
target.field(MetricColumns.BUSINESS_TRANSACTION_ERROR_DURATION_SUM.getName(), source.getBusinessTransactionErrorDurationSum());
target.field(MetricColumns.BUSINESS_TRANSACTION_AVERAGE_DURATION.getName(), source.getBusinessTransactionAverageDuration());
target.put(MetricColumns.MQ_TRANSACTION_CALLS.getName(), source.getMqTransactionCalls());
target.put(MetricColumns.MQ_TRANSACTION_ERROR_CALLS.getName(), source.getMqTransactionErrorCalls());
target.put(MetricColumns.MQ_TRANSACTION_DURATION_SUM.getName(), source.getMqTransactionDurationSum());
target.put(MetricColumns.MQ_TRANSACTION_ERROR_DURATION_SUM.getName(), source.getMqTransactionErrorDurationSum());
target.put(MetricColumns.MQ_TRANSACTION_AVERAGE_DURATION.getName(), source.getMqTransactionAverageDuration());
target.field(MetricColumns.MQ_TRANSACTION_CALLS.getName(), source.getMqTransactionCalls());
target.field(MetricColumns.MQ_TRANSACTION_ERROR_CALLS.getName(), source.getMqTransactionErrorCalls());
target.field(MetricColumns.MQ_TRANSACTION_DURATION_SUM.getName(), source.getMqTransactionDurationSum());
target.field(MetricColumns.MQ_TRANSACTION_ERROR_DURATION_SUM.getName(), source.getMqTransactionErrorDurationSum());
target.field(MetricColumns.MQ_TRANSACTION_AVERAGE_DURATION.getName(), source.getMqTransactionAverageDuration());
}
}

View File

@ -33,6 +33,10 @@ public class StorageModuleEsConfig extends ElasticSearchClientConfig {
private int hourMetricDataTTL = 36;
private int dayMetricDataTTL = 45;
private int monthMetricDataTTL = 18;
private int bulkActions = 2000;
private int bulkSize = 20;
private int flushInterval = 10;
private int concurrentRequests = 2;
int getIndexShardsNumber() {
return indexShardsNumber;
@ -97,4 +101,36 @@ public class StorageModuleEsConfig extends ElasticSearchClientConfig {
void setMonthMetricDataTTL(int monthMetricDataTTL) {
this.monthMetricDataTTL = monthMetricDataTTL == 0 ? 18 : monthMetricDataTTL;
}
public int getBulkActions() {
return bulkActions;
}
public void setBulkActions(int bulkActions) {
this.bulkActions = bulkActions == 0 ? 2000 : bulkActions;
}
public int getBulkSize() {
return bulkSize;
}
public void setBulkSize(int bulkSize) {
this.bulkSize = bulkSize == 0 ? 20 : bulkSize;
}
public int getFlushInterval() {
return flushInterval;
}
public void setFlushInterval(int flushInterval) {
this.flushInterval = flushInterval == 0 ? 10 : flushInterval;
}
public int getConcurrentRequests() {
return concurrentRequests;
}
public void setConcurrentRequests(int concurrentRequests) {
this.concurrentRequests = concurrentRequests == 0 ? 2 : concurrentRequests;
}
}

View File

@ -48,7 +48,7 @@ import org.apache.skywalking.apm.collector.storage.dao.rtd.*;
import org.apache.skywalking.apm.collector.storage.dao.smp.*;
import org.apache.skywalking.apm.collector.storage.dao.srmp.*;
import org.apache.skywalking.apm.collector.storage.dao.ui.*;
import org.apache.skywalking.apm.collector.storage.es.base.dao.BatchEsDAO;
import org.apache.skywalking.apm.collector.storage.es.base.dao.BatchProcessEsDAO;
import org.apache.skywalking.apm.collector.storage.es.base.define.ElasticSearchStorageInstaller;
import org.apache.skywalking.apm.collector.storage.es.dao.*;
import org.apache.skywalking.apm.collector.storage.es.dao.acp.*;
@ -105,7 +105,7 @@ public class StorageModuleEsProvider extends ModuleProvider {
elasticSearchClient = new ElasticSearchClient(config.getClusterName(), config.getClusterTransportSniffer(), config.getClusterNodes(), nameSpace);
this.registerServiceImplementation(ITTLConfigService.class, new TTLConfigService(config));
this.registerServiceImplementation(IBatchDAO.class, new BatchEsDAO(elasticSearchClient));
this.registerServiceImplementation(IBatchDAO.class, new BatchProcessEsDAO(elasticSearchClient, config.getBulkActions(), config.getBulkSize(), config.getFlushInterval(), config.getConcurrentRequests()));
registerCacheDAO();
registerRegisterDAO();

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.apm.collector.storage.es.base.dao;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.data.StreamData;
@ -25,6 +26,7 @@ import org.apache.skywalking.apm.collector.storage.base.dao.IPersistenceDAO;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.index.reindex.BulkByScrollResponse;
import org.slf4j.*;
@ -56,17 +58,17 @@ public abstract class AbstractPersistenceEsDAO<STREAM_DATA extends StreamData> e
}
}
protected abstract Map<String, Object> esStreamDataToEsData(STREAM_DATA streamData);
protected abstract XContentBuilder esStreamDataToEsData(STREAM_DATA streamData) throws IOException;
@Override
public final IndexRequestBuilder prepareBatchInsert(STREAM_DATA streamData) {
Map<String, Object> source = esStreamDataToEsData(streamData);
public final IndexRequestBuilder prepareBatchInsert(STREAM_DATA streamData) throws IOException {
XContentBuilder source = esStreamDataToEsData(streamData);
return getClient().prepareIndex(tableName(), streamData.getId()).setSource(source);
}
@Override
public final UpdateRequestBuilder prepareBatchUpdate(STREAM_DATA streamData) {
Map<String, Object> source = esStreamDataToEsData(streamData);
public final UpdateRequestBuilder prepareBatchUpdate(STREAM_DATA streamData) throws IOException {
XContentBuilder source = esStreamDataToEsData(streamData);
return getClient().prepareUpdate(tableName(), streamData.getId()).setDoc(source);
}

View File

@ -1,72 +0,0 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.apm.collector.storage.es.base.dao;
import java.util.List;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.BatchParameter;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.core.util.CollectionUtils;
import org.apache.skywalking.apm.collector.storage.base.dao.IBatchDAO;
import org.elasticsearch.action.bulk.BulkItemResponse;
import org.elasticsearch.action.bulk.BulkRequestBuilder;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* @author peng-yongsheng
*/
public class BatchEsDAO extends EsDAO implements IBatchDAO {
private final Logger logger = LoggerFactory.getLogger(BatchEsDAO.class);
public BatchEsDAO(ElasticSearchClient client) {
super(client);
}
@GraphComputingMetric(name = "/persistence/batchPersistence/")
@Override public void batchPersistence(@BatchParameter List<?> batchCollection) {
if (logger.isDebugEnabled()) {
logger.debug("bulk data size: {}", batchCollection.size());
}
if (CollectionUtils.isNotEmpty(batchCollection)) {
BulkRequestBuilder bulkRequest = getClient().prepareBulk();
batchCollection.forEach(builder -> {
if (builder instanceof IndexRequestBuilder) {
bulkRequest.add((IndexRequestBuilder)builder);
}
if (builder instanceof UpdateRequestBuilder) {
bulkRequest.add((UpdateRequestBuilder)builder);
}
});
BulkResponse bulkResponse = bulkRequest.execute().actionGet();
if (bulkResponse.hasFailures()) {
logger.error(bulkResponse.buildFailureMessage());
for (BulkItemResponse itemResponse : bulkResponse.getItems()) {
logger.error("Bulk request failure, index: {}, id: {}", itemResponse.getIndex(), itemResponse.getId());
}
}
}
}
}

View File

@ -0,0 +1,120 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.apm.collector.storage.es.base.dao;
import java.lang.reflect.Field;
import java.util.List;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.UnexpectedException;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.core.util.CollectionUtils;
import org.apache.skywalking.apm.collector.storage.base.dao.IBatchDAO;
import org.elasticsearch.action.bulk.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.client.Client;
import org.elasticsearch.common.unit.*;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class BatchProcessEsDAO extends EsDAO implements IBatchDAO {
private static final Logger logger = LoggerFactory.getLogger(BatchProcessEsDAO.class);
private BulkProcessor bulkProcessor;
private final int bulkActions;
private final int bulkSize;
private final int flushInterval;
private final int concurrentRequests;
public BatchProcessEsDAO(ElasticSearchClient client, int bulkActions, int bulkSize, int flushInterval,
int concurrentRequests) {
super(client);
this.bulkActions = bulkActions;
this.bulkSize = bulkSize;
this.flushInterval = flushInterval;
this.concurrentRequests = concurrentRequests;
}
@GraphComputingMetric(name = "/persistence/batchPersistence/")
@Override public void batchPersistence(List<?> batchCollection) {
if (bulkProcessor == null) {
this.bulkProcessor = createBulkProcessor();
}
if (logger.isDebugEnabled()) {
logger.debug("bulk data size: {}", batchCollection.size());
}
if (CollectionUtils.isNotEmpty(batchCollection)) {
batchCollection.forEach(builder -> {
if (builder instanceof IndexRequestBuilder) {
this.bulkProcessor.add(((IndexRequestBuilder)builder).request());
}
if (builder instanceof UpdateRequestBuilder) {
this.bulkProcessor.add(((UpdateRequestBuilder)builder).request());
}
});
}
}
private BulkProcessor createBulkProcessor() {
ElasticSearchClient elasticSearchClient = getClient();
Client client;
try {
Field field = elasticSearchClient.getClass().getDeclaredField("client");
field.setAccessible(true);
client = (Client)field.get(elasticSearchClient);
} catch (Exception e) {
logger.error(e.getMessage(), e);
throw new UnexpectedException(e.getMessage());
}
return BulkProcessor.builder(
client,
new BulkProcessor.Listener() {
@Override
public void beforeBulk(long executionId,
BulkRequest request) {
}
@Override
public void afterBulk(long executionId,
BulkRequest request,
BulkResponse response) {
}
@Override
public void afterBulk(long executionId,
BulkRequest request,
Throwable failure) {
logger.error("{} data bulk failed, reason: {}", request.numberOfActions(), failure);
}
})
.setBulkActions(bulkActions)
.setBulkSize(new ByteSizeValue(bulkSize, ByteSizeUnit.MB))
.setFlushInterval(TimeValue.timeValueSeconds(flushInterval))
.setConcurrentRequests(concurrentRequests)
.setBackoffPolicy(BackoffPolicy.exponentialBackoff(TimeValue.timeValueMillis(100), 3))
.build();
}
}

View File

@ -18,7 +18,8 @@
package org.apache.skywalking.apm.collector.storage.es.dao;
import java.util.*;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.storage.dao.IGlobalTracePersistenceDAO;
@ -26,6 +27,7 @@ import org.apache.skywalking.apm.collector.storage.es.base.dao.AbstractPersisten
import org.apache.skywalking.apm.collector.storage.table.global.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
/**
* @author peng-yongsheng
@ -48,12 +50,13 @@ public class GlobalTraceEsPersistenceDAO extends AbstractPersistenceEsDAO<Global
return globalTrace;
}
@Override protected Map<String, Object> esStreamDataToEsData(GlobalTrace streamData) {
Map<String, Object> target = new HashMap<>();
target.put(GlobalTraceTable.SEGMENT_ID.getName(), streamData.getSegmentId());
target.put(GlobalTraceTable.TRACE_ID.getName(), streamData.getTraceId());
target.put(GlobalTraceTable.TIME_BUCKET.getName(), streamData.getTimeBucket());
return target;
@Override protected XContentBuilder esStreamDataToEsData(GlobalTrace streamData) throws IOException {
return XContentFactory.jsonBuilder()
.startObject()
.field(GlobalTraceTable.SEGMENT_ID.getName(), streamData.getSegmentId())
.field(GlobalTraceTable.TRACE_ID.getName(), streamData.getTraceId())
.field(GlobalTraceTable.TIME_BUCKET.getName(), streamData.getTimeBucket())
.endObject();
}
@Override protected String timeBucketColumnNameForDelete() {

View File

@ -18,7 +18,8 @@
package org.apache.skywalking.apm.collector.storage.es.dao;
import java.util.*;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.UnexpectedException;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
@ -28,6 +29,7 @@ import org.apache.skywalking.apm.collector.storage.table.register.*;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
import org.slf4j.*;
/**
@ -63,9 +65,11 @@ public class InstanceHeartBeatEsPersistenceDAO extends EsDAO implements IInstanc
throw new UnexpectedException("Received an instance heart beat message under instance id= " + data.getId() + " , which doesn't exist.");
}
@Override public UpdateRequestBuilder prepareBatchUpdate(Instance data) {
Map<String, Object> source = new HashMap<>();
source.put(InstanceTable.HEARTBEAT_TIME.getName(), data.getHeartBeatTime());
@Override public UpdateRequestBuilder prepareBatchUpdate(Instance data) throws IOException {
XContentBuilder source = XContentFactory.jsonBuilder().startObject()
.field(InstanceTable.HEARTBEAT_TIME.getName(), data.getHeartBeatTime())
.endObject();
return getClient().prepareUpdate(InstanceTable.TABLE, data.getId()).setDoc(source);
}

View File

@ -18,14 +18,14 @@
package org.apache.skywalking.apm.collector.storage.es.dao;
import com.google.gson.Gson;
import java.util.*;
import java.io.IOException;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.storage.dao.ISegmentDurationPersistenceDAO;
import org.apache.skywalking.apm.collector.storage.es.base.dao.EsDAO;
import org.apache.skywalking.apm.collector.storage.table.segment.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.index.reindex.BulkByScrollResponse;
import org.slf4j.*;
@ -35,9 +35,7 @@ import org.slf4j.*;
*/
public class SegmentDurationEsPersistenceDAO extends EsDAO implements ISegmentDurationPersistenceDAO<IndexRequestBuilder, UpdateRequestBuilder, SegmentDuration> {
private final Logger logger = LoggerFactory.getLogger(SegmentDurationEsPersistenceDAO.class);
private final Gson gson = new Gson();
private static final Logger logger = LoggerFactory.getLogger(SegmentDurationEsPersistenceDAO.class);
public SegmentDurationEsPersistenceDAO(ElasticSearchClient client) {
super(client);
@ -54,16 +52,18 @@ public class SegmentDurationEsPersistenceDAO extends EsDAO implements ISegmentDu
}
@Override
public IndexRequestBuilder prepareBatchInsert(SegmentDuration data) {
Map<String, Object> target = new HashMap<>();
target.put(SegmentDurationTable.SEGMENT_ID.getName(), data.getSegmentId());
target.put(SegmentDurationTable.APPLICATION_ID.getName(), data.getApplicationId());
target.put(SegmentDurationTable.SERVICE_NAME.getName(), gson.toJson(data.getServiceName()));
target.put(SegmentDurationTable.DURATION.getName(), data.getDuration());
target.put(SegmentDurationTable.START_TIME.getName(), data.getStartTime());
target.put(SegmentDurationTable.END_TIME.getName(), data.getEndTime());
target.put(SegmentDurationTable.IS_ERROR.getName(), data.getIsError());
target.put(SegmentDurationTable.TIME_BUCKET.getName(), data.getTimeBucket());
public IndexRequestBuilder prepareBatchInsert(SegmentDuration data) throws IOException {
XContentBuilder target = XContentFactory.jsonBuilder().startObject()
.field(SegmentDurationTable.SEGMENT_ID.getName(), data.getSegmentId())
.field(SegmentDurationTable.APPLICATION_ID.getName(), data.getApplicationId())
.array(SegmentDurationTable.SERVICE_NAME.getName(), data.getServiceName())
.field(SegmentDurationTable.DURATION.getName(), data.getDuration())
.field(SegmentDurationTable.START_TIME.getName(), data.getStartTime())
.field(SegmentDurationTable.END_TIME.getName(), data.getEndTime())
.field(SegmentDurationTable.IS_ERROR.getName(), data.getIsError())
.field(SegmentDurationTable.TIME_BUCKET.getName(), data.getTimeBucket())
.endObject();
return getClient().prepareIndex(SegmentDurationTable.TABLE, data.getId()).setSource(target);
}

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.apm.collector.storage.es.dao;
import java.io.IOException;
import java.util.*;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
@ -26,6 +27,7 @@ import org.apache.skywalking.apm.collector.storage.es.base.dao.AbstractPersisten
import org.apache.skywalking.apm.collector.storage.table.segment.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
/**
* @author peng-yongsheng
@ -47,11 +49,11 @@ public class SegmentEsPersistenceDAO extends AbstractPersistenceEsDAO<Segment> i
return segment;
}
@Override protected Map<String, Object> esStreamDataToEsData(Segment streamData) {
Map<String, Object> target = new HashMap<>();
target.put(SegmentTable.DATA_BINARY.getName(), new String(Base64.getEncoder().encode(streamData.getDataBinary())));
target.put(SegmentTable.TIME_BUCKET.getName(), streamData.getTimeBucket());
return target;
@Override protected XContentBuilder esStreamDataToEsData(Segment streamData) throws IOException {
return XContentFactory.jsonBuilder().startObject()
.field(SegmentTable.DATA_BINARY.getName(), new String(Base64.getEncoder().encode(streamData.getDataBinary())))
.field(SegmentTable.TIME_BUCKET.getName(), streamData.getTimeBucket())
.endObject();
}
@Override protected String timeBucketColumnNameForDelete() {

View File

@ -18,7 +18,8 @@
package org.apache.skywalking.apm.collector.storage.es.dao;
import java.util.*;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.UnexpectedException;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
@ -28,6 +29,7 @@ import org.apache.skywalking.apm.collector.storage.table.register.*;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
import org.slf4j.*;
/**
@ -43,18 +45,24 @@ public class ServiceNameHeartBeatEsPersistenceDAO extends EsDAO implements IServ
@GraphComputingMetric(name = "/persistence/get/" + ServiceNameTable.TABLE + "/heartbeat")
@Override public ServiceName get(String id) {
GetResponse getResponse = getClient().prepareGet(ServiceNameTable.TABLE, id).get();
String[] includeSources = {ServiceNameTable.HEARTBEAT_TIME.getName()};
GetResponse getResponse = getClient().prepareGet(ServiceNameTable.TABLE, id).setFetchSource(includeSources, null).get();
if (getResponse.isExists()) {
Map<String, Object> source = getResponse.getSource();
ServiceName serviceName = new ServiceName();
serviceName.setId(id);
serviceName.setServiceId(((Number)source.get(ServiceNameTable.SERVICE_ID.getName())).intValue());
serviceName.setServiceId(Integer.valueOf(id));
serviceName.setHeartBeatTime(((Number)source.get(ServiceNameTable.HEARTBEAT_TIME.getName())).longValue());
logger.debug("service id: {} is exists", id);
if (logger.isDebugEnabled()) {
logger.debug("service id: {} is exists", id);
}
return serviceName;
} else {
logger.debug("service id: {} is not exists", id);
if (logger.isDebugEnabled()) {
logger.debug("service id: {} is not exists", id);
}
return null;
}
}
@ -63,11 +71,15 @@ public class ServiceNameHeartBeatEsPersistenceDAO extends EsDAO implements IServ
throw new UnexpectedException("Received an service name heart beat message under service id= " + data.getId() + " , which doesn't exist.");
}
@Override public UpdateRequestBuilder prepareBatchUpdate(ServiceName data) {
logger.info("service name heart beat, service id: {}, heart beat time: {}", data.getId(), data.getHeartBeatTime());
@Override public UpdateRequestBuilder prepareBatchUpdate(ServiceName data) throws IOException {
if (logger.isDebugEnabled()) {
logger.debug("service name heart beat, service id: {}, heart beat time: {}", data.getId(), data.getHeartBeatTime());
}
XContentBuilder source = XContentFactory.jsonBuilder().startObject()
.field(ServiceNameTable.HEARTBEAT_TIME.getName(), data.getHeartBeatTime())
.endObject();
Map<String, Object> source = new HashMap<>();
source.put(ServiceNameTable.HEARTBEAT_TIME.getName(), data.getHeartBeatTime());
return getClient().prepareUpdate(ServiceNameTable.TABLE, data.getId()).setDoc(source);
}

View File

@ -18,13 +18,13 @@
package org.apache.skywalking.apm.collector.storage.es.dao.acp;
import java.util.HashMap;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.storage.es.base.dao.AbstractPersistenceEsDAO;
import org.apache.skywalking.apm.collector.storage.table.application.ApplicationComponent;
import org.apache.skywalking.apm.collector.storage.table.application.ApplicationComponentTable;
import org.apache.skywalking.apm.collector.storage.table.application.*;
import org.elasticsearch.common.xcontent.*;
/**
* @author peng-yongsheng
@ -49,15 +49,14 @@ public abstract class AbstractApplicationComponentEsPersistenceDAO extends Abstr
return applicationComponent;
}
@Override protected final Map<String, Object> esStreamDataToEsData(ApplicationComponent streamData) {
Map<String, Object> target = new HashMap<>();
target.put(ApplicationComponentTable.METRIC_ID.getName(), streamData.getMetricId());
@Override protected final XContentBuilder esStreamDataToEsData(ApplicationComponent streamData) throws IOException {
return XContentFactory.jsonBuilder().startObject()
.field(ApplicationComponentTable.METRIC_ID.getName(), streamData.getMetricId())
target.put(ApplicationComponentTable.COMPONENT_ID.getName(), streamData.getComponentId());
target.put(ApplicationComponentTable.APPLICATION_ID.getName(), streamData.getApplicationId());
target.put(ApplicationComponentTable.TIME_BUCKET.getName(), streamData.getTimeBucket());
return target;
.field(ApplicationComponentTable.COMPONENT_ID.getName(), streamData.getComponentId())
.field(ApplicationComponentTable.APPLICATION_ID.getName(), streamData.getApplicationId())
.field(ApplicationComponentTable.TIME_BUCKET.getName(), streamData.getTimeBucket())
.endObject();
}
@GraphComputingMetric(name = "/persistence/get/" + ApplicationComponentTable.TABLE)

View File

@ -18,13 +18,13 @@
package org.apache.skywalking.apm.collector.storage.es.dao.alarm;
import java.util.HashMap;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.storage.es.base.dao.AbstractPersistenceEsDAO;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationAlarmList;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationAlarmListTable;
import org.apache.skywalking.apm.collector.storage.table.alarm.*;
import org.elasticsearch.common.xcontent.*;
/**
* @author peng-yongsheng
@ -52,17 +52,17 @@ public abstract class AbstractApplicationAlarmListEsPersistenceDAO extends Abstr
return applicationAlarmList;
}
@Override protected final Map<String, Object> esStreamDataToEsData(ApplicationAlarmList streamData) {
Map<String, Object> target = new HashMap<>();
target.put(ApplicationAlarmListTable.METRIC_ID.getName(), streamData.getMetricId());
target.put(ApplicationAlarmListTable.APPLICATION_ID.getName(), streamData.getApplicationId());
target.put(ApplicationAlarmListTable.SOURCE_VALUE.getName(), streamData.getSourceValue());
@Override protected final XContentBuilder esStreamDataToEsData(ApplicationAlarmList streamData) throws IOException {
return XContentFactory.jsonBuilder().startObject()
.field(ApplicationAlarmListTable.METRIC_ID.getName(), streamData.getMetricId())
.field(ApplicationAlarmListTable.APPLICATION_ID.getName(), streamData.getApplicationId())
.field(ApplicationAlarmListTable.SOURCE_VALUE.getName(), streamData.getSourceValue())
target.put(ApplicationAlarmListTable.ALARM_TYPE.getName(), streamData.getAlarmType());
target.put(ApplicationAlarmListTable.ALARM_CONTENT.getName(), streamData.getAlarmContent());
.field(ApplicationAlarmListTable.ALARM_TYPE.getName(), streamData.getAlarmType())
.field(ApplicationAlarmListTable.ALARM_CONTENT.getName(), streamData.getAlarmContent())
target.put(ApplicationAlarmListTable.TIME_BUCKET.getName(), streamData.getTimeBucket());
return target;
.field(ApplicationAlarmListTable.TIME_BUCKET.getName(), streamData.getTimeBucket())
.endObject();
}
@GraphComputingMetric(name = "/persistence/get/" + ApplicationAlarmListTable.TABLE)

View File

@ -18,16 +18,16 @@
package org.apache.skywalking.apm.collector.storage.es.dao.alarm;
import java.util.HashMap;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.storage.dao.alarm.IApplicationAlarmPersistenceDAO;
import org.apache.skywalking.apm.collector.storage.es.base.dao.AbstractPersistenceEsDAO;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationAlarm;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationAlarmTable;
import org.apache.skywalking.apm.collector.storage.table.alarm.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
/**
* @author peng-yongsheng
@ -54,16 +54,16 @@ public class ApplicationAlarmEsPersistenceDAO extends AbstractPersistenceEsDAO<A
return instanceAlarm;
}
@Override protected Map<String, Object> esStreamDataToEsData(ApplicationAlarm streamData) {
Map<String, Object> target = new HashMap<>();
target.put(ApplicationAlarmTable.APPLICATION_ID.getName(), streamData.getApplicationId());
target.put(ApplicationAlarmTable.SOURCE_VALUE.getName(), streamData.getSourceValue());
@Override protected XContentBuilder esStreamDataToEsData(ApplicationAlarm streamData) throws IOException {
return XContentFactory.jsonBuilder().startObject()
.field(ApplicationAlarmTable.APPLICATION_ID.getName(), streamData.getApplicationId())
.field(ApplicationAlarmTable.SOURCE_VALUE.getName(), streamData.getSourceValue())
target.put(ApplicationAlarmTable.ALARM_TYPE.getName(), streamData.getAlarmType());
target.put(ApplicationAlarmTable.ALARM_CONTENT.getName(), streamData.getAlarmContent());
.field(ApplicationAlarmTable.ALARM_TYPE.getName(), streamData.getAlarmType())
.field(ApplicationAlarmTable.ALARM_CONTENT.getName(), streamData.getAlarmContent())
target.put(ApplicationAlarmTable.LAST_TIME_BUCKET.getName(), streamData.getLastTimeBucket());
return target;
.field(ApplicationAlarmTable.LAST_TIME_BUCKET.getName(), streamData.getLastTimeBucket())
.endObject();
}
@Override protected String timeBucketColumnNameForDelete() {

View File

@ -18,16 +18,16 @@
package org.apache.skywalking.apm.collector.storage.es.dao.alarm;
import java.util.HashMap;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.storage.dao.alarm.IApplicationReferenceAlarmPersistenceDAO;
import org.apache.skywalking.apm.collector.storage.es.base.dao.AbstractPersistenceEsDAO;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationReferenceAlarm;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationReferenceAlarmTable;
import org.apache.skywalking.apm.collector.storage.table.alarm.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
/**
* @author peng-yongsheng
@ -55,17 +55,17 @@ public class ApplicationReferenceAlarmEsPersistenceDAO extends AbstractPersisten
return applicationReferenceAlarm;
}
@Override protected Map<String, Object> esStreamDataToEsData(ApplicationReferenceAlarm streamData) {
Map<String, Object> target = new HashMap<>();
target.put(ApplicationReferenceAlarmTable.FRONT_APPLICATION_ID.getName(), streamData.getFrontApplicationId());
target.put(ApplicationReferenceAlarmTable.BEHIND_APPLICATION_ID.getName(), streamData.getBehindApplicationId());
target.put(ApplicationReferenceAlarmTable.SOURCE_VALUE.getName(), streamData.getSourceValue());
@Override protected XContentBuilder esStreamDataToEsData(ApplicationReferenceAlarm streamData) throws IOException {
return XContentFactory.jsonBuilder().startObject()
.field(ApplicationReferenceAlarmTable.FRONT_APPLICATION_ID.getName(), streamData.getFrontApplicationId())
.field(ApplicationReferenceAlarmTable.BEHIND_APPLICATION_ID.getName(), streamData.getBehindApplicationId())
.field(ApplicationReferenceAlarmTable.SOURCE_VALUE.getName(), streamData.getSourceValue())
target.put(ApplicationReferenceAlarmTable.ALARM_TYPE.getName(), streamData.getAlarmType());
target.put(ApplicationReferenceAlarmTable.ALARM_CONTENT.getName(), streamData.getAlarmContent());
.field(ApplicationReferenceAlarmTable.ALARM_TYPE.getName(), streamData.getAlarmType())
.field(ApplicationReferenceAlarmTable.ALARM_CONTENT.getName(), streamData.getAlarmContent())
target.put(ApplicationReferenceAlarmTable.LAST_TIME_BUCKET.getName(), streamData.getLastTimeBucket());
return target;
.field(ApplicationReferenceAlarmTable.LAST_TIME_BUCKET.getName(), streamData.getLastTimeBucket())
.endObject();
}
@Override protected String timeBucketColumnNameForDelete() {

View File

@ -18,16 +18,16 @@
package org.apache.skywalking.apm.collector.storage.es.dao.alarm;
import java.util.HashMap;
import java.io.IOException;
import java.util.Map;
import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.apm.collector.core.annotations.trace.GraphComputingMetric;
import org.apache.skywalking.apm.collector.storage.dao.alarm.IApplicationReferenceAlarmListPersistenceDAO;
import org.apache.skywalking.apm.collector.storage.es.base.dao.AbstractPersistenceEsDAO;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationReferenceAlarmList;
import org.apache.skywalking.apm.collector.storage.table.alarm.ApplicationReferenceAlarmListTable;
import org.apache.skywalking.apm.collector.storage.table.alarm.*;
import org.elasticsearch.action.index.IndexRequestBuilder;
import org.elasticsearch.action.update.UpdateRequestBuilder;
import org.elasticsearch.common.xcontent.*;
/**
* @author peng-yongsheng
@ -55,17 +55,18 @@ public class ApplicationReferenceAlarmListEsPersistenceDAO extends AbstractPersi
return applicationReferenceAlarmList;
}
@Override protected Map<String, Object> esStreamDataToEsData(ApplicationReferenceAlarmList streamData) {
Map<String, Object> target = new HashMap<>();
target.put(ApplicationReferenceAlarmListTable.FRONT_APPLICATION_ID.getName(), streamData.getFrontApplicationId());
target.put(ApplicationReferenceAlarmListTable.BEHIND_APPLICATION_ID.getName(), streamData.getBehindApplicationId());
target.put(ApplicationReferenceAlarmListTable.SOURCE_VALUE.getName(), streamData.getSourceValue());
@Override
protected XContentBuilder esStreamDataToEsData(ApplicationReferenceAlarmList streamData) throws IOException {
return XContentFactory.jsonBuilder().startObject()
.field(ApplicationReferenceAlarmListTable.FRONT_APPLICATION_ID.getName(), streamData.getFrontApplicationId())
.field(ApplicationReferenceAlarmListTable.BEHIND_APPLICATION_ID.getName(), streamData.getBehindApplicationId())
.field(ApplicationReferenceAlarmListTable.SOURCE_VALUE.getName(), streamData.getSourceValue())
target.put(ApplicationReferenceAlarmListTable.ALARM_TYPE.getName(), streamData.getAlarmType());
target.put(ApplicationReferenceAlarmListTable.ALARM_CONTENT.getName(), streamData.getAlarmContent());
.field(ApplicationReferenceAlarmListTable.ALARM_TYPE.getName(), streamData.getAlarmType())
.field(ApplicationReferenceAlarmListTable.ALARM_CONTENT.getName(), streamData.getAlarmContent())
target.put(ApplicationReferenceAlarmListTable.TIME_BUCKET.getName(), streamData.getTimeBucket());
return target;
.field(ApplicationReferenceAlarmListTable.TIME_BUCKET.getName(), streamData.getTimeBucket())
.endObject();
}
@Override protected String timeBucketColumnNameForDelete() {

Some files were not shown because too many files have changed in this diff Show More