Skip to content

Commit 371ae83

Browse files
committed
Add context extraction and enrichment for DynamoDB, Kinesis, and S3 events in OpenTelemetry tracing
1 parent 288878e commit 371ae83

10 files changed

Lines changed: 331 additions & 140 deletions

File tree

powertools-tracing-opentelemetry/src/main/java/software/amazon/lambda/powertools/tracing/opentelemetry/TracingOpenTelemetry.java

Lines changed: 6 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -40,14 +40,8 @@ public final class TracingOpenTelemetry {
4040
private final LambdaEventContextExtractorResolver eventContextExtractorResolver;
4141

4242
private TracingOpenTelemetry(Builder builder) {
43-
this.tracer = Objects.requireNonNull(
44-
builder.tracer,
45-
"tracer must not be null"
46-
);
47-
this.propagator = Objects.requireNonNull(
48-
builder.propagator,
49-
"propagator must not be null"
50-
);
43+
this.tracer = Objects.requireNonNull(builder.tracer, "tracer must not be null");
44+
this.propagator = Objects.requireNonNull(builder.propagator, "propagator must not be null");
5145
this.eventContextExtractorResolver = Objects.requireNonNull(
5246
builder.eventContextExtractorResolver,
5347
"eventContextExtractorResolver must not be null"
@@ -67,16 +61,11 @@ public TracingOpenTelemetry(Tracer tracer) {
6761
public TracingOpenTelemetry(
6862
Tracer tracer,
6963
TextMapPropagator propagator,
70-
LambdaEventContextExtractorResolver eventContextExtractorResolver) {
64+
LambdaEventContextExtractorResolver eventContextExtractorResolver
65+
) {
7166

72-
this.tracer = Objects.requireNonNull(
73-
tracer,
74-
"tracer must not be null"
75-
);
76-
this.propagator = Objects.requireNonNull(
77-
propagator,
78-
"propagator must not be null"
79-
);
67+
this.tracer = Objects.requireNonNull(tracer, "tracer must not be null");
68+
this.propagator = Objects.requireNonNull(propagator, "propagator must not be null");
8069
this.eventContextExtractorResolver = Objects.requireNonNull(
8170
eventContextExtractorResolver,
8271
"eventContextExtractorResolver must not be null"
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
package software.amazon.lambda.powertools.tracing.opentelemetry.context;
2+
3+
import com.amazonaws.services.lambda.runtime.events.DynamodbEvent;
4+
import io.opentelemetry.api.trace.Span;
5+
import io.opentelemetry.context.Context;
6+
import io.opentelemetry.context.propagation.TextMapPropagator;
7+
import java.util.Objects;
8+
9+
public final class DynamoDbTraceContextExtractor implements LambdaEventContextExtractor {
10+
11+
@Override
12+
public boolean supports(Object event) {
13+
return event instanceof DynamodbEvent;
14+
}
15+
16+
@Override
17+
public Context extract(Object event, Context parentContext, TextMapPropagator propagator) {
18+
19+
/*
20+
* DynamoDB Streams records do not expose message attributes
21+
* that can be used for W3C trace context propagation.
22+
*
23+
* Do not assume that traceparent is stored inside the DynamoDB
24+
* record payload. Propagation through DynamoDB Streams should be
25+
* defined by a dedicated propagation strategy if supported in
26+
* the future.
27+
*/
28+
return parentContext;
29+
}
30+
31+
@Override
32+
public void enrichSpan(Object event, Span span) {
33+
34+
DynamodbEvent dynamoDBEvent = (DynamodbEvent) event;
35+
36+
if (dynamoDBEvent.getRecords() == null || dynamoDBEvent.getRecords().isEmpty()) {
37+
return;
38+
}
39+
40+
DynamodbEvent.DynamodbStreamRecord record = dynamoDBEvent.getRecords()
41+
.stream()
42+
.filter(Objects::nonNull)
43+
.findFirst()
44+
.orElse(null);
45+
46+
if (record == null) {
47+
return;
48+
}
49+
50+
span.setAttribute("messaging.system", "aws.dynamodb");
51+
52+
span.setAttribute("messaging.batch.message_count", dynamoDBEvent.getRecords().size());
53+
if (record.getEventSourceARN() != null) {
54+
span.setAttribute("messaging.destination.name", extractStreamName(record.getEventSourceARN()));
55+
}
56+
}
57+
58+
private String extractStreamName(String streamArn) {
59+
int separator = streamArn.lastIndexOf('/');
60+
61+
return separator >= 0
62+
? streamArn.substring(separator + 1)
63+
: streamArn;
64+
}
65+
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
package software.amazon.lambda.powertools.tracing.opentelemetry.context;
2+
3+
import com.amazonaws.services.lambda.runtime.events.KinesisEvent;
4+
import io.opentelemetry.api.trace.Span;
5+
import io.opentelemetry.context.Context;
6+
import io.opentelemetry.context.propagation.TextMapPropagator;
7+
8+
public final class KinesisTraceContextExtractor
9+
implements LambdaEventContextExtractor {
10+
11+
@Override
12+
public boolean supports(Object event) {
13+
return event instanceof KinesisEvent;
14+
}
15+
16+
@Override
17+
public Context extract(Object event, Context parentContext, TextMapPropagator propagator) {
18+
19+
/*
20+
* Kinesis records do not expose message attributes
21+
* that can be used for W3C trace context propagation.
22+
*
23+
* Do not assume that traceparent is stored inside the Kinesis
24+
* record payload. Propagation through Kinesis should be
25+
* defined by a dedicated propagation strategy if supported in
26+
* the future.
27+
*/
28+
29+
return parentContext;
30+
}
31+
32+
@Override
33+
public void enrichSpan(Object event, Span span) {
34+
35+
KinesisEvent kinesisEvent = (KinesisEvent) event;
36+
37+
if (kinesisEvent.getRecords() == null || kinesisEvent.getRecords().isEmpty()) {
38+
return;
39+
}
40+
41+
KinesisEvent.KinesisEventRecord firstRecord = kinesisEvent.getRecords().get(0);
42+
43+
if (firstRecord == null || firstRecord.getKinesis() == null) {
44+
return;
45+
}
46+
47+
KinesisEvent.Record kinesis = firstRecord.getKinesis();
48+
49+
span.setAttribute("messaging.system", "aws.kinesis");
50+
51+
if (kinesis.getPartitionKey() != null) {
52+
span.setAttribute("messaging.partition_key", kinesis.getPartitionKey());
53+
}
54+
55+
if (kinesis.getSequenceNumber() != null) {
56+
span.setAttribute("messaging.message.id", kinesis.getSequenceNumber());
57+
}
58+
59+
if (kinesis.getApproximateArrivalTimestamp() != null) {
60+
span.setAttribute("messaging.message.receive.timestamp",
61+
kinesis.getApproximateArrivalTimestamp().getTime());
62+
}
63+
64+
if (firstRecord.getEventSourceARN() != null) {
65+
span.setAttribute("messaging.destination.name", extractStreamName(firstRecord.getEventSourceARN()));
66+
}
67+
}
68+
69+
70+
private String extractStreamName(String arn) {
71+
int separator = arn.lastIndexOf('/');
72+
73+
return separator >= 0
74+
? arn.substring(separator + 1)
75+
: arn;
76+
}
77+
}

powertools-tracing-opentelemetry/src/main/java/software/amazon/lambda/powertools/tracing/opentelemetry/context/LambdaEventContextExtractorResolver.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,10 @@ public static LambdaEventContextExtractorResolver create() {
1919
List.of(
2020
new ApiGatewayTraceContextExtractor(),
2121
new SqsTraceContextExtractor(),
22-
new SnsTraceContextExtractor()
22+
new SnsTraceContextExtractor(),
23+
new KinesisTraceContextExtractor(),
24+
new DynamoDbTraceContextExtractor(),
25+
new S3TraceContextExtractor()
2326
)
2427
);
2528
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
package software.amazon.lambda.powertools.tracing.opentelemetry.context;
2+
3+
import com.amazonaws.services.lambda.runtime.events.S3Event;
4+
import com.amazonaws.services.lambda.runtime.events.models.s3.S3EventNotification;
5+
import io.opentelemetry.api.trace.Span;
6+
import io.opentelemetry.context.Context;
7+
import io.opentelemetry.context.propagation.TextMapPropagator;
8+
import java.util.Objects;
9+
10+
public final class S3TraceContextExtractor implements LambdaEventContextExtractor {
11+
12+
@Override
13+
public boolean supports(Object event) {
14+
return event instanceof S3Event;
15+
}
16+
17+
@Override
18+
public Context extract(Object event, Context parentContext, TextMapPropagator propagator) {
19+
20+
/*
21+
* S3 event notifications do not expose message attributes
22+
* equivalent to SQS/SNS that can be passed directly to a
23+
* TextMapPropagator.
24+
*
25+
* Do not assume that traceparent/tracestate are embedded
26+
* inside the S3 event payload.
27+
*/
28+
return parentContext;
29+
}
30+
31+
@Override
32+
public void enrichSpan(Object event, Span span) {
33+
34+
S3Event s3Event = (S3Event) event;
35+
36+
if (s3Event.getRecords() == null || s3Event.getRecords().isEmpty()) {
37+
return;
38+
}
39+
40+
span.setAttribute("messaging.system", "aws.s3");
41+
42+
span.setAttribute("messaging.batch.message_count", s3Event.getRecords().size());
43+
44+
S3EventNotification.S3EventNotificationRecord record =
45+
s3Event.getRecords()
46+
.stream()
47+
.filter(Objects::nonNull)
48+
.findFirst()
49+
.orElse(null);
50+
51+
if (record == null || record.getS3() == null) {
52+
return;
53+
}
54+
55+
if (record.getS3().getBucket() != null
56+
&& record.getS3().getBucket().getName() != null) {
57+
58+
span.setAttribute("messaging.destination.name", record.getS3().getBucket().getName());
59+
}
60+
61+
if (record.getEventName() != null) {
62+
span.setAttribute("messaging.event.type", record.getEventName());
63+
}
64+
}
65+
}

powertools-tracing-opentelemetry/src/main/java/software/amazon/lambda/powertools/tracing/opentelemetry/context/SnsTraceContextExtractor.java

Lines changed: 41 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -24,28 +24,43 @@ public Context extract(Object event, Context parentContext, TextMapPropagator pr
2424
return parentContext;
2525
}
2626

27-
SNSEvent.SNSRecord record = snsEvent.getRecords().get(0);
28-
29-
Map<String, SNSEvent.MessageAttribute> attributes = record.getSNS().getMessageAttributes();
30-
31-
if (attributes == null || attributes.isEmpty()) {
32-
return parentContext;
27+
for (SNSEvent.SNSRecord record : snsEvent.getRecords()) {
28+
29+
if (record == null || record.getSNS() == null) {
30+
continue;
31+
}
32+
33+
Map<String, SNSEvent.MessageAttribute> attributes = record.getSNS().getMessageAttributes();
34+
35+
if (attributes == null || attributes.isEmpty()) {
36+
continue;
37+
}
38+
39+
Map<String, String> propagationAttributes = attributes.entrySet()
40+
.stream()
41+
.filter(entry -> entry.getValue() != null)
42+
.filter(entry -> entry.getValue().getValue() != null)
43+
.collect(Collectors.toMap(
44+
Map.Entry::getKey,
45+
entry -> entry.getValue().getValue()
46+
));
47+
48+
if (propagationAttributes.isEmpty()) {
49+
continue;
50+
}
51+
52+
Context extractedContext = propagator.extract(
53+
parentContext,
54+
propagationAttributes,
55+
OpenTelemetryProvider.textMapGetter()
56+
);
57+
58+
if (extractedContext != parentContext) {
59+
return extractedContext;
60+
}
3361
}
3462

35-
Map<String, String> propagationAttributes = attributes.entrySet()
36-
.stream()
37-
.filter(entry -> entry.getValue() != null)
38-
.filter(entry -> entry.getValue().getValue() != null)
39-
.collect(Collectors.toMap(
40-
Map.Entry::getKey,
41-
entry -> entry.getValue().getValue()
42-
));
43-
44-
return propagator.extract(
45-
parentContext,
46-
propagationAttributes,
47-
OpenTelemetryProvider.textMapGetter()
48-
);
63+
return parentContext;
4964
}
5065

5166
@Override
@@ -57,18 +72,18 @@ public void enrichSpan(Object event, Span span) {
5772
return;
5873
}
5974

60-
SNSEvent.SNSRecord record = snsEvent.getRecords().get(0);
75+
SNSEvent.SNSRecord record = snsEvent.getRecords()
76+
.stream()
77+
.filter(r -> r != null && r.getSNS() != null)
78+
.findFirst()
79+
.orElse(null);
6180

62-
if (record.getSNS() == null) {
81+
if (record == null) {
6382
return;
6483
}
6584

6685
span.setAttribute("messaging.system", "aws.sns");
6786

68-
if (record.getSNS().getMessageId() != null) {
69-
span.setAttribute("messaging.message.id", record.getSNS().getMessageId());
70-
}
71-
7287
if (record.getSNS().getTopicArn() != null) {
7388
span.setAttribute("messaging.destination.name", extractTopicName(record.getSNS().getTopicArn()));
7489
}

0 commit comments

Comments
 (0)