This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new da7b1e01a1 [Test][E2E] Wait for Kafka exactly-once test sends (#11231)
da7b1e01a1 is described below
commit da7b1e01a16968d87fae55b3dfb12bb2975410f0
Author: Daniel <[email protected]>
AuthorDate: Sun Aug 16 22:42:25 2026 +0800
[Test][E2E] Wait for Kafka exactly-once test sends (#11231)
Co-authored-by: DanielLeens <[email protected]>
---
.../org/apache/seatunnel/e2e/connector/kafka/KafkaIT.java | 15 +++------------
1 file changed, 3 insertions(+), 12 deletions(-)
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-kafka-e2e/src/test/java/org/apache/seatunnel/e2e/connector/kafka/KafkaIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-kafka-e2e/src/test/java/org/apache/seatunnel/e2e/connector/kafka/KafkaIT.java
index 2a7faa6af4..626330bc8c 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-kafka-e2e/src/test/java/org/apache/seatunnel/e2e/connector/kafka/KafkaIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-kafka-e2e/src/test/java/org/apache/seatunnel/e2e/connector/kafka/KafkaIT.java
@@ -1990,10 +1990,7 @@ public class KafkaIT extends TestSuiteBase implements
TestResource {
final String jobId = Long.toUnsignedString(System.nanoTime());
long sinkStartOffset = endOffsetOnP0(consumerTopic);
for (int i = 0; i < 10; i++) {
- ProducerRecord<byte[], byte[]> record =
- new ProducerRecord<>(producerTopic, null,
sourceData.getBytes());
- producer.send(record);
- producer.flush();
+ sendTextRecordAndWait(producerTopic, sourceData);
}
// async execute
CompletableFuture.supplyAsync(
@@ -2025,10 +2022,7 @@ public class KafkaIT extends TestSuiteBase implements
TestResource {
long restoreStartOffset = endOffsetOnP0(consumerTopic);
for (int i = 0; i < 10; i++) {
- ProducerRecord<byte[], byte[]> record =
- new ProducerRecord<>(producerTopic, null,
sourceDataRestore.getBytes());
- producer.send(record);
- producer.flush();
+ sendTextRecordAndWait(producerTopic, sourceDataRestore);
}
CompletableFuture.runAsync(
@@ -2122,10 +2116,7 @@ public class KafkaIT extends TestSuiteBase implements
TestResource {
String sourceData = "Seatunnel Exactly Once Example";
long sinkStartOffset = endOffsetOnP0(consumerTopic);
for (int i = 0; i < 10; i++) {
- ProducerRecord<byte[], byte[]> record =
- new ProducerRecord<>(producerTopic, null,
sourceData.getBytes());
- producer.send(record);
- producer.flush();
+ sendTextRecordAndWait(producerTopic, sourceData);
}
Long endOffset;
KafkaConsumer<String, String> consumer = null;