diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java index ddb794808c..12208d3de4 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java @@ -193,8 +193,6 @@ public class KafkaDslTests { assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(i + 1); assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0); assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo((long) i); - assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME"); - assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048633L); assertThat(headers.get("foo")).isEqualTo("bar"); } @@ -211,8 +209,6 @@ public class KafkaDslTests { assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(i + 1); assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0); assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo((long) i); - assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME"); - assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048644L); } Message message = MessageBuilder.withPayload("BAR").setHeader(KafkaHeaders.TOPIC, TEST_TOPIC2).build(); @@ -360,8 +356,7 @@ public class KafkaDslTests { .enrichHeaders(h -> h.header(KafkaIntegrationHeaders.FUTURE_TOKEN, "foo")) .publishSubscribeChannel(c -> c .subscribe(sf -> sf.handle( - kafkaMessageHandler(producerFactory(), TEST_TOPIC1) - .timestampExpression("T(Long).valueOf('1487694048633')"), + kafkaMessageHandler(producerFactory(), TEST_TOPIC1), e -> e.id("kafkaProducer1"))) .subscribe(sf -> sf.handle(kafkaMessageHandlerTopic2, e -> e.id("kafkaProducer2"))) ); @@ -370,8 +365,7 @@ public class KafkaDslTests { @Bean public KafkaProducerMessageHandlerSpec kafkaMessageHandlerTopic2() { return kafkaMessageHandler(producerFactory(), TEST_TOPIC2) - .flush(msg -> true) - .timestamp(m -> 1487694048644L); + .flush(msg -> true); } private KafkaProducerMessageHandlerSpec kafkaMessageHandler(