Introduct DataCarrier to avoid service blocked (#5660)

This commit is contained in:
Daming 2020-10-14 08:12:22 +08:00 committed by GitHub
parent 2b9de2007f
commit 7881c57023
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
1 changed files with 50 additions and 9 deletions

View File

@ -18,6 +18,8 @@
package org.apache.skywalking.apm.agent.core.kafka;
import java.util.List;
import java.util.Objects;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.utils.Bytes;
@ -30,18 +32,26 @@ import org.apache.skywalking.apm.agent.core.context.trace.TraceSegment;
import org.apache.skywalking.apm.agent.core.logging.api.ILog;
import org.apache.skywalking.apm.agent.core.logging.api.LogManager;
import org.apache.skywalking.apm.agent.core.remote.TraceSegmentServiceClient;
import org.apache.skywalking.apm.commons.datacarrier.DataCarrier;
import org.apache.skywalking.apm.commons.datacarrier.buffer.BufferStrategy;
import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer;
import org.apache.skywalking.apm.network.language.agent.v3.SegmentObject;
import static org.apache.skywalking.apm.agent.core.conf.Config.Buffer.BUFFER_SIZE;
import static org.apache.skywalking.apm.agent.core.conf.Config.Buffer.CHANNEL_SIZE;
/**
* A tracing segment data reporter.
* A tracing segment data reporter.
*/
@OverrideImplementor(TraceSegmentServiceClient.class)
public class KafkaTraceSegmentServiceClient implements BootService, TracingContextListener {
public class KafkaTraceSegmentServiceClient implements BootService, IConsumer<TraceSegment>, TracingContextListener {
private static final ILog LOGGER = LogManager.getLogger(KafkaTraceSegmentServiceClient.class);
private String topic;
private KafkaProducer<String, Bytes> producer;
private volatile DataCarrier<TraceSegment> carrier;
@Override
public void prepare() {
topic = KafkaReporterPluginConfig.Plugin.Kafka.TOPIC_SEGMENT;
@ -49,6 +59,10 @@ public class KafkaTraceSegmentServiceClient implements BootService, TracingConte
@Override
public void boot() {
carrier = new DataCarrier<>(CHANNEL_SIZE, BUFFER_SIZE);
carrier.setBufferStrategy(BufferStrategy.IF_POSSIBLE);
carrier.consume(this, 1);
producer = ServiceManager.INSTANCE.findService(KafkaProducerManager.class).getProducer();
}
@ -60,6 +74,39 @@ public class KafkaTraceSegmentServiceClient implements BootService, TracingConte
@Override
public void shutdown() {
TracingContext.ListenerManager.remove(this);
carrier.shutdownConsumers();
}
@Override
public void init() {
}
@Override
public void consume(final List<TraceSegment> data) {
data.forEach(traceSegment -> {
SegmentObject upstreamSegment = traceSegment.transform();
ProducerRecord<String, Bytes> record = new ProducerRecord<>(
topic,
upstreamSegment.getTraceSegmentId(),
Bytes.wrap(upstreamSegment.toByteArray())
);
producer.send(record, (m, e) -> {
if (Objects.nonNull(e)) {
LOGGER.error("Failed to report TraceSegment.", e);
}
});
});
}
@Override
public void onError(final List<TraceSegment> data, final Throwable t) {
LOGGER.error(t, "Try to send {} trace segments to collector, with unexpected exception.", data.size());
}
@Override
public void onExit() {
}
@Override
@ -72,13 +119,7 @@ public class KafkaTraceSegmentServiceClient implements BootService, TracingConte
LOGGER.debug("Trace[TraceId={}] is ignored.", traceSegment.getTraceSegmentId());
return;
}
SegmentObject upstreamSegment = traceSegment.transform();
ProducerRecord<String, Bytes> record = new ProducerRecord<>(
topic,
upstreamSegment.getTraceSegmentId(),
Bytes.wrap(upstreamSegment.toByteArray())
);
producer.send(record);
carrier.produce(traceSegment);
}
}