55
66package io .opentelemetry .instrumentation .api .incubator .semconv .messaging ;
77
8+ import static io .opentelemetry .instrumentation .api .internal .SemconvStability .emitOldMessagingSemconv ;
9+ import static io .opentelemetry .instrumentation .api .internal .SemconvStability .emitStableMessagingSemconv ;
10+ import static io .opentelemetry .semconv .ErrorAttributes .ERROR_TYPE ;
811import static java .util .concurrent .TimeUnit .SECONDS ;
912import static java .util .logging .Level .FINE ;
1013
2326import io .opentelemetry .instrumentation .api .instrumenter .OperationMetrics ;
2427import io .opentelemetry .instrumentation .api .internal .OperationMetricsUtil ;
2528import java .util .logging .Logger ;
29+ import javax .annotation .Nullable ;
2630
2731/**
2832 * {@link OperationListener} which keeps track of <a
29- * href="https://github.com/open-telemetry/semantic-conventions/blob/v1.26 .0/docs/messaging/messaging-metrics.md#consumer-metrics">Consumer
33+ * href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43 .0/docs/messaging/messaging-metrics.md#consumer-metrics">consumer
3034 * metrics</a>.
3135 */
3236public final class MessagingConsumerMetrics implements OperationListener {
@@ -35,34 +39,37 @@ public final class MessagingConsumerMetrics implements OperationListener {
3539 // copied from MessagingIncubatingAttributes
3640 private static final AttributeKey <Long > MESSAGING_BATCH_MESSAGE_COUNT =
3741 AttributeKey .longKey ("messaging.batch.message_count" );
42+ private static final AttributeKey <String > MESSAGING_OPERATION_TYPE =
43+ AttributeKey .stringKey ("messaging.operation.type" );
3844 private static final ContextKey <MessagingConsumerMetrics .State > MESSAGING_CONSUMER_METRICS_STATE =
3945 ContextKey .named ("messaging-consumer-metrics-state" );
4046 private static final Logger logger = Logger .getLogger (MessagingConsumerMetrics .class .getName ());
4147
42- private final DoubleHistogram receiveDurationHistogram ;
43- private final LongCounter receiveMessageCount ;
44-
45- private MessagingConsumerMetrics (Meter meter ) {
46- DoubleHistogramBuilder durationBuilder =
47- meter
48- .histogramBuilder ("messaging.receive.duration" )
49- .setDescription ("Measures the duration of receive operation." )
50- .setExplicitBucketBoundariesAdvice (MessagingMetricsAdvice .DURATION_SECONDS_BUCKETS )
51- .setUnit ("s" );
52- MessagingMetricsAdvice .applyReceiveDurationAdvice (durationBuilder );
53- receiveDurationHistogram = durationBuilder .build ();
54-
55- LongCounterBuilder longCounterBuilder =
56- meter
57- .counterBuilder ("messaging.receive.messages" )
58- .setDescription ("Measures the number of received messages." )
59- .setUnit ("{message}" );
60- MessagingMetricsAdvice .applyReceiveMessagesAdvice (longCounterBuilder );
61- receiveMessageCount = longCounterBuilder .build ();
48+ @ Nullable private final DoubleHistogram receiveDurationHistogram ;
49+ @ Nullable private final LongCounter receiveMessageCount ;
50+ @ Nullable private final DoubleHistogram clientOperationDurationHistogram ;
51+ @ Nullable private final LongCounter consumedMessagesCounter ;
52+
53+ private MessagingConsumerMetrics (Meter meter , boolean supportsStableSemconv ) {
54+ boolean emitOldSemconv = !supportsStableSemconv || emitOldMessagingSemconv ();
55+ boolean emitStableSemconv = supportsStableSemconv && emitStableMessagingSemconv ();
56+ receiveDurationHistogram = emitOldSemconv ? buildReceiveDuration (meter ) : null ;
57+ receiveMessageCount = emitOldSemconv ? buildReceiveMessages (meter ) : null ;
58+ clientOperationDurationHistogram =
59+ emitStableSemconv ? buildClientOperationDuration (meter ) : null ;
60+ consumedMessagesCounter = emitStableSemconv ? buildConsumedMessages (meter ) : null ;
6261 }
6362
63+ /** Returns metrics for extractors configured with {@link MessageOperation}. */
6464 public static OperationMetrics get () {
65- return OperationMetricsUtil .create ("messaging consumer" , MessagingConsumerMetrics ::new );
65+ return OperationMetricsUtil .create (
66+ "messaging consumer" , meter -> new MessagingConsumerMetrics (meter , false ));
67+ }
68+
69+ /** Returns metrics for extractors configured with {@link MessagingOperationType}. */
70+ public static OperationMetrics getForOperationType () {
71+ return OperationMetricsUtil .create (
72+ "messaging consumer" , meter -> new MessagingConsumerMetrics (meter , true ));
6673 }
6774
6875 @ Override
@@ -85,21 +92,97 @@ public void onEnd(Context context, Attributes endAttributes, long endNanos) {
8592 }
8693
8794 Attributes attributes = state .startAttributes ().toBuilder ().putAll (endAttributes ).build ();
88- receiveDurationHistogram .record (
89- (endNanos - state .startTimeNanos ()) / NANOS_PER_S , attributes , context );
95+ double duration = (endNanos - state .startTimeNanos ()) / NANOS_PER_S ;
96+ if (receiveDurationHistogram != null ) {
97+ receiveDurationHistogram .record (duration , attributes , context );
98+ }
99+ // Metric view attribute advice can only select keys statically. The concrete destination name
100+ // must be omitted when a template is available or the destination is temporary or anonymous,
101+ // so this conditional requirement must be enforced before recording.
102+ Attributes filteredAttributes =
103+ clientOperationDurationHistogram != null || consumedMessagesCounter != null
104+ ? MessagingMetricsAdvice .filterAttributes (attributes )
105+ : attributes ;
106+ String operationType = attributes .get (MESSAGING_OPERATION_TYPE );
107+ if (clientOperationDurationHistogram != null
108+ && !MessagingOperationType .PROCESS .value ().equals (operationType )) {
109+ clientOperationDurationHistogram .record (duration , filteredAttributes , context );
110+ }
90111
91- long receiveMessagesCount = getReceiveMessagesCount (state .startAttributes (), endAttributes );
92- receiveMessageCount .add (receiveMessagesCount , attributes , context );
112+ Long batchMessageCount = getBatchMessageCount (state .startAttributes (), endAttributes );
113+ if (receiveMessageCount != null ) {
114+ receiveMessageCount .add (
115+ batchMessageCount == null ? 1 : batchMessageCount , attributes , context );
116+ }
117+ if (consumedMessagesCounter != null
118+ && MessagingOperationType .RECEIVE .value ().equals (operationType )) {
119+ long consumedMessagesCount = getConsumedMessagesCount (attributes , batchMessageCount );
120+ if (consumedMessagesCount > 0 ) {
121+ consumedMessagesCounter .add (consumedMessagesCount , filteredAttributes , context );
122+ }
123+ }
93124 }
94125
95- private static long getReceiveMessagesCount (Attributes ... attributesList ) {
126+ @ Nullable
127+ private static Long getBatchMessageCount (Attributes ... attributesList ) {
96128 for (Attributes attributes : attributesList ) {
97129 Long value = attributes .get (MESSAGING_BATCH_MESSAGE_COUNT );
98130 if (value != null ) {
99131 return value ;
100132 }
101133 }
102- return 1 ;
134+ return null ;
135+ }
136+
137+ private static long getConsumedMessagesCount (
138+ Attributes attributes , @ Nullable Long batchMessageCount ) {
139+ if (batchMessageCount != null ) {
140+ return batchMessageCount ;
141+ }
142+ return attributes .get (ERROR_TYPE ) == null ? 1 : 0 ;
143+ }
144+
145+ private static DoubleHistogram buildReceiveDuration (Meter meter ) {
146+ DoubleHistogramBuilder builder =
147+ meter
148+ .histogramBuilder ("messaging.receive.duration" )
149+ .setDescription ("Measures the duration of receive operation." )
150+ .setExplicitBucketBoundariesAdvice (MessagingMetricsAdvice .DURATION_SECONDS_BUCKETS )
151+ .setUnit ("s" );
152+ MessagingMetricsAdvice .applyOldDurationAdvice (builder );
153+ return builder .build ();
154+ }
155+
156+ private static LongCounter buildReceiveMessages (Meter meter ) {
157+ LongCounterBuilder builder =
158+ meter
159+ .counterBuilder ("messaging.receive.messages" )
160+ .setDescription ("Measures the number of received messages." )
161+ .setUnit ("{message}" );
162+ MessagingMetricsAdvice .applyOldMessagesAdvice (builder );
163+ return builder .build ();
164+ }
165+
166+ private static DoubleHistogram buildClientOperationDuration (Meter meter ) {
167+ DoubleHistogramBuilder builder =
168+ meter
169+ .histogramBuilder ("messaging.client.operation.duration" )
170+ .setDescription (
171+ "Duration of messaging operation initiated by a producer or consumer client." )
172+ .setExplicitBucketBoundariesAdvice (MessagingMetricsAdvice .DURATION_SECONDS_BUCKETS )
173+ .setUnit ("s" );
174+ MessagingMetricsAdvice .applyClientOperationDurationAdvice (builder );
175+ return builder .build ();
176+ }
177+
178+ private static LongCounter buildConsumedMessages (Meter meter ) {
179+ LongCounterBuilder builder =
180+ meter
181+ .counterBuilder ("messaging.client.consumed.messages" )
182+ .setDescription ("Number of messages that were delivered to the application." )
183+ .setUnit ("{message}" );
184+ MessagingMetricsAdvice .applyConsumedMessagesAdvice (builder );
185+ return builder .build ();
103186 }
104187
105188 @ AutoValue
0 commit comments