This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] 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 09941a1781 [Improve][Connector-V2] Make Couchbase sink readiness
timeout configurable (#12168)
09941a1781 is described below
commit 09941a178186ca004eb55af377b9757d8dab1871
Author: Goutam Adwant <[email protected]>
AuthorDate: Fri Sep 11 02:49:45 2026 +0000
[Improve][Connector-V2] Make Couchbase sink readiness timeout configurable
(#12168)
Signed-off-by: Goutam Adwant <[email protected]>
---
docs/en/connectors/sink/Couchbase.md | 17 +++
docs/zh/connectors/sink/Couchbase.md | 14 ++
.../couchbase/config/CouchbaseSinkOptions.java | 9 ++
.../couchbase/sink/CouchbaseSinkFactory.java | 5 +
.../seatunnel/couchbase/sink/CouchbaseWriter.java | 7 +-
.../couchbase/sink/CouchbaseWriterOptions.java | 25 +++
.../couchbase/sink/CouchbaseSinkFactoryTest.java | 157 +++++++++++++++++++
.../sink/CouchbaseWriterConstructorLeakTest.java | 39 +++++
.../couchbase/sink/CouchbaseWriterOptionsTest.java | 77 ++++++++++
.../connector/couchbase/CouchbaseReadinessIT.java | 167 +++++++++++++++++++++
10 files changed, 516 insertions(+), 1 deletion(-)
diff --git a/docs/en/connectors/sink/Couchbase.md
b/docs/en/connectors/sink/Couchbase.md
index 1304718d51..57767d6f7c 100644
--- a/docs/en/connectors/sink/Couchbase.md
+++ b/docs/en/connectors/sink/Couchbase.md
@@ -78,12 +78,29 @@ Couchbase stores JSON documents. The connector maps
SeaTunnel types to JSON valu
| bucket | String | Yes | - | Target
bucket name. |
| scope | String | No | `_default` | Target scope
name within the bucket. |
| collection | String | Yes | - | Target
collection name. |
+| ready.timeout | Integer | No | `30` | Maximum
seconds to wait for the target bucket to become ready during writer
initialization. Must be greater than zero. |
| primary-key | `List<String>` | No | - | Field names
used to build the document key (length-prefixed encoding: `<len>:<value>`
components separated by `#`). A random UUID is used when not set. |
| upsert-enable | Boolean | No | `false` | Enable
upsert (insert-or-replace) mode. When `false`, duplicate keys will cause an
error. |
| buffer-flush.max-rows | Integer | No | `1000` | Maximum rows
to buffer before a batch write is triggered. Use `-1` to disable. |
| retry.max | Integer | No | `3` | Maximum
retry attempts on transient write failure. |
| retry.interval | Long | No | `1000` | Base
milliseconds for linear retry delay. Attempt `n` waits `retry.interval × n` ms.
|
+### Startup readiness
+
+`ready.timeout` controls the bucket-readiness wait during writer
initialization. The default
+remains 30 seconds. For a cluster that needs more time to become available,
set a larger positive
+value, for example `ready.timeout = 60`.
+
+The value is in seconds, not milliseconds. No connector-specific upper limit
is enforced;
+choose the smallest budget that covers the cluster's observed recovery time.
Excessively large
+values can delay writer-initialization failure when the bucket remains
unavailable.
+
+The Couchbase SDK handles connection attempts within this wait; the connector
does not add an
+outer bootstrap retry loop. `retry.max` and `retry.interval` still apply only
to writes. This
+option does not change individual SDK operation timeouts or the engine's
job-startup timeout.
+An expired readiness wait still fails writer initialization and disconnects
the client.
+Increasing the budget does not correct invalid credentials, incorrect
addresses or a missing bucket.
+
## Security
### TLS / encrypted transport
diff --git a/docs/zh/connectors/sink/Couchbase.md
b/docs/zh/connectors/sink/Couchbase.md
index 43883e9544..b19e806395 100644
--- a/docs/zh/connectors/sink/Couchbase.md
+++ b/docs/zh/connectors/sink/Couchbase.md
@@ -73,12 +73,26 @@ sh bin/install-plugin.sh ${version}
| bucket | String | 是 | - | 目标 Bucket
名称。 |
| scope | String | 否 | `_default` | Bucket 中的目标
Scope 名称。 |
| collection | String | 是 | - | 目标
Collection 名称。 |
+| ready.timeout | Integer | 否 | `30` | 写入器初始化时等待目标
bucket 就绪的最长时间(秒),必须大于零。 |
| primary-key | `List<String>` | 否 | - |
用于构建文档键的字段名列表(长度前缀编码:`<长度>:<值>` 分量以 `#` 分隔)。未设置时使用随机 UUID。 |
| upsert-enable | Boolean | 否 | `false` | 是否启用
Upsert(插入或替换)模式。为 `false` 时,重复键将报错。 |
| buffer-flush.max-rows | Integer | 否 | `1000` |
触发批量写入的最大缓冲行数。设为 `-1` 禁用。 |
| retry.max | Integer | 否 | `3` |
写入失败时的最大重试次数。 |
| retry.interval | Long | 否 | `1000` |
线性退避基础间隔(毫秒)。第 n 次重试等待 `retry.interval × n` 毫秒。 |
+### 启动就绪等待
+
+`ready.timeout` 控制写入器初始化时等待目标 bucket 就绪的时间,默认仍为 30 秒。
+对于需要更长时间才能恢复可用的集群,可配置更大的正数,例如 `ready.timeout = 60`。
+
+该值的单位是秒,而不是毫秒。连接器不额外限制上限,应根据集群实际恢复时间选择足够的最小值。
+当 bucket 持续不可用时,过大的值会延迟写入器初始化失败的报告。
+
+Couchbase SDK 在此等待期间处理连接尝试,连接器不会额外添加启动重试循环。
+`retry.max` 和 `retry.interval` 仍仅用于写入重试。此选项不会修改 SDK 单次操作的超时时间,
+也不会修改引擎的作业启动超时时间。等待超时后,写入器初始化仍会失败并断开客户端连接。
+增加等待时间不能修复无效凭据、错误地址或不存在的 bucket。
+
## 安全性
### TLS / 加密传输
diff --git
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
index 4d09a802c5..9596c32dad 100644
---
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
+++
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
@@ -25,6 +25,15 @@ import java.util.List;
/** Configuration options specific to the Couchbase sink connector. */
public class CouchbaseSinkOptions extends CouchbaseConfig {
+ /** Maximum time to wait for bucket readiness during writer
initialization, in seconds. */
+ public static final Option<Integer> READY_TIMEOUT =
+ Options.key("ready.timeout")
+ .intType()
+ .defaultValue(30)
+ .withDescription(
+ "The timeout in seconds for waiting until the
target bucket is ready"
+ + " during writer initialization. Must be
greater than zero.");
+
/**
* Maximum number of rows buffered before a batch write is triggered.
*
diff --git
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
index 630580311e..a4ebbc265b 100644
---
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
+++
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
@@ -18,6 +18,7 @@
package org.apache.seatunnel.connectors.seatunnel.couchbase.sink;
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.Conditions;
import org.apache.seatunnel.api.configuration.util.OptionRule;
import org.apache.seatunnel.api.table.catalog.CatalogTable;
import org.apache.seatunnel.api.table.catalog.TableIdentifier;
@@ -60,6 +61,9 @@ public class CouchbaseSinkFactory implements TableSinkFactory
{
CouchbaseSinkOptions.RETRY_INTERVAL,
CouchbaseSinkOptions.UPSERT_ENABLE,
CouchbaseSinkOptions.PRIMARY_KEY)
+ .optional(
+ CouchbaseSinkOptions.READY_TIMEOUT,
+
Conditions.greaterThan(CouchbaseSinkOptions.READY_TIMEOUT, 0))
.build();
}
@@ -88,6 +92,7 @@ public class CouchbaseSinkFactory implements TableSinkFactory
{
.withUsername(config.get(CouchbaseSinkOptions.USERNAME))
.withPassword(config.get(CouchbaseSinkOptions.PASSWORD))
.withBucket(config.get(CouchbaseSinkOptions.BUCKET))
+
.withReadyTimeout(config.get(CouchbaseSinkOptions.READY_TIMEOUT))
.withScope(config.get(CouchbaseSinkOptions.SCOPE))
.withCollection(config.get(CouchbaseSinkOptions.COLLECTION));
diff --git
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
index 5da720e4ad..7b4b4e9c48 100644
---
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
+++
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
@@ -115,7 +115,12 @@ public class CouchbaseWriter implements
SinkWriter<SeaTunnelRow, Void, Void> {
options.getPassword());
Collection resolvedCollection;
try {
-
connectedCluster.bucket(options.getBucket()).waitUntilReady(Duration.ofSeconds(30));
+ log.debug(
+ "Waiting up to {} seconds for Couchbase bucket readiness",
+ options.getReadyTimeout());
+ connectedCluster
+ .bucket(options.getBucket())
+
.waitUntilReady(Duration.ofSeconds(options.getReadyTimeout()));
resolvedCollection =
connectedCluster
.bucket(options.getBucket())
diff --git
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
index 324176d276..fe75caf449 100644
---
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
+++
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
@@ -17,6 +17,8 @@
package org.apache.seatunnel.connectors.seatunnel.couchbase.sink;
+import
org.apache.seatunnel.connectors.seatunnel.couchbase.config.CouchbaseSinkOptions;
+
import lombok.Getter;
import java.io.Serializable;
@@ -42,6 +44,7 @@ public class CouchbaseWriterOptions implements Serializable {
private final String[] primaryKey;
private final int retryMax;
private final long retryInterval;
+ private final int readyTimeout;
private CouchbaseWriterOptions(Builder builder) {
this.connectionString = builder.connectionString;
@@ -55,12 +58,18 @@ public class CouchbaseWriterOptions implements Serializable
{
this.primaryKey = builder.primaryKey;
this.retryMax = builder.retryMax;
this.retryInterval = builder.retryInterval;
+ this.readyTimeout = builder.readyTimeout;
}
public static Builder builder() {
return new Builder();
}
+ /** Retains the previous readiness budget for options serialized before
this field existed. */
+ public int getReadyTimeout() {
+ return readyTimeout == 0 ?
CouchbaseSinkOptions.READY_TIMEOUT.defaultValue() : readyTimeout;
+ }
+
/** Fluent builder for {@link CouchbaseWriterOptions}. */
public static class Builder {
private String connectionString;
@@ -74,6 +83,7 @@ public class CouchbaseWriterOptions implements Serializable {
private String[] primaryKey = new String[0];
private int retryMax = 3;
private long retryInterval = 1000L;
+ private int readyTimeout =
CouchbaseSinkOptions.READY_TIMEOUT.defaultValue();
public Builder withConnectionString(String connectionString) {
this.connectionString = connectionString;
@@ -130,6 +140,21 @@ public class CouchbaseWriterOptions implements
Serializable {
return this;
}
+ /**
+ * Sets the bucket-readiness budget used during writer initialization.
+ *
+ * @param readyTimeout positive readiness timeout in seconds
+ * @return this builder
+ * @throws IllegalArgumentException if the timeout is zero or negative
+ */
+ public Builder withReadyTimeout(int readyTimeout) {
+ if (readyTimeout <= 0) {
+ throw new IllegalArgumentException("'ready.timeout' must be
greater than zero.");
+ }
+ this.readyTimeout = readyTimeout;
+ return this;
+ }
+
public CouchbaseWriterOptions build() {
return new CouchbaseWriterOptions(this);
}
diff --git
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactoryTest.java
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactoryTest.java
new file mode 100644
index 0000000000..e0547924e8
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactoryTest.java
@@ -0,0 +1,157 @@
+/*
+ * 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.seatunnel.couchbase.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.TableIdentifier;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import com.couchbase.client.java.Bucket;
+import com.couchbase.client.java.Cluster;
+import com.couchbase.client.java.Collection;
+import com.couchbase.client.java.Scope;
+
+import java.time.Duration;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class CouchbaseSinkFactoryTest {
+
+ @Test
+ void testConfiguredReadinessTimeoutReachesWriter() throws Exception {
+ Map<String, Object> config = baseConfig();
+ config.put("ready.timeout", 60);
+ verifyReadinessTimeout(config, Duration.ofSeconds(60));
+ }
+
+ @Test
+ void testDefaultReadinessTimeoutRemainsThirtySeconds() throws Exception {
+ verifyReadinessTimeout(baseConfig(), Duration.ofSeconds(30));
+ }
+
+ @ParameterizedTest
+ @ValueSource(ints = {0, -1})
+ void testNonPositiveReadinessTimeoutRejectedBeforeConnecting(int timeout) {
+ Map<String, Object> config = baseConfig();
+ config.put("ready.timeout", timeout);
+ try (MockedStatic<Cluster> staticCluster =
Mockito.mockStatic(Cluster.class)) {
+ OptionValidationException error =
+ assertThrows(OptionValidationException.class, () ->
createSink(config));
+ assertTrue(error.getMessage().contains("ready.timeout"));
+ staticCluster.verifyNoInteractions();
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(ints = {0, -1})
+ void testDirectFactoryAlsoRejectsNonPositiveTimeout(int timeout) {
+ Map<String, Object> config = baseConfig();
+ config.put("ready.timeout", timeout);
+ try (MockedStatic<Cluster> staticCluster =
Mockito.mockStatic(Cluster.class)) {
+ IllegalArgumentException error =
+ assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ new CouchbaseSinkFactory()
+ .createSink(
+ new
TableSinkFactoryContext(
+ null,
+
ReadonlyConfig.fromMap(config),
+
getClass().getClassLoader())));
+ assertTrue(error.getMessage().contains("ready.timeout"));
+ staticCluster.verifyNoInteractions();
+ }
+ }
+
+ @Test
+ void testWriteRetriesDoNotChangeReadinessTimeout() throws Exception {
+ Map<String, Object> config = baseConfig();
+ config.put("retry.max", 10);
+ config.put("retry.interval", 5000L);
+ verifyReadinessTimeout(config, Duration.ofSeconds(30));
+ }
+
+ private void verifyReadinessTimeout(Map<String, Object> config, Duration
expected)
+ throws Exception {
+ Cluster cluster = mock(Cluster.class);
+ Bucket bucket = mock(Bucket.class);
+ Scope scope = mock(Scope.class);
+ Collection collection = mock(Collection.class);
+ when(cluster.bucket("test_bucket")).thenReturn(bucket);
+ when(bucket.scope("_default")).thenReturn(scope);
+ when(scope.collection("_default")).thenReturn(collection);
+
+ try (MockedStatic<Cluster> staticCluster =
Mockito.mockStatic(Cluster.class)) {
+ staticCluster
+ .when(() -> Cluster.connect("couchbase://localhost",
"user", "pass"))
+ .thenReturn(cluster);
+ CouchbaseWriter writer = createSink(config).createWriter(null);
+ try {
+ verify(bucket).waitUntilReady(expected);
+ } finally {
+ writer.close();
+ }
+ verify(cluster).disconnect();
+ }
+ }
+
+ private CouchbaseSink createSink(Map<String, Object> config) {
+ CouchbaseSinkFactory factory = new CouchbaseSinkFactory();
+ ReadonlyConfig options = ReadonlyConfig.fromMap(config);
+ ConfigValidator.of(options).validate(factory.optionRule());
+ CatalogTable table =
+ CatalogTable.of(
+ TableIdentifier.of("catalog", "database", "table"),
+ TableSchema.builder().build(),
+ Collections.emptyMap(),
+ Collections.emptyList(),
+ "");
+ return (CouchbaseSink)
+ factory.createSink(
+ new TableSinkFactoryContext(
+ table, options,
getClass().getClassLoader()))
+ .createSink();
+ }
+
+ private Map<String, Object> baseConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put("connection.string", "couchbase://localhost");
+ config.put("username", "user");
+ config.put("password", "pass");
+ config.put("bucket", "test_bucket");
+ config.put("collection", "_default");
+ return config;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
index 7d5ce27ef5..df73d5ac12 100644
---
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
+++
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
@@ -23,16 +23,22 @@ import
org.apache.seatunnel.api.table.catalog.TableIdentifier;
import org.apache.seatunnel.api.table.catalog.TableSchema;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.MockedStatic;
import org.mockito.Mockito;
+import com.couchbase.client.core.error.AuthenticationFailureException;
+import com.couchbase.client.core.error.UnambiguousTimeoutException;
import com.couchbase.client.java.Bucket;
import com.couchbase.client.java.Cluster;
import com.couchbase.client.java.Collection;
import com.couchbase.client.java.Scope;
import java.time.Duration;
+import java.util.stream.Stream;
+import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
@@ -59,6 +65,39 @@ import static org.mockito.Mockito.when;
*/
class CouchbaseWriterConstructorLeakTest {
+ @ParameterizedTest
+ @MethodSource("readinessFailures")
+ void testReadinessFailureIsNotRetriedOrMasked(RuntimeException failure) {
+ Bucket bucket = mock(Bucket.class);
+ doThrow(failure).when(bucket).waitUntilReady(any(Duration.class));
+ Cluster cluster = mockCluster(bucket);
+ RuntimeException disconnectFailure = new RuntimeException("disconnect
failed");
+ doThrow(disconnectFailure).when(cluster).disconnect();
+
+ try (MockedStatic<Cluster> staticCluster =
Mockito.mockStatic(Cluster.class)) {
+ staticCluster
+ .when(() -> Cluster.connect(anyString(), anyString(),
anyString()))
+ .thenReturn(cluster);
+ RuntimeException thrown =
+ assertThrows(
+ RuntimeException.class,
+ () ->
+ new CouchbaseWriter(
+ minimalOptions(),
minimalCatalogTable(), null));
+ assertSame(failure, thrown);
+ assertSame(disconnectFailure, thrown.getSuppressed()[0]);
+ verify(bucket).waitUntilReady(Duration.ofSeconds(30));
+ verify(cluster).disconnect();
+ }
+ }
+
+ private static Stream<RuntimeException> readinessFailures() {
+ return Stream.of(
+ new UnambiguousTimeoutException("readiness timed out", null),
+ new AuthenticationFailureException("authentication failed",
null, null),
+ new RuntimeException(new InterruptedException("readiness
interrupted")));
+ }
+
/** Minimal {@link CouchbaseWriterOptions} pointing at a fake cluster. */
private static CouchbaseWriterOptions minimalOptions() {
return CouchbaseWriterOptions.builder()
diff --git
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptionsTest.java
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptionsTest.java
new file mode 100644
index 0000000000..da2d7098a2
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptionsTest.java
@@ -0,0 +1,77 @@
+/*
+ * 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.seatunnel.couchbase.sink;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.ObjectInputStream;
+import java.io.ObjectOutputStream;
+import java.util.Base64;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CouchbaseWriterOptionsTest {
+
+ // Serialized default options from before readyTimeout was added
(serialVersionUID = 1).
+ private static final String LEGACY_OPTIONS =
+
"rO0ABXNyAE9vcmcuYXBhY2hlLnNlYXR1bm5lbC5jb25uZWN0b3JzLnNlYXR1bm5lbC5jb3VjaGJhc2Uuc2luay5Db3VjaGJhc2VXcml0ZXJPcHRpb25zAAAAAAAAAAECAAtJAAlmbHVzaFNpemVKAA1yZXRyeUludGVydmFsSQAIcmV0cnlNYXhaAAx1cHNlcnRFbmFibGVMAAZidWNrZXR0ABJMamF2YS9sYW5nL1N0cmluZztMAApjb2xsZWN0aW9ucQB+AAFMABBjb25uZWN0aW9uU3RyaW5ncQB+AAFMAAhwYXNzd29yZHEAfgABWwAKcHJpbWFyeUtleXQAE1tMamF2YS9sYW5nL1N0cmluZztMAAVzY29wZXEAfgABTAAIdXNlcm5hbWVxAH4AAXhwAAAD6AAAAAAAAAPoAAAAAwBwcHBwdXIAE1tMamF2YS5sYW5nLlN0cmluZzut0lbn6R17RwI
[...]
+
+ @Test
+ void testDefaultReadinessTimeout() {
+ assertEquals(30,
CouchbaseWriterOptions.builder().build().getReadyTimeout());
+ }
+
+ @ParameterizedTest
+ @ValueSource(ints = {0, -1})
+ void testBuilderRejectsNonPositiveReadinessTimeout(int timeout) {
+ IllegalArgumentException error =
+ assertThrows(
+ IllegalArgumentException.class,
+ () ->
CouchbaseWriterOptions.builder().withReadyTimeout(timeout));
+ assertTrue(error.getMessage().contains("ready.timeout"));
+ }
+
+ @Test
+ void testConfiguredReadinessTimeoutSurvivesSerialization() throws
Exception {
+ ByteArrayOutputStream bytes = new ByteArrayOutputStream();
+ try (ObjectOutputStream output = new ObjectOutputStream(bytes)) {
+
output.writeObject(CouchbaseWriterOptions.builder().withReadyTimeout(60).build());
+ }
+ assertEquals(60, deserialize(bytes.toByteArray()).getReadyTimeout());
+ }
+
+ @Test
+ void testLegacySerializedOptionsRetainThirtySecondTimeout() throws
Exception {
+ CouchbaseWriterOptions options =
deserialize(Base64.getDecoder().decode(LEGACY_OPTIONS));
+ assertEquals(30, options.getReadyTimeout());
+ assertEquals(3, options.getRetryMax());
+ assertEquals("_default", options.getScope());
+ }
+
+ private CouchbaseWriterOptions deserialize(byte[] bytes) throws Exception {
+ try (ObjectInputStream input = new ObjectInputStream(new
ByteArrayInputStream(bytes))) {
+ return (CouchbaseWriterOptions) input.readObject();
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseReadinessIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseReadinessIT.java
new file mode 100644
index 0000000000..372a15855a
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseReadinessIT.java
@@ -0,0 +1,167 @@
+/*
+ * 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.e2e.connector.couchbase;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
+import org.apache.seatunnel.api.table.catalog.TableIdentifier;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.connectors.seatunnel.couchbase.sink.CouchbaseSink;
+import
org.apache.seatunnel.connectors.seatunnel.couchbase.sink.CouchbaseSinkFactory;
+import
org.apache.seatunnel.connectors.seatunnel.couchbase.sink.CouchbaseWriter;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.api.Timeout;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.couchbase.BucketDefinition;
+import org.testcontainers.couchbase.CouchbaseContainer;
+import org.testcontainers.couchbase.CouchbaseService;
+import org.testcontainers.utility.DockerImageName;
+
+import com.couchbase.client.core.error.UnambiguousTimeoutException;
+import com.couchbase.client.java.Cluster;
+import lombok.extern.slf4j.Slf4j;
+
+import java.time.Duration;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/** Factory-level readiness tests; no engine containers are needed for this
client contract. */
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@Timeout(value = 2, unit = TimeUnit.MINUTES, threadMode =
Timeout.ThreadMode.SAME_THREAD)
+@Slf4j
+class CouchbaseReadinessIT {
+
+ private final CouchbaseContainer server =
+ new CouchbaseContainer(
+ DockerImageName.parse(
+ System.getProperty(
+ "couchbase.test.image",
+
"couchbase/server:community-7.1.1")))
+ .withNetworkMode("bridge")
+ .withCredentials("Administrator", "password")
+ .withEnabledServices(CouchbaseService.KV)
+ .withBucket(
+ new BucketDefinition("readiness")
+ .withQuota(128)
+ .withReplicas(0)
+ .withPrimaryIndex(false))
+ .withStartupTimeout(Duration.ofMinutes(3))
+ .withStartupAttempts(3)
+ .withLogConsumer(new
Slf4jLogConsumer(log).withPrefix("couchbase-readiness"));
+
+ private Cluster verification;
+
+ @BeforeAll
+ void startServer() {
+ server.start();
+ verification =
+ Cluster.connect(
+ server.getConnectionString(), server.getUsername(),
server.getPassword());
+
verification.bucket("readiness").waitUntilReady(Duration.ofSeconds(60));
+ }
+
+ @AfterAll
+ void closeServer() {
+ try {
+ if (verification != null) {
+ verification.disconnect();
+ }
+ } finally {
+ server.stop();
+ }
+ }
+
+ @Test
+ void testReadinessTimeoutAndRecovery() throws Exception {
+ CouchbaseSink unavailableSink = createSink(5, server.getPassword());
+
server.getDockerClient().pauseContainerCmd(server.getContainerId()).exec();
+ try {
+ assertThrows(
+ UnambiguousTimeoutException.class, () ->
unavailableSink.createWriter(null));
+ } finally {
+
server.getDockerClient().unpauseContainerCmd(server.getContainerId()).exec();
+ }
+
+ CouchbaseWriter writer = createSink(60,
server.getPassword()).createWriter(null);
+ try {
+ writer.write(new SeaTunnelRow(new Object[] {"recovered"}));
+ writer.prepareCommit();
+ assertEquals(
+ "recovered",
+ verification
+ .bucket("readiness")
+ .defaultCollection()
+ .get("9:recovered")
+ .contentAsObject()
+ .getString("id"));
+ } finally {
+ writer.close();
+ }
+ }
+
+ @Test
+ void testInvalidCredentialsStillFailReadiness() {
+ CouchbaseSink sink = createSink(5, "incorrect-password");
+ assertThrows(UnambiguousTimeoutException.class, () ->
sink.createWriter(null));
+ }
+
+ private CouchbaseSink createSink(int timeout, String password) {
+ Map<String, Object> config = new HashMap<>();
+ config.put("connection.string", server.getConnectionString());
+ config.put("username", server.getUsername());
+ config.put("password", password);
+ config.put("bucket", "readiness");
+ config.put("collection", "_default");
+ config.put("primary-key", Collections.singletonList("id"));
+ config.put("upsert-enable", true);
+ config.put("ready.timeout", timeout);
+ CatalogTable table =
+ CatalogTable.of(
+ TableIdentifier.of("catalog", "database", "table"),
+ TableSchema.builder()
+ .column(
+ PhysicalColumn.of(
+ "id", BasicType.STRING_TYPE,
64L, false, null, ""))
+ .build(),
+ Collections.emptyMap(),
+ Collections.emptyList(),
+ "");
+ CouchbaseSinkFactory factory = new CouchbaseSinkFactory();
+ ReadonlyConfig options = ReadonlyConfig.fromMap(config);
+ ConfigValidator.of(options).validate(factory.optionRule());
+ return (CouchbaseSink)
+ factory.createSink(
+ new TableSinkFactoryContext(
+ table, options,
getClass().getClassLoader()))
+ .createSink();
+ }
+}