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 6dc4ccd238 [Fix][Connector-V2] Handle BigQuery append offset conflicts
correctly (#11561)
6dc4ccd238 is described below
commit 6dc4ccd238605ac55ae817517f8a94dafd1e00af
Author: Doyeon Kim <[email protected]>
AuthorDate: Mon Jul 27 23:50:26 2026 +0900
[Fix][Connector-V2] Handle BigQuery append offset conflicts correctly
(#11561)
---
.../bigquery/sink/BigQuerySinkBatchWriter.java | 58 ++-----
.../bigquery/sink/BigQuerySinkBatchWriterTest.java | 186 +++++++++++++++++++++
2 files changed, 200 insertions(+), 44 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-bigquery/src/main/java/org/apache/seatunnel/connectors/bigquery/sink/BigQuerySinkBatchWriter.java
b/seatunnel-connectors-v2/connector-bigquery/src/main/java/org/apache/seatunnel/connectors/bigquery/sink/BigQuerySinkBatchWriter.java
index 162dd4822d..f9d13449d4 100644
---
a/seatunnel-connectors-v2/connector-bigquery/src/main/java/org/apache/seatunnel/connectors/bigquery/sink/BigQuerySinkBatchWriter.java
+++
b/seatunnel-connectors-v2/connector-bigquery/src/main/java/org/apache/seatunnel/connectors/bigquery/sink/BigQuerySinkBatchWriter.java
@@ -67,8 +67,8 @@ public class BigQuerySinkBatchWriter extends
AbstractBigQuerySinkWriter {
AppendRowsResponse response = future.get(60, TimeUnit.SECONDS);
if (response.hasError()) {
- if (isOffsetConflict(response)) {
- recreateBatchStreamAndRetry(dataToSend);
+ if (isAlreadyExists(response)) {
+ markAppendAsSuccessful(dataToSend);
return;
}
throw new BigQueryConnectorException(
@@ -83,8 +83,8 @@ public class BigQuerySinkBatchWriter extends
AbstractBigQuerySinkWriter {
buffer = dataToSend;
throw new
BigQueryConnectorException(BigQueryConnectorErrorCode.APPEND_ROWS_FAILED, e);
} catch (Exception e) {
- if (isOffsetConflict(e)) {
- recreateBatchStreamAndRetry(dataToSend);
+ if (isAlreadyExists(e)) {
+ markAppendAsSuccessful(dataToSend);
return;
}
buffer = dataToSend;
@@ -92,64 +92,34 @@ public class BigQuerySinkBatchWriter extends
AbstractBigQuerySinkWriter {
}
}
- private void recreateBatchStreamAndRetry(JSONArray dataToSend) {
- log.warn(
- "Detected BigQuery buffered stream offset conflict. "
- + "Recreating buffered stream and retrying append.");
- recreateBatchStream();
- try {
- ApiFuture<AppendRowsResponse> future =
streamWriter.append(dataToSend);
- AppendRowsResponse response = future.get(60, TimeUnit.SECONDS);
-
- if (response.hasError()) {
- throw new BigQueryConnectorException(
- BigQueryConnectorErrorCode.APPEND_ROWS_FAILED,
- response.getError().getMessage());
- }
-
- streamWriter.onAppendSuccess(dataToSend.length());
- log.info("Successfully appended {} rows.", dataToSend.length());
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- buffer = dataToSend;
- throw new
BigQueryConnectorException(BigQueryConnectorErrorCode.APPEND_ROWS_FAILED, e);
- } catch (Exception e) {
- buffer = dataToSend;
- throw new
BigQueryConnectorException(BigQueryConnectorErrorCode.APPEND_ROWS_FAILED, e);
- }
+ private void markAppendAsSuccessful(JSONArray dataToSend) {
+ streamWriter.onAppendSuccess(dataToSend.length());
+ log.info(
+ "BigQuery already accepted the append at the requested offset;
"
+ + "advancing the local offset by {} rows.",
+ dataToSend.length());
}
- private boolean isOffsetConflict(AppendRowsResponse response) {
+ private boolean isAlreadyExists(AppendRowsResponse response) {
if (response == null || !response.hasError()) {
return false;
}
- int code = response.getError().getCode();
- return code == Code.ALREADY_EXISTS_VALUE || code ==
Code.OUT_OF_RANGE_VALUE;
+ return response.getError().getCode() == Code.ALREADY_EXISTS_VALUE;
}
- private boolean isOffsetConflict(Throwable throwable) {
+ private boolean isAlreadyExists(Throwable throwable) {
Throwable current = throwable;
while (current != null) {
if (current instanceof ApiException) {
StatusCode.Code code = ((ApiException)
current).getStatusCode().getCode();
- return code == StatusCode.Code.ALREADY_EXISTS
- || code == StatusCode.Code.OUT_OF_RANGE;
+ return code == StatusCode.Code.ALREADY_EXISTS;
}
current = current.getCause();
}
return false;
}
- private void recreateBatchStream() {
- try {
- streamWriter.close();
- } catch (Exception e) {
- log.warn("Failed to close stale BigQuery buffered stream writer",
e);
- }
- streamWriter = BigQueryBatchWriter.of(client, config);
- }
-
@Override
public void write(SeaTunnelRow element) {
buffer.put(serializer.convert(element, false));
diff --git
a/seatunnel-connectors-v2/connector-bigquery/src/test/java/org/apache/seatunnel/connectors/bigquery/sink/BigQuerySinkBatchWriterTest.java
b/seatunnel-connectors-v2/connector-bigquery/src/test/java/org/apache/seatunnel/connectors/bigquery/sink/BigQuerySinkBatchWriterTest.java
new file mode 100644
index 0000000000..a3749abb34
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-bigquery/src/test/java/org/apache/seatunnel/connectors/bigquery/sink/BigQuerySinkBatchWriterTest.java
@@ -0,0 +1,186 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.bigquery.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import
org.apache.seatunnel.connectors.bigquery.exception.BigQueryConnectorException;
+import org.apache.seatunnel.connectors.bigquery.sink.writer.BigQueryWriter;
+
+import org.json.JSONArray;
+import org.junit.jupiter.api.Test;
+
+import com.google.api.core.ApiFuture;
+import com.google.api.core.ApiFutures;
+import com.google.api.gax.rpc.AlreadyExistsException;
+import com.google.api.gax.rpc.OutOfRangeException;
+import com.google.api.gax.rpc.StatusCode;
+import com.google.cloud.bigquery.storage.v1.AppendRowsResponse;
+import com.google.rpc.Code;
+import com.google.rpc.Status;
+
+import java.io.IOException;
+import java.util.Collections;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+class BigQuerySinkBatchWriterTest {
+
+ @Test
+ void testAlreadyExistsResponseAdvancesOffsetWithoutRecreatingStream() {
+ TestingBigQueryWriter streamWriter =
+ new TestingBigQueryWriter(
+
ApiFutures.immediateFuture(errorResponse(Code.ALREADY_EXISTS)));
+ BigQuerySinkBatchWriter sinkWriter = createSinkWriter(streamWriter);
+ sinkWriter.buffer.put("row-1");
+ sinkWriter.buffer.put("row-2");
+
+ sinkWriter.flush();
+
+ assertEquals(1, streamWriter.appendCount);
+ assertEquals(2, streamWriter.successfulRowCount);
+ assertEquals(0, streamWriter.closeCount);
+ assertEquals(0, sinkWriter.buffer.length());
+ assertSame(streamWriter, sinkWriter.streamWriter);
+ }
+
+ @Test
+ void testAlreadyExistsExceptionAdvancesOffsetWithoutRecreatingStream() {
+ StatusCode statusCode = statusCode(StatusCode.Code.ALREADY_EXISTS);
+ TestingBigQueryWriter streamWriter =
+ new TestingBigQueryWriter(
+ ApiFutures.immediateFailedFuture(
+ new AlreadyExistsException(
+ new IOException("append response was
lost"),
+ statusCode,
+ false)));
+ BigQuerySinkBatchWriter sinkWriter = createSinkWriter(streamWriter);
+ sinkWriter.buffer.put("row-1");
+ sinkWriter.buffer.put("row-2");
+
+ sinkWriter.flush();
+
+ assertEquals(1, streamWriter.appendCount);
+ assertEquals(2, streamWriter.successfulRowCount);
+ assertEquals(0, streamWriter.closeCount);
+ assertEquals(0, sinkWriter.buffer.length());
+ assertSame(streamWriter, sinkWriter.streamWriter);
+ }
+
+ @Test
+ void testOutOfRangeResponseFailsAndRetainsBuffer() {
+ TestingBigQueryWriter streamWriter =
+ new TestingBigQueryWriter(
+
ApiFutures.immediateFuture(errorResponse(Code.OUT_OF_RANGE)));
+ BigQuerySinkBatchWriter sinkWriter = createSinkWriter(streamWriter);
+ sinkWriter.buffer.put("row-1");
+ sinkWriter.buffer.put("row-2");
+
+ assertThrows(BigQueryConnectorException.class, sinkWriter::flush);
+
+ assertEquals(1, streamWriter.appendCount);
+ assertEquals(0, streamWriter.successfulRowCount);
+ assertEquals(0, streamWriter.closeCount);
+ assertEquals(2, sinkWriter.buffer.length());
+ assertSame(streamWriter, sinkWriter.streamWriter);
+ }
+
+ @Test
+ void testOutOfRangeExceptionFailsAndRetainsBuffer() {
+ StatusCode statusCode = statusCode(StatusCode.Code.OUT_OF_RANGE);
+ TestingBigQueryWriter streamWriter =
+ new TestingBigQueryWriter(
+ ApiFutures.immediateFailedFuture(
+ new OutOfRangeException(
+ new IOException("offset is beyond the
stream end"),
+ statusCode,
+ false)));
+ BigQuerySinkBatchWriter sinkWriter = createSinkWriter(streamWriter);
+ sinkWriter.buffer.put("row");
+
+ assertThrows(BigQueryConnectorException.class, sinkWriter::flush);
+
+ assertEquals(1, streamWriter.appendCount);
+ assertEquals(0, streamWriter.successfulRowCount);
+ assertEquals(0, streamWriter.closeCount);
+ assertEquals(1, sinkWriter.buffer.length());
+ assertSame(streamWriter, sinkWriter.streamWriter);
+ }
+
+ private static BigQuerySinkBatchWriter createSinkWriter(BigQueryWriter
streamWriter) {
+ return new BigQuerySinkBatchWriter(
+ ReadonlyConfig.fromMap(Collections.emptyMap()), streamWriter,
null, null);
+ }
+
+ private static AppendRowsResponse errorResponse(Code code) {
+ return AppendRowsResponse.newBuilder()
+ .setError(
+ Status.newBuilder()
+ .setCode(code.getNumber())
+ .setMessage(code.name())
+ .build())
+ .build();
+ }
+
+ private static StatusCode statusCode(StatusCode.Code code) {
+ return new StatusCode() {
+ @Override
+ public StatusCode.Code getCode() {
+ return code;
+ }
+
+ @Override
+ public Object getTransportCode() {
+ return null;
+ }
+ };
+ }
+
+ private static class TestingBigQueryWriter implements BigQueryWriter {
+ private final ApiFuture<AppendRowsResponse> result;
+ private int appendCount;
+ private int successfulRowCount;
+ private int closeCount;
+
+ private TestingBigQueryWriter(ApiFuture<AppendRowsResponse> result) {
+ this.result = result;
+ }
+
+ @Override
+ public ApiFuture<AppendRowsResponse> append(JSONArray jsonArr) {
+ appendCount++;
+ return result;
+ }
+
+ @Override
+ public void onAppendSuccess(int rowCount) {
+ successfulRowCount += rowCount;
+ }
+
+ @Override
+ public void close() {
+ closeCount++;
+ }
+
+ @Override
+ public String getStreamName() {
+ return "test-stream";
+ }
+ }
+}