Attempt to mitigate KafkaDslTests without timestamp

This commit is contained in:
Artem Bilan
2024-12-17 13:49:44 -05:00
parent 5e3d8d890a
commit af4277e2f4

View File

@@ -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<String> 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<Integer, String, ?> kafkaMessageHandlerTopic2() {
return kafkaMessageHandler(producerFactory(), TEST_TOPIC2)
.flush(msg -> true)
.timestamp(m -> 1487694048644L);
.flush(msg -> true);
}
private KafkaProducerMessageHandlerSpec<Integer, String, ?> kafkaMessageHandler(