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 d2ec8ce3c1 [Fix][Connector-V2] Fix Kafka exactly-once sink losing
first record on checkpoint (#11541)
d2ec8ce3c1 is described below
commit d2ec8ce3c107e685d85893e9a51aeaa777fb4f3d
Author: FatMan <[email protected]>
AuthorDate: Fri Aug 14 15:24:42 2026 +0800
[Fix][Connector-V2] Fix Kafka exactly-once sink losing first record on
checkpoint (#11541)
Co-authored-by: chen <[email protected]>
Co-authored-by: davidzollo <[email protected]>
---
docs/en/connectors/sink/Kafka.md | 11 +
docs/zh/connectors/sink/Kafka.md | 9 +
.../kafka/exception/KafkaConnectorErrorCode.java | 6 +-
.../kafka/sink/KafkaTransactionSender.java | 83 +++++-
.../kafka/sink/KafkaTransactionSenderTest.java | 282 +++++++++++++++++++--
.../seatunnel/e2e/connector/kafka/KafkaIT.java | 9 +
6 files changed, 382 insertions(+), 18 deletions(-)
diff --git a/docs/en/connectors/sink/Kafka.md b/docs/en/connectors/sink/Kafka.md
index 5b800fc314..40c8b2b0a1 100644
--- a/docs/en/connectors/sink/Kafka.md
+++ b/docs/en/connectors/sink/Kafka.md
@@ -473,6 +473,17 @@ sink {
Ensure the Kafka broker has transactions enabled and that
`transaction.timeout.ms` is aligned with your checkpoint interval.
+Under `EXACTLY_ONCE`, a failed send fails the checkpoint instead of silently
dropping records. Two
+errors can be reported in that situation:
+
+| Code | Name | Meaning
| What to do
|
+|----------|-------------------------|-------------------------------------------------------------------------------|--------------------------------------------------------------------------------------------------------------|
+| KAFKA-08 | TRANSACTION_NOT_STARTED | The transaction carries records but
Kafka never registered it on the broker. | Check broker availability and
whether `transaction.timeout.ms` is shorter than the checkpoint interval.
|
+| KAFKA-09 | PRODUCE_DATA_FAILED | A record of the transaction failed to
be sent asynchronously. | Read the exception cause; retriable
causes usually recover on checkpoint retry, others need broker-side work. |
+
+Both errors abort the current transaction, so the affected records are re-sent
from the last
+completed checkpoint rather than lost.
+
### How do I configure SASL/Kerberos authentication?
Pass broker authentication settings via `kafka.*` properties:
diff --git a/docs/zh/connectors/sink/Kafka.md b/docs/zh/connectors/sink/Kafka.md
index d7c17065ee..1a17c2057e 100644
--- a/docs/zh/connectors/sink/Kafka.md
+++ b/docs/zh/connectors/sink/Kafka.md
@@ -473,6 +473,15 @@ sink {
确保 Kafka Broker 开启了事务支持,且 `transaction.timeout.ms` 与 checkpoint 间隔相匹配。
+在 `EXACTLY_ONCE` 语义下,发送失败会让 checkpoint 失败,而不是静默丢弃数据。此时可能出现两种错误:
+
+| 错误码 | 名称 | 含义
| 处理建议 |
+|----------|-------------------------|-------------------------------------------|------------------------------------------------------------------|
+| KAFKA-08 | TRANSACTION_NOT_STARTED | 事务中已有数据,但 Kafka 始终未在 Broker 端完成该事务的注册。
| 检查 Broker 是否可用,以及 `transaction.timeout.ms` 是否小于 checkpoint 间隔。 |
+| KAFKA-09 | PRODUCE_DATA_FAILED | 事务中的某条数据异步发送失败。
| 查看异常 cause;可重试的异常通常在 checkpoint 重试后恢复,其他异常需要排查 Broker 端问题。 |
+
+两种错误都会中止当前事务,受影响的数据会从上一个已完成的 checkpoint 重新发送,不会丢失。
+
### 如何配置 SASL/Kerberos 认证?
```hocon
diff --git
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/exception/KafkaConnectorErrorCode.java
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/exception/KafkaConnectorErrorCode.java
index cb5c09764d..e9eabe024e 100644
---
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/exception/KafkaConnectorErrorCode.java
+++
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/exception/KafkaConnectorErrorCode.java
@@ -29,7 +29,11 @@ public enum KafkaConnectorErrorCode implements
SeaTunnelErrorCode {
CONSUME_THREAD_RUN_ERROR(
"KAFKA-05", "Error occurred when the kafka consumer thread was
running"),
CONSUME_DATA_FAILED("KAFKA-06", "Kafka failed to consume data"),
- CONSUMER_CLOSE_FAILED("KAFKA-07", "Kafka failed to close consumer");
+ CONSUMER_CLOSE_FAILED("KAFKA-07", "Kafka failed to close consumer"),
+ TRANSACTION_NOT_STARTED(
+ "KAFKA-08",
+ "Kafka transaction still reported as not started after flushing
pending sends"),
+ PRODUCE_DATA_FAILED("KAFKA-09", "Kafka failed to produce data");
private final String code;
private final String description;
diff --git
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSender.java
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSender.java
index cc6f71f4d6..0f7bd65c52 100644
---
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSender.java
+++
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSender.java
@@ -20,11 +20,14 @@ package
org.apache.seatunnel.connectors.seatunnel.kafka.sink;
import org.apache.seatunnel.shade.com.google.common.collect.Lists;
import org.apache.seatunnel.connectors.seatunnel.kafka.KafkaClientUtils;
+import
org.apache.seatunnel.connectors.seatunnel.kafka.exception.KafkaConnectorErrorCode;
+import
org.apache.seatunnel.connectors.seatunnel.kafka.exception.KafkaConnectorException;
import org.apache.seatunnel.connectors.seatunnel.kafka.state.KafkaCommitInfo;
import org.apache.seatunnel.connectors.seatunnel.kafka.state.KafkaSinkState;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.clients.producer.RecordMetadata;
import lombok.extern.slf4j.Slf4j;
@@ -32,6 +35,7 @@ import java.time.Duration;
import java.util.List;
import java.util.Optional;
import java.util.Properties;
+import java.util.concurrent.atomic.AtomicReference;
import static
org.apache.seatunnel.connectors.seatunnel.kafka.sink.KafkaSinkWriter.generateTransactionId;
@@ -48,6 +52,15 @@ public class KafkaTransactionSender<K, V> implements
KafkaProduceSender<K, V> {
private String transactionId;
private final String transactionPrefix;
private final Properties kafkaProperties;
+
+ /**
+ * Holds the first asynchronous send failure of the current transaction.
It is written from the
+ * producer's sender thread through the send callback and read by the task
thread, so it must
+ * stay thread-safe. It is scoped to a single transaction and therefore
reset by {@link
+ * #beginTransaction(String)}.
+ */
+ private final AtomicReference<Exception> asyncSendException = new
AtomicReference<>();
+
private int recordNumInTransaction = 0;
public KafkaTransactionSender(String transactionPrefix, Properties
kafkaProperties) {
@@ -57,30 +70,96 @@ public class KafkaTransactionSender<K, V> implements
KafkaProduceSender<K, V> {
@Override
public void send(ProducerRecord<K, V> producerRecord) {
- kafkaProducer.send(producerRecord);
+ // Surface an already recorded asynchronous failure on the write path.
The current
+ // transaction can no longer be committed, so buffering and
transmitting more records for it
+ // would only waste producer memory and network bandwidth until the
next checkpoint.
+ checkAsyncSendException();
+ kafkaProducer.send(producerRecord, this::onSendCompleted);
recordNumInTransaction++;
}
+ /**
+ * Records the first asynchronous send failure of the current transaction
so that {@link
+ * #prepareCommit()} can fail the checkpoint instead of committing a
partial transaction.
+ *
+ * <p>Invoked on the producer's sender thread.
+ */
+ private void onSendCompleted(RecordMetadata metadata, Exception exception)
{
+ if (exception == null) {
+ return;
+ }
+ if (!asyncSendException.compareAndSet(null, exception)) {
+ // Only the first failure becomes the checkpoint failure cause.
Log the later ones so a
+ // broker-side incident affecting several partitions can still be
diagnosed.
+ log.warn(
+ "Suppressed an additional asynchronous send failure of
Kafka transaction [{}]",
+ transactionId,
+ exception);
+ }
+ }
+
@Override
public void beginTransaction(String transactionId) {
this.transactionId = transactionId;
this.kafkaProducer = getTransactionProducer(transactionId);
kafkaProducer.beginTransaction();
+ // Reset the per-transaction state. A new transaction always runs on a
newly created
+ // producer, so a failure recorded for the previous transaction no
longer applies. Keeping
+ // it would turn a single transient send error into a permanent
checkpoint failure loop.
recordNumInTransaction = 0;
+ asyncSendException.set(null);
}
@Override
public Optional<KafkaCommitInfo> prepareCommit() {
+ // Flush pending async sends before capturing the transaction state.
Kafka only marks the
+ // transaction as started once the AddPartitionsToTxn request has been
acknowledged by the
+ // broker, and that request is issued asynchronously by the producer's
sender thread. If a
+ // checkpoint reaches this point before the first record's transaction
registration
+ // completes, isTxnStarted() would still return false and the
resulting commit info would
+ // instruct the committer to skip EndTxn, leaving the transaction to
time out and its
+ // records permanently invisible to read_committed consumers.
+ kafkaProducer.flush();
+ checkAsyncSendException();
+ boolean txnStarted = kafkaProducer.isTxnStarted();
+ if (recordNumInTransaction > 0 && !txnStarted) {
+ // Records were sent in this transaction but Kafka still reports
it as not started even
+ // after flushing, meaning the transaction registration never
completed. Committing with
+ // txnStarted=false would make the committer skip EndTxn and drop
these records, so fail
+ // fast and let the checkpoint abort this transaction instead of
silently producing a
+ // lossy commit info.
+ throw new KafkaConnectorException(
+ KafkaConnectorErrorCode.TRANSACTION_NOT_STARTED,
+ String.format(
+ "Kafka transaction [%s] has %d record(s) but is
still reported as not "
+ + "started after flushing pending sends.
The transaction "
+ + "registration did not complete. Refusing
to commit to avoid "
+ + "data loss.",
+ transactionId, recordNumInTransaction));
+ }
KafkaCommitInfo kafkaCommitInfo =
new KafkaCommitInfo(
transactionId,
kafkaProperties,
this.kafkaProducer.getProducerId(),
this.kafkaProducer.getEpoch(),
- this.kafkaProducer.isTxnStarted());
+ txnStarted);
return Optional.of(kafkaCommitInfo);
}
+ private void checkAsyncSendException() {
+ Exception exception = asyncSendException.get();
+ if (exception != null) {
+ throw new KafkaConnectorException(
+ KafkaConnectorErrorCode.PRODUCE_DATA_FAILED,
+ String.format(
+ "Kafka transaction [%s] failed to send one or more
of its %d record(s) "
+ + "asynchronously.",
+ transactionId, recordNumInTransaction),
+ exception);
+ }
+ }
+
@Override
public void abortTransaction() {
kafkaProducer.abortTransaction();
diff --git
a/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSenderTest.java
b/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSenderTest.java
index 8b34fc8237..372579d9f4 100644
---
a/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSenderTest.java
+++
b/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSenderTest.java
@@ -17,7 +17,12 @@
package org.apache.seatunnel.connectors.seatunnel.kafka.sink;
-import org.apache.kafka.clients.producer.ProducerConfig;
+import
org.apache.seatunnel.connectors.seatunnel.kafka.exception.KafkaConnectorErrorCode;
+import
org.apache.seatunnel.connectors.seatunnel.kafka.exception.KafkaConnectorException;
+import org.apache.seatunnel.connectors.seatunnel.kafka.state.KafkaCommitInfo;
+
+import org.apache.kafka.clients.producer.Callback;
+import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@@ -25,12 +30,180 @@ import org.mockito.Mockito;
import java.time.Duration;
import java.util.ArrayDeque;
+import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Deque;
import java.util.List;
+import java.util.Optional;
import java.util.Properties;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
class KafkaTransactionSenderTest {
+ private static final String TRANSACTION_PREFIX = "SeaTunnel0001";
+ private static final String TRANSACTION_ID = TRANSACTION_PREFIX + "-1";
+ private static final String NEXT_TRANSACTION_ID = TRANSACTION_PREFIX +
"-2";
+ private static final String TOPIC = "test-topic";
+ private static final long PRODUCER_ID = 1001L;
+ private static final short EPOCH = 5;
+
+ /**
+ * Reproduces the reported race: Kafka reports the transaction as not
started until the pending
+ * AddPartitionsToTxn request has been acknowledged, which flush() forces.
prepareCommit() must
+ * flush first so the captured commit info carries txnStarted=true and the
committer performs
+ * EndTxn.
+ */
+ @Test
+ void prepareCommitFlushesBeforeCapturingTransactionState() {
+ ProducerStub producer = new ProducerStub();
+ producer.startTransactionOnFlush();
+
+ TestingKafkaTransactionSender sender = createSender(producer);
+ sender.beginTransaction(TRANSACTION_ID);
+ sender.send(record());
+
+ Optional<KafkaCommitInfo> commitInfo = sender.prepareCommit();
+
+ Assertions.assertTrue(commitInfo.isPresent());
+ Assertions.assertTrue(
+ commitInfo.get().isTxnStarted(),
+ "txnStarted must be captured after flush so the committer
performs EndTxn");
+ Assertions.assertEquals(TRANSACTION_ID,
commitInfo.get().getTransactionId());
+ Assertions.assertEquals(PRODUCER_ID, commitInfo.get().getProducerId());
+ Assertions.assertEquals(EPOCH, commitInfo.get().getEpoch());
+ verify(producer.mock, times(1)).flush();
+ }
+
+ /**
+ * When the transaction genuinely carries records but Kafka still reports
it as not started even
+ * after flushing, prepareCommit() must fail fast rather than emit a
commit info that instructs
+ * the committer to skip EndTxn and silently drop those records.
+ */
+ @Test
+ void prepareCommitFailsWhenRecordsSentButTransactionNotStarted() {
+ ProducerStub producer = new ProducerStub();
+
+ TestingKafkaTransactionSender sender = createSender(producer);
+ sender.beginTransaction(TRANSACTION_ID);
+ sender.send(record());
+ sender.send(record());
+
+ KafkaConnectorException exception =
+ Assertions.assertThrows(KafkaConnectorException.class,
sender::prepareCommit);
+
+ Assertions.assertEquals(
+ KafkaConnectorErrorCode.TRANSACTION_NOT_STARTED.getCode(),
+ exception.getSeaTunnelErrorCode().getCode());
+ verify(producer.mock, times(1)).flush();
+ }
+
+ /**
+ * A transaction can already be marked as started while a record still
fails to be sent. The
+ * failure is only reported once flush() completes the pending send, so
prepareCommit() must
+ * check for it after flushing instead of committing the remaining
successful records.
+ */
+ @Test
+ void prepareCommitFailsWhenAsyncSendFailsAfterTransactionStarted() {
+ ProducerStub producer = new ProducerStub();
+ producer.transactionStarted();
+ RuntimeException asyncSendFailure = new RuntimeException("async send
failed");
+ producer.failPendingSendsWith(asyncSendFailure);
+
+ TestingKafkaTransactionSender sender = createSender(producer);
+ sender.beginTransaction(TRANSACTION_ID);
+ sender.send(record());
+
+ KafkaConnectorException exception =
+ Assertions.assertThrows(KafkaConnectorException.class,
sender::prepareCommit);
+
+ Assertions.assertEquals(
+ KafkaConnectorErrorCode.PRODUCE_DATA_FAILED.getCode(),
+ exception.getSeaTunnelErrorCode().getCode());
+ Assertions.assertSame(asyncSendFailure, exception.getCause());
+ verify(producer.mock, times(1)).flush();
+ }
+
+ /**
+ * An empty transaction (no records sent) legitimately reports
txnStarted=false and must not be
+ * treated as an error; the commit info simply carries txnStarted=false.
+ */
+ @Test
+ void prepareCommitAllowsEmptyTransactionWithoutRecords() {
+ ProducerStub producer = new ProducerStub();
+
+ TestingKafkaTransactionSender sender = createSender(producer);
+ sender.beginTransaction(TRANSACTION_ID);
+
+ Optional<KafkaCommitInfo> commitInfo = sender.prepareCommit();
+
+ Assertions.assertTrue(commitInfo.isPresent());
+ Assertions.assertFalse(commitInfo.get().isTxnStarted());
+ verify(producer.mock, times(1)).flush();
+ }
+
+ /**
+ * Once a failure has been recorded the current transaction can no longer
be committed, so the
+ * write path must reject further records instead of buffering them until
the next checkpoint.
+ */
+ @Test
+ void sendFailsFastAfterAsyncSendFailureIsRecorded() {
+ ProducerStub producer = new ProducerStub();
+
+ TestingKafkaTransactionSender sender = createSender(producer);
+ sender.beginTransaction(TRANSACTION_ID);
+ sender.send(record());
+ producer.completePendingSends(new RuntimeException("async send
failed"));
+
+ KafkaConnectorException exception =
+ Assertions.assertThrows(KafkaConnectorException.class, () ->
sender.send(record()));
+
+ Assertions.assertEquals(
+ KafkaConnectorErrorCode.PRODUCE_DATA_FAILED.getCode(),
+ exception.getSeaTunnelErrorCode().getCode());
+ Assertions.assertEquals(
+ 1, producer.sendCount, "no further record may be handed to the
producer");
+ }
+
+ /**
+ * An asynchronous send failure is scoped to the transaction that produced
it. After the engine
+ * aborts that transaction and opens a new one on the same sender,
checkpoints must succeed
+ * again; otherwise one transient broker error would block every later
checkpoint.
+ *
+ * <p>The recovered transaction is deliberately empty and reports
txnStarted=false, so it also
+ * proves the record counter was reset: a stale counter would raise
TRANSACTION_NOT_STARTED.
+ */
+ @Test
+ void senderRecoversAfterFailedTransaction() {
+ ProducerStub failingProducer = new ProducerStub();
+ failingProducer.transactionStarted();
+ failingProducer.failPendingSendsWith(new RuntimeException("async send
failed"));
+ ProducerStub healthyProducer = new ProducerStub();
+
+ TestingKafkaTransactionSender sender = createSender(failingProducer,
healthyProducer);
+ sender.beginTransaction(TRANSACTION_ID);
+ sender.send(record());
+ Assertions.assertThrows(KafkaConnectorException.class,
sender::prepareCommit);
+ sender.abortTransaction();
+
+ sender.beginTransaction(NEXT_TRANSACTION_ID);
+
+ Optional<KafkaCommitInfo> commitInfo =
Assertions.assertDoesNotThrow(sender::prepareCommit);
+ Assertions.assertTrue(commitInfo.isPresent());
+ Assertions.assertEquals(NEXT_TRANSACTION_ID,
commitInfo.get().getTransactionId());
+ Assertions.assertFalse(commitInfo.get().isTxnStarted());
+ }
+
+ /**
+ * Each transactional ID must be fenced by its own producer: changing the
ID on a reused
+ * producer can retain a non-zero epoch indefinitely and spin the cleanup
loop forever.
+ */
@Test
void abortTransactionUsesFreshProducerForEachTransactionalId() {
KafkaInternalProducer<byte[], byte[]> existingTransaction =
@@ -46,40 +219,119 @@ class KafkaTransactionSenderTest {
sender.abortTransaction(7L);
Assertions.assertEquals(
- Arrays.asList("test-prefix-7", "test-prefix-8"),
sender.createdTransactionIds);
+ Arrays.asList(TRANSACTION_PREFIX + "-7", TRANSACTION_PREFIX +
"-8"),
+ sender.createdTransactionIds);
Mockito.verify(existingTransaction).close(Duration.ZERO);
Mockito.verify(unusedTransaction).close(Duration.ZERO);
}
+ private static ProducerRecord<byte[], byte[]> record() {
+ return new ProducerRecord<>(TOPIC, new byte[] {1});
+ }
+
+ /**
+ * Builds a sender that hands out the given producer stubs, one per {@code
beginTransaction}.
+ */
+ private TestingKafkaTransactionSender createSender(ProducerStub...
producers) {
+ KafkaInternalProducer<byte[], byte[]>[] mocks = new
KafkaInternalProducer[producers.length];
+ for (int i = 0; i < producers.length; i++) {
+ mocks[i] = producers[i].mock;
+ }
+ return new TestingKafkaTransactionSender(mocks);
+ }
+
+ /**
+ * A sender that returns pre-built producers instead of connecting to a
broker, so the real
+ * transaction lifecycle can be driven from a unit test.
+ */
private static class TestingKafkaTransactionSender
extends KafkaTransactionSender<byte[], byte[]> {
- private final ArrayDeque<KafkaInternalProducer<byte[], byte[]>>
producers;
- private final List<String> createdTransactionIds = new
java.util.ArrayList<>();
+ private final Deque<KafkaInternalProducer<byte[], byte[]>> producers;
+ private final List<String> createdTransactionIds = new ArrayList<>();
@SafeVarargs
private TestingKafkaTransactionSender(KafkaInternalProducer<byte[],
byte[]>... producers) {
- super("test-prefix", kafkaProperties());
+ super(TRANSACTION_PREFIX, new Properties());
this.producers = new ArrayDeque<>(Arrays.asList(producers));
}
@Override
protected KafkaInternalProducer<byte[], byte[]>
createTransactionProducer(
String transactionId) {
+ Assertions.assertFalse(
+ producers.isEmpty(), "unexpected producer creation for " +
transactionId);
createdTransactionIds.add(transactionId);
return producers.removeFirst();
}
+ }
+
+ /**
+ * A mocked {@link KafkaInternalProducer} whose {@code flush()} completes
the callbacks of
+ * previously submitted sends, mirroring Kafka's guarantee that flush
returns only once all
+ * buffered records have completed either successfully or exceptionally.
+ */
+ private static final class ProducerStub {
+
+ private final KafkaInternalProducer<byte[], byte[]> mock;
+ private final List<Callback> pendingCallbacks = new ArrayList<>();
+ private final AtomicBoolean txnStarted = new AtomicBoolean(false);
+
+ private boolean startTransactionOnFlush;
+ private Exception sendFailure;
+ private int sendCount;
+
+ @SuppressWarnings("unchecked")
+ private ProducerStub() {
+ this.mock = mock(KafkaInternalProducer.class);
+ when(mock.getProducerId()).thenReturn(PRODUCER_ID);
+ when(mock.getEpoch()).thenReturn(EPOCH);
+ when(mock.isTxnStarted()).thenAnswer(invocation ->
txnStarted.get());
+ doAnswer(
+ invocation -> {
+
pendingCallbacks.add(invocation.getArgument(1));
+ sendCount++;
+ return null;
+ })
+ .when(mock)
+ .send(any(), any());
+ doAnswer(
+ invocation -> {
+ completePendingSends(sendFailure);
+ if (startTransactionOnFlush) {
+ txnStarted.set(true);
+ }
+ return null;
+ })
+ .when(mock)
+ .flush();
+ }
+
+ /** Reports the transaction as started from the beginning. */
+ private void transactionStarted() {
+ txnStarted.set(true);
+ }
+
+ /**
+ * Defers the transaction registration to {@code flush()}, as the
broker acknowledging
+ * AddPartitionsToTxn does.
+ */
+ private void startTransactionOnFlush() {
+ startTransactionOnFlush = true;
+ }
+
+ /** Makes pending and subsequent sends complete exceptionally when
flushed. */
+ private void failPendingSendsWith(Exception exception) {
+ sendFailure = exception;
+ }
- private static Properties kafkaProperties() {
- Properties properties = new Properties();
- properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
"localhost:9092");
- properties.put(
- ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
-
"org.apache.kafka.common.serialization.ByteArraySerializer");
- properties.put(
- ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
-
"org.apache.kafka.common.serialization.ByteArraySerializer");
- return properties;
+ /** Invokes the pending send callbacks, as the producer's sender
thread would. */
+ private void completePendingSends(Exception exception) {
+ List<Callback> callbacks = new ArrayList<>(pendingCallbacks);
+ pendingCallbacks.clear();
+ for (Callback callback : callbacks) {
+ callback.onCompletion(null, exception);
+ }
}
}
}
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 ea1fc779d9..2a7faa6af4 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
@@ -2056,6 +2056,15 @@ public class KafkaIT extends TestSuiteBase implements
TestResource {
sourceDataRestore)));
}
+ /**
+ * Regression guard for issue #11534: under EXACTLY_ONCE semantics a
checkpoint could capture
+ * the transaction state before the first record's asynchronous send had
registered its
+ * partitions with the broker, causing the committer to skip EndTxn and
drop that record
+ * permanently. The checkData assertion below fails on both a lost record
(matched < 10) and
+ * a duplicate (matched > 10), so it exercises the exactly-once
guarantee in both directions
+ * once the send is flushed before the transaction state is captured in
+ * KafkaTransactionSender#prepareCommit.
+ */
@TestTemplate
@DisabledOnContainer(
type = EngineType.SPARK,