Skip to content

Commit

Permalink
Update attribute key of rocketmq's message tag (open-telemetry#6677)
Browse files Browse the repository at this point in the history
  • Loading branch information
aaron-ai authored and LironKS committed Dec 4, 2022
1 parent ee6ab10 commit f708d03
Show file tree
Hide file tree
Showing 3 changed files with 12 additions and 14 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import io.opentelemetry.api.common.AttributesBuilder;
import io.opentelemetry.context.Context;
import io.opentelemetry.instrumentation.api.instrumenter.AttributesExtractor;
import io.opentelemetry.semconv.trace.attributes.SemanticAttributes;
import java.net.SocketAddress;
import javax.annotation.Nullable;
import org.apache.rocketmq.common.message.MessageExt;
Expand All @@ -17,8 +18,6 @@ enum RocketMqConsumerExperimentalAttributeExtractor
implements AttributesExtractor<MessageExt, Void> {
INSTANCE;

private static final AttributeKey<String> MESSAGING_ROCKETMQ_TAGS =
AttributeKey.stringKey("messaging.rocketmq.tags");
private static final AttributeKey<Long> MESSAGING_ROCKETMQ_QUEUE_ID =
AttributeKey.longKey("messaging.rocketmq.queue_id");
private static final AttributeKey<Long> MESSAGING_ROCKETMQ_QUEUE_OFFSET =
Expand All @@ -30,7 +29,7 @@ enum RocketMqConsumerExperimentalAttributeExtractor
public void onStart(AttributesBuilder attributes, Context parentContext, MessageExt msg) {
String tags = msg.getTags();
if (tags != null) {
attributes.put(MESSAGING_ROCKETMQ_TAGS, tags);
attributes.put(SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG, tags);
}
attributes.put(MESSAGING_ROCKETMQ_QUEUE_ID, msg.getQueueId());
attributes.put(MESSAGING_ROCKETMQ_QUEUE_OFFSET, msg.getQueueOffset());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,15 +9,14 @@
import io.opentelemetry.api.common.AttributesBuilder;
import io.opentelemetry.context.Context;
import io.opentelemetry.instrumentation.api.instrumenter.AttributesExtractor;
import io.opentelemetry.semconv.trace.attributes.SemanticAttributes;
import javax.annotation.Nullable;
import org.apache.rocketmq.client.hook.SendMessageContext;

enum RocketMqProducerExperimentalAttributeExtractor
implements AttributesExtractor<SendMessageContext, Void> {
INSTANCE;

private static final AttributeKey<String> MESSAGING_ROCKETMQ_TAGS =
AttributeKey.stringKey("messaging.rocketmq.tags");
private static final AttributeKey<String> MESSAGING_ROCKETMQ_BROKER_ADDRESS =
AttributeKey.stringKey("messaging.rocketmq.broker_address");
private static final AttributeKey<String> MESSAGING_ROCKETMQ_SEND_RESULT =
Expand All @@ -29,7 +28,7 @@ public void onStart(
if (request.getMessage() != null) {
String tags = request.getMessage().getTags();
if (tags != null) {
attributes.put(MESSAGING_ROCKETMQ_TAGS, tags);
attributes.put(SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG, tags);
}
}
String brokerAddr = request.getBrokerAddr();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_DESTINATION" sharedTopic
"$SemanticAttributes.MESSAGING_DESTINATION_KIND" "topic"
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagA"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagA"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.send_result" "SEND_OK"
"messaging.payload" String
Expand All @@ -120,7 +120,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_OPERATION" "process"
"$SemanticAttributes.MESSAGING_MESSAGE_PAYLOAD_SIZE_BYTES" Long
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagA"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagA"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.queue_id" Long
"messaging.rocketmq.queue_offset" Long
Expand Down Expand Up @@ -161,7 +161,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_DESTINATION" sharedTopic
"$SemanticAttributes.MESSAGING_DESTINATION_KIND" "topic"
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagA"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagA"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.send_result" "SEND_OK"
"messaging.payload" String
Expand All @@ -178,7 +178,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_OPERATION" "process"
"$SemanticAttributes.MESSAGING_MESSAGE_PAYLOAD_SIZE_BYTES" Long
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagA"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagA"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.queue_id" Long
"messaging.rocketmq.queue_offset" Long
Expand Down Expand Up @@ -267,7 +267,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_OPERATION" "process"
"$SemanticAttributes.MESSAGING_MESSAGE_PAYLOAD_SIZE_BYTES" Long
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagA"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagA"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.queue_id" Long
"messaging.rocketmq.queue_offset" Long
Expand All @@ -286,7 +286,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_OPERATION" "process"
"$SemanticAttributes.MESSAGING_MESSAGE_PAYLOAD_SIZE_BYTES" Long
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagB"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagB"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.queue_id" Long
"messaging.rocketmq.queue_offset" Long
Expand Down Expand Up @@ -331,7 +331,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_DESTINATION" sharedTopic
"$SemanticAttributes.MESSAGING_DESTINATION_KIND" "topic"
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagA"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagA"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.send_result" "SEND_OK"
"messaging.header.test_message_header" { it == ["test"] }
Expand All @@ -348,7 +348,7 @@ abstract class AbstractRocketMqClientTest extends InstrumentationSpecification {
"$SemanticAttributes.MESSAGING_OPERATION" "process"
"$SemanticAttributes.MESSAGING_MESSAGE_PAYLOAD_SIZE_BYTES" Long
"$SemanticAttributes.MESSAGING_MESSAGE_ID" String
"messaging.rocketmq.tags" "TagA"
"$SemanticAttributes.MESSAGING_ROCKETMQ_MESSAGE_TAG" "TagA"
"messaging.rocketmq.broker_address" String
"messaging.rocketmq.queue_id" Long
"messaging.rocketmq.queue_offset" Long
Expand Down

0 comments on commit f708d03

Please sign in to comment.