Skip to content

Commit 77b35db

Browse files
authored
Telemetry for OTLP traces/metrics/logs (#12057)
Telemetry for OTLP traces/metrics/logs: * Datadog tracer health metrics * otel.metrics_export_{attempts,successes,failures} * otel.log_records OTLP telemetry metrics are tagged by protocol and encoding (per signal) Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Remove HealthMetrics for OTLP traces (since HealthMetrics relies on StatsD client) and use telemetry instead Simplify OTLP telemetry Stage telemetry metrics separate to collection Co-authored-by: stuart.mcculloch <stuart.mcculloch@datadoghq.com>
1 parent 409bc13 commit 77b35db

23 files changed

Lines changed: 530 additions & 75 deletions

File tree

dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
package datadog.trace.common.writer;
22

3+
import datadog.trace.api.telemetry.OtlpTelemetry;
34
import datadog.trace.core.CoreSpan;
45
import datadog.trace.core.otlp.common.OtlpPayload;
56
import datadog.trace.core.otlp.common.OtlpSender;
@@ -37,7 +38,9 @@ public void flush() {
3738
try {
3839
OtlpPayload payload = collector.collectTraces();
3940
if (payload != OtlpPayload.EMPTY) {
40-
sender.send(payload);
41+
OtlpTelemetry.getInstance().onTracesExportAttempt();
42+
RemoteApi.Response response = sender.send(payload);
43+
OtlpTelemetry.getInstance().onTracesExportComplete(response.success());
4144
}
4245
} catch (RuntimeException e) { // don't catch severe Errors
4346
log.debug("Failed to send OTLP payload", e);
@@ -46,7 +49,7 @@ public void flush() {
4649

4750
@Override
4851
public void onDroppedTrace(int spanCount) {
49-
// TODO: surface drop counts via HealthMetrics
52+
// no telemetry currently tracked for dropped traces
5053
}
5154

5255
@Override

dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpWriter.java

Lines changed: 3 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -38,11 +38,10 @@ public static OtlpWriterBuilder builder() {
3838
TraceProcessingWorker worker,
3939
PayloadDispatcher dispatcher,
4040
OtlpSender sender,
41-
HealthMetrics healthMetrics,
4241
int flushTimeout,
4342
TimeUnit flushTimeoutUnit,
4443
boolean alwaysFlush) {
45-
super(worker, dispatcher, healthMetrics, flushTimeout, flushTimeoutUnit, alwaysFlush);
44+
super(worker, dispatcher, HealthMetrics.NO_OP, flushTimeout, flushTimeoutUnit, alwaysFlush);
4645
this.sender = sender;
4746
}
4847

@@ -64,7 +63,6 @@ public static class OtlpWriterBuilder {
6463
private OtlpConfig.Protocol protocol = OtlpConfig.Protocol.HTTP_PROTOBUF;
6564
private OtlpConfig.Compression compression = OtlpConfig.Compression.NONE;
6665
private int traceBufferSize = BUFFER_SIZE;
67-
private HealthMetrics healthMetrics = HealthMetrics.NO_OP;
6866
private int flushIntervalMilliseconds = 1000;
6967
private int flushTimeout = 1;
7068
private TimeUnit flushTimeoutUnit = TimeUnit.SECONDS;
@@ -102,11 +100,6 @@ public OtlpWriterBuilder traceBufferSize(int traceBufferSize) {
102100
return this;
103101
}
104102

105-
public OtlpWriterBuilder healthMetrics(HealthMetrics healthMetrics) {
106-
this.healthMetrics = healthMetrics;
107-
return this;
108-
}
109-
110103
public OtlpWriterBuilder flushIntervalMilliseconds(int flushIntervalMilliseconds) {
111104
this.flushIntervalMilliseconds = flushIntervalMilliseconds;
112105
return this;
@@ -151,7 +144,7 @@ public OtlpWriter build() {
151144
final TraceProcessingWorker worker =
152145
new TraceProcessingWorker(
153146
traceBufferSize,
154-
healthMetrics,
147+
HealthMetrics.NO_OP,
155148
dispatcher,
156149
DroppingPolicy.DISABLED,
157150
Prioritization.FAST_LANE,
@@ -160,7 +153,7 @@ public OtlpWriter build() {
160153
singleSpanSampler);
161154

162155
return new OtlpWriter(
163-
worker, dispatcher, sender, healthMetrics, flushTimeout, flushTimeoutUnit, alwaysFlush);
156+
worker, dispatcher, sender, flushTimeout, flushTimeoutUnit, alwaysFlush);
164157
}
165158
}
166159
}

dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -76,7 +76,6 @@ public static Writer createWriter(
7676
.protocol(config.getOtlpTracesProtocol())
7777
.compression(config.getOtlpTracesCompression())
7878
.timeoutMillis(config.getOtlpTracesTimeout())
79-
.healthMetrics(healthMetrics)
8079
.spanSamplingRules(singleSpanSampler)
8180
.flushIntervalMilliseconds(flushIntervalMilliseconds)
8281
.build();

dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpGrpcSender.java

Lines changed: 3 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -2,18 +2,16 @@
22

33
import static datadog.communication.http.OkHttpUtils.buildHttp2Client;
44
import static datadog.communication.http.OkHttpUtils.isPlainHttp;
5-
import static datadog.communication.http.OkHttpUtils.sendWithRetries;
65

76
import datadog.communication.http.HttpRetryPolicy;
87
import datadog.logging.RatelimitedLogger;
98
import datadog.trace.api.config.OtlpConfig.Compression;
10-
import java.io.IOException;
9+
import datadog.trace.common.writer.RemoteApi;
1110
import java.util.Map;
1211
import java.util.concurrent.TimeUnit;
1312
import okhttp3.HttpUrl;
1413
import okhttp3.OkHttpClient;
1514
import okhttp3.Request;
16-
import okhttp3.Response;
1715
import org.slf4j.Logger;
1816
import org.slf4j.LoggerFactory;
1917

@@ -55,20 +53,8 @@ public OtlpGrpcSender(
5553
}
5654

5755
@Override
58-
public void send(OtlpPayload payload) {
59-
Request request = makeRequest(payload);
60-
try (Response response = sendWithRetries(client, retryPolicy, request)) {
61-
if (!response.isSuccessful()) {
62-
RATELIMITED_LOGGER.warn(
63-
"OTLP export to {} failed with status {}: {}",
64-
request.url(),
65-
response.code(),
66-
response.message());
67-
}
68-
} catch (IOException e) {
69-
RATELIMITED_LOGGER.warn(
70-
"OTLP export to {} failed with exception: {}", request.url(), e.toString());
71-
}
56+
public RemoteApi.Response send(OtlpPayload payload) {
57+
return OtlpSenderSupport.send(client, retryPolicy, makeRequest(payload), RATELIMITED_LOGGER);
7258
}
7359

7460
@Override

dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpHttpSender.java

Lines changed: 3 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -2,18 +2,16 @@
22

33
import static datadog.communication.http.OkHttpUtils.buildHttpClient;
44
import static datadog.communication.http.OkHttpUtils.isPlainHttp;
5-
import static datadog.communication.http.OkHttpUtils.sendWithRetries;
65

76
import datadog.communication.http.HttpRetryPolicy;
87
import datadog.logging.RatelimitedLogger;
98
import datadog.trace.api.config.OtlpConfig.Compression;
10-
import java.io.IOException;
9+
import datadog.trace.common.writer.RemoteApi;
1110
import java.util.Map;
1211
import java.util.concurrent.TimeUnit;
1312
import okhttp3.HttpUrl;
1413
import okhttp3.OkHttpClient;
1514
import okhttp3.Request;
16-
import okhttp3.Response;
1715
import org.slf4j.Logger;
1816
import org.slf4j.LoggerFactory;
1917

@@ -59,20 +57,8 @@ public HttpUrl url() {
5957
}
6058

6159
@Override
62-
public void send(OtlpPayload payload) {
63-
Request request = makeRequest(payload);
64-
try (Response response = sendWithRetries(client, retryPolicy, request)) {
65-
if (!response.isSuccessful()) {
66-
RATELIMITED_LOGGER.warn(
67-
"OTLP export to {} failed with status {}: {}",
68-
request.url(),
69-
response.code(),
70-
response.message());
71-
}
72-
} catch (IOException e) {
73-
RATELIMITED_LOGGER.warn(
74-
"OTLP export to {} failed with exception: {}", request.url(), e.toString());
75-
}
60+
public RemoteApi.Response send(OtlpPayload payload) {
61+
return OtlpSenderSupport.send(client, retryPolicy, makeRequest(payload), RATELIMITED_LOGGER);
7662
}
7763

7864
@Override
Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,10 @@
11
package datadog.trace.core.otlp.common;
22

3+
import datadog.trace.common.writer.RemoteApi;
4+
35
/** Sends chunks of OTLP data. */
46
public interface OtlpSender {
5-
void send(OtlpPayload payload);
7+
RemoteApi.Response send(OtlpPayload payload);
68

79
void shutdown();
810
}
Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
package datadog.trace.core.otlp.common;
2+
3+
import static datadog.communication.http.OkHttpUtils.sendWithRetries;
4+
import static datadog.trace.common.writer.RemoteApi.Response.failed;
5+
import static datadog.trace.common.writer.RemoteApi.Response.success;
6+
7+
import datadog.communication.http.HttpRetryPolicy;
8+
import datadog.logging.RatelimitedLogger;
9+
import datadog.trace.common.writer.RemoteApi;
10+
import java.io.IOException;
11+
import okhttp3.OkHttpClient;
12+
13+
/** Shared request execution and response handling for {@link OtlpSender} implementations. */
14+
final class OtlpSenderSupport {
15+
private OtlpSenderSupport() {}
16+
17+
/** Executes the given request with retries, logging failures via the rate-limited logger. */
18+
static RemoteApi.Response send(
19+
OkHttpClient client,
20+
HttpRetryPolicy.Factory retryPolicy,
21+
okhttp3.Request request,
22+
RatelimitedLogger ratelimitedLogger) {
23+
try (okhttp3.Response response = sendWithRetries(client, retryPolicy, request)) {
24+
if (response.isSuccessful()) {
25+
return success(response.code());
26+
}
27+
ratelimitedLogger.warn(
28+
"OTLP export to {} failed with status {}: {}",
29+
request.url(),
30+
response.code(),
31+
response.message());
32+
return failed(response.code());
33+
} catch (IOException e) {
34+
ratelimitedLogger.warn(
35+
"OTLP export to {} failed with exception: {}", request.url(), e.toString());
36+
return failed(e);
37+
}
38+
}
39+
}

dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsCollector.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,4 +7,7 @@ public abstract class OtlpLogsCollector {
77

88
/** Waits for logs to be batched within the given interval. */
99
public abstract OtlpPayload waitForLogs(int intervalMillis);
10+
11+
/** Number of log records collected. */
12+
public abstract int getLogRecordCount();
1013
}

dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsJsonCollector.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ public final class OtlpLogsJsonCollector extends OtlpLogsCollector
3636
private JsonWriter writer;
3737
private boolean anyLogRecordWritten;
3838
private boolean logRecordStarted;
39+
private int logRecordCount;
3940

4041
private final LazyJsonArray attributesArray = new LazyJsonArray();
4142

@@ -65,6 +66,8 @@ OtlpPayload collectLogs(ObjIntConsumer<OtlpLogsVisitor> processor, int intervalM
6566

6667
/** Prepare temporary elements to collect logs data. */
6768
private void start() {
69+
logRecordCount = 0;
70+
6871
writer = new JsonWriter();
6972
writer.beginObject();
7073
writer.name("resourceLogs").beginArray();
@@ -86,6 +89,11 @@ private void stop() {
8689
currentScope = null;
8790
}
8891

92+
@Override
93+
public int getLogRecordCount() {
94+
return logRecordCount;
95+
}
96+
8997
@Override
9098
public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) {
9199
if (currentScope != null) {
@@ -119,6 +127,7 @@ public void visitLogRecord(OtlpLogRecord logRecord) {
119127

120128
logRecordStarted = false;
121129
anyLogRecordWritten = true;
130+
logRecordCount++;
122131
}
123132

124133
// opens the log record object on first attribute or value written for it

dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsProtoCollector.java

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ public final class OtlpLogsProtoCollector extends OtlpLogsCollector
4141
// total number of chunked bytes at different nesting levels
4242
private int payloadBytes;
4343
private int scopedBytes;
44+
private int logRecordCount;
4445

4546
private OtelInstrumentationScope currentScope;
4647

@@ -68,6 +69,7 @@ OtlpPayload collectLogs(ObjIntConsumer<OtlpLogsVisitor> processor, int intervalM
6869

6970
/** Prepare temporary elements to collect logs data. */
7071
private void start() {
72+
logRecordCount = 0;
7173

7274
// remove stale entries from caches
7375
OtlpCommonProto.recalibrateCaches();
@@ -84,6 +86,11 @@ private void stop() {
8486
currentScope = null;
8587
}
8688

89+
@Override
90+
public int getLogRecordCount() {
91+
return logRecordCount;
92+
}
93+
8794
@Override
8895
public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) {
8996
if (currentScope != null) {
@@ -96,6 +103,7 @@ public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) {
96103
@Override
97104
public void visitLogRecord(OtlpLogRecord logRecord) {
98105
scopedBytes += recordLogRecordMessage(buf, logRecord, protobuf);
106+
logRecordCount++;
99107
}
100108

101109
@Override

0 commit comments

Comments
 (0)