diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/SimpleKafkaHeaderMapper.java b/spring-kafka/src/main/java/org/springframework/kafka/support/SimpleKafkaHeaderMapper.java index 1dad51fb..8d1f440c 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/SimpleKafkaHeaderMapper.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/SimpleKafkaHeaderMapper.java @@ -45,7 +45,9 @@ public class SimpleKafkaHeaderMapper extends AbstractKafkaHeaderMapper { * consumer/producer records. */ public SimpleKafkaHeaderMapper() { - super(); + super("!" + MessageHeaders.ID, + "!" + MessageHeaders.TIMESTAMP, + "*"); } /** diff --git a/spring-kafka/src/test/java/org/springframework/kafka/support/SimpleKafkaHeaderMapperTests.java b/spring-kafka/src/test/java/org/springframework/kafka/support/SimpleKafkaHeaderMapperTests.java index 8308c0fd..a2a29f6a 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/support/SimpleKafkaHeaderMapperTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/support/SimpleKafkaHeaderMapperTests.java @@ -38,7 +38,7 @@ public class SimpleKafkaHeaderMapperTests { @Test public void testSpecificStringConvert() { - SimpleKafkaHeaderMapper mapper = new SimpleKafkaHeaderMapper("*"); + SimpleKafkaHeaderMapper mapper = new SimpleKafkaHeaderMapper(); Map rawMappedHeaders = new HashMap<>(); rawMappedHeaders.put("thisOnesAString", true); rawMappedHeaders.put("thisOnesBytes", false); @@ -64,7 +64,7 @@ public class SimpleKafkaHeaderMapperTests { @Test public void testNotStringConvert() { - SimpleKafkaHeaderMapper mapper = new SimpleKafkaHeaderMapper("*"); + SimpleKafkaHeaderMapper mapper = new SimpleKafkaHeaderMapper(); Map rawMappedHeaders = new HashMap<>(); rawMappedHeaders.put("thisOnesBytes", false); mapper.setRawMappedHeaders(rawMappedHeaders); @@ -87,7 +87,7 @@ public class SimpleKafkaHeaderMapperTests { @Test public void testAlwaysStringConvert() { - SimpleKafkaHeaderMapper mapper = new SimpleKafkaHeaderMapper("*"); + SimpleKafkaHeaderMapper mapper = new SimpleKafkaHeaderMapper(); mapper.setMapAllStringsOut(true); Map rawMappedHeaders = new HashMap<>(); rawMappedHeaders.put("thisOnesBytes", false); @@ -111,4 +111,22 @@ public class SimpleKafkaHeaderMapperTests { entry("neverConverted", "baz".getBytes())); } + @Test + public void testDefaultHeaderPatterns() { + SimpleKafkaHeaderMapper mapper = new SimpleKafkaHeaderMapper(); + mapper.setMapAllStringsOut(true); + Map headersMap = new HashMap<>(); + headersMap.put(MessageHeaders.ID, "foo".getBytes()); + headersMap.put(MessageHeaders.TIMESTAMP, "bar"); + headersMap.put("thisOnePresent", "baz"); + MessageHeaders headers = new MessageHeaders(headersMap); + Headers target = new RecordHeaders(); + mapper.fromHeaders(headers, target); + assertThat(target).contains( + new RecordHeader("thisOnePresent", "baz".getBytes())); + headersMap.clear(); + mapper.toHeaders(target, headersMap); + assertThat(headersMap).contains( + entry("thisOnePresent", "baz".getBytes())); + } }