Handle sleuth headers sent from spring-kafka as byte[] in Spring Cloud (#925)
Similar to https://github.com/openzipkin/brave/blob/master/instrumentation/kafka-clients/src/main/java/brave/kafka/clients/KafkaPropagation.java
This commit is contained in:
@@ -29,6 +29,7 @@ import org.springframework.messaging.support.NativeMessageHeaderAccessor;
|
||||
import org.springframework.util.LinkedMultiValueMap;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import static java.nio.charset.StandardCharsets.UTF_8;
|
||||
import static org.springframework.messaging.support.NativeMessageHeaderAccessor.NATIVE_HEADERS;
|
||||
|
||||
/**
|
||||
@@ -127,7 +128,13 @@ enum MessageHeaderPropagation
|
||||
}
|
||||
}
|
||||
Object result = accessor.getHeader(key);
|
||||
return result != null ? result.toString() : null;
|
||||
if (result != null) {
|
||||
if (result instanceof byte[]) {
|
||||
return new String((byte[]) result, UTF_8);
|
||||
}
|
||||
return result.toString();
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
static void removeAnyTraceHeaders(MessageHeaderAccessor accessor,
|
||||
|
||||
@@ -19,8 +19,13 @@ package org.springframework.cloud.sleuth.instrument.messaging;
|
||||
import java.util.Collections;
|
||||
|
||||
import brave.propagation.Propagation;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
import org.springframework.messaging.support.MessageHeaderAccessor;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.junit.Assert.*;
|
||||
|
||||
public class MessageHeaderPropagationTest
|
||||
extends PropagationSetterTest<MessageHeaderAccessor, String> {
|
||||
MessageHeaderAccessor carrier = new MessageHeaderAccessor();
|
||||
@@ -43,4 +48,31 @@ public class MessageHeaderPropagationTest
|
||||
Collections.singleton(result.toString()) :
|
||||
Collections.emptyList();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetByteArrayValue() {
|
||||
MessageHeaderAccessor carrier = carrier();
|
||||
carrier.setHeader("X-B3-TraceId", "48485a3953bb6124".getBytes());
|
||||
carrier.setHeader("X-B3-TraceId", "48485a3953bb6124000000".getBytes());
|
||||
String value = MessageHeaderPropagation.INSTANCE.get(carrier, "X-B3-TraceId");
|
||||
assertEquals("48485a3953bb6124000000", value);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetStringValue() {
|
||||
MessageHeaderAccessor carrier = carrier();
|
||||
carrier.setHeader("X-B3-TraceId", "48485a3953bb6124");
|
||||
carrier.setHeader("X-B3-TraceId", "48485a3953bb61240000000");
|
||||
String value = MessageHeaderPropagation.INSTANCE.get(carrier, "X-B3-TraceId");
|
||||
assertEquals("48485a3953bb61240000000", value);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetNullValue() {
|
||||
MessageHeaderAccessor carrier = carrier();
|
||||
carrier.setHeader("X-B3-TraceId", "48485a3953bb6124");
|
||||
carrier.setHeader("X-B3-TraceId", "48485a3953bb61240000000");
|
||||
String value = MessageHeaderPropagation.INSTANCE.get(carrier, "non existent key");
|
||||
assertNull(value);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user