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 03ef9f1c00 [Improve][Connector-V2][AmazonDynamoDB] Migrate max retries
validation to OptionRule (#11821)
03ef9f1c00 is described below
commit 03ef9f1c0010a60eb29490ceead1309ad825c555
Author: Goutam Adwant <[email protected]>
AuthorDate: Wed Sep 2 02:38:10 2026 +0000
[Improve][Connector-V2][AmazonDynamoDB] Migrate max retries validation to
OptionRule (#11821)
Signed-off-by: goutamadwant <[email protected]>
---
docs/en/connectors/sink/AmazonDynamoDB.md | 3 +-
docs/zh/connectors/sink/AmazonDynamoDB.md | 3 +-
.../config/AmazonDynamoDBConfig.java | 6 --
.../sink/AmazonDynamoDBSinkFactory.java | 3 +-
.../AmazonDynamoDBSinkFactoryTest.java | 104 +++++++++++++++++++++
5 files changed, 110 insertions(+), 9 deletions(-)
diff --git a/docs/en/connectors/sink/AmazonDynamoDB.md
b/docs/en/connectors/sink/AmazonDynamoDB.md
index 8567832c4d..93261bedb5 100644
--- a/docs/en/connectors/sink/AmazonDynamoDB.md
+++ b/docs/en/connectors/sink/AmazonDynamoDB.md
@@ -33,7 +33,7 @@ The target table must already exist. The connector writes
each row as a DynamoDB
| table | string | yes | - | DynamoDB table
name to write to. |
| batch_size | int | no | 25 | Records buffered
for one batch write request. |
| multi_table_sink_replica | int | no | - | Sink writer
replicas for each table. |
-| max_retries | int | no | 10 | Retries for
unprocessed items. |
+| max_retries | int | no | 10 | Retries for
unprocessed items. Must be at least `0`. |
| retry_base_delay_ms | long | no | 100 | Initial retry
backoff delay in milliseconds. |
| retry_max_delay_ms | long | no | 5000 | Maximum retry
backoff delay in milliseconds. |
| common-options | object | no | - | Sink plugin common
parameters. |
@@ -76,6 +76,7 @@ Optional common sink option used by multi-table sink jobs.
For details, see [Sin
### max_retries [int]
The maximum number of retries when DynamoDB returns unprocessed items from a
batch write request.
+A value of `0` disables retries. The value must not be negative.
### retry_base_delay_ms [long]
diff --git a/docs/zh/connectors/sink/AmazonDynamoDB.md
b/docs/zh/connectors/sink/AmazonDynamoDB.md
index 10bfbb06dc..8d0fde0d4d 100644
--- a/docs/zh/connectors/sink/AmazonDynamoDB.md
+++ b/docs/zh/connectors/sink/AmazonDynamoDB.md
@@ -33,7 +33,7 @@ Amazon DynamoDB 写入连接器用于将 SeaTunnel 数据行写入 DynamoDB 表
| table | string | 是 | - | 要写入的 DynamoDB 表名。 |
| batch_size | int | 否 | 25 | 一次批量写入请求缓存的记录数。 |
| multi_table_sink_replica | int | 否 | - | 每张表对应的 Sink Writer 副本数。 |
-| max_retries | int | 否 | 10 | 未处理 item 的最大重试次数。 |
+| max_retries | int | 否 | 10 | 未处理 item 的最大重试次数,必须大于等于 `0`。 |
| retry_base_delay_ms | long | 否 | 100 | 初始重试等待时间,单位毫秒。 |
| retry_max_delay_ms | long | 否 | 5000 | 最大重试等待时间,单位毫秒。 |
| common-options | object | 否 | - | Sink 插件通用参数。 |
@@ -76,6 +76,7 @@ DynamoDB batch write 每次最多支持 25 条写请求,所以默认值为 `25
### max_retries [int]
当 DynamoDB 在批量写入结果中返回未处理 item 时,最多重试的次数。
+设置为 `0` 表示不重试,该值不能为负数。
### retry_base_delay_ms [long]
diff --git
a/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/config/AmazonDynamoDBConfig.java
b/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/config/AmazonDynamoDBConfig.java
index d6554fe6e7..09f3c94ae4 100644
---
a/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/config/AmazonDynamoDBConfig.java
+++
b/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/config/AmazonDynamoDBConfig.java
@@ -64,12 +64,6 @@ public class AmazonDynamoDBConfig implements Serializable {
this.scanItemLimit =
config.get(AmazonDynamoDBSourceOptions.SCAN_ITEM_LIMIT);
this.parallelScanThreads =
config.get(AmazonDynamoDBSourceOptions.PARALLEL_SCAN_THREADS);
this.maxRetries = config.get(AmazonDynamoDBSinkOptions.MAX_RETRIES);
- if (this.maxRetries < 0) {
- throw new IllegalArgumentException(
- String.format(
- "max_retries must be a non-negative integer, but
got: %d",
- this.maxRetries));
- }
this.retryBaseDelayMs =
config.get(AmazonDynamoDBSinkOptions.RETRY_BASE_DELAY_MS);
this.retryMaxDelayMs =
config.get(AmazonDynamoDBSinkOptions.RETRY_MAX_DELAY_MS);
}
diff --git
a/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/sink/AmazonDynamoDBSinkFactory.java
b/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/sink/AmazonDynamoDBSinkFactory.java
index 6dcafdc8d3..0ee85f32f0 100644
---
a/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/sink/AmazonDynamoDBSinkFactory.java
+++
b/seatunnel-connectors-v2/connector-amazondynamodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/sink/AmazonDynamoDBSinkFactory.java
@@ -17,6 +17,7 @@
package org.apache.seatunnel.connectors.seatunnel.amazondynamodb.sink;
+import org.apache.seatunnel.api.configuration.util.Conditions;
import org.apache.seatunnel.api.configuration.util.OptionRule;
import org.apache.seatunnel.api.options.SinkConnectorCommonOptions;
import org.apache.seatunnel.api.table.connector.TableSink;
@@ -51,9 +52,9 @@ public class AmazonDynamoDBSinkFactory implements
TableSinkFactory {
.optional(
BATCH_SIZE,
SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA,
- MAX_RETRIES,
RETRY_BASE_DELAY_MS,
RETRY_MAX_DELAY_MS)
+ .optional(MAX_RETRIES, Conditions.greaterOrEqual(MAX_RETRIES,
0))
.build();
}
diff --git
a/seatunnel-connectors-v2/connector-amazondynamodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/AmazonDynamoDBSinkFactoryTest.java
b/seatunnel-connectors-v2/connector-amazondynamodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/AmazonDynamoDBSinkFactoryTest.java
new file mode 100644
index 0000000000..c28de5a71a
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-amazondynamodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/amazondynamodb/AmazonDynamoDBSinkFactoryTest.java
@@ -0,0 +1,104 @@
+/*
+ * 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.amazondynamodb;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import
org.apache.seatunnel.connectors.seatunnel.amazondynamodb.sink.AmazonDynamoDBSinkFactory;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import static
org.apache.seatunnel.connectors.seatunnel.amazondynamodb.config.AmazonDynamoDBSinkOptions.ACCESS_KEY_ID;
+import static
org.apache.seatunnel.connectors.seatunnel.amazondynamodb.config.AmazonDynamoDBSinkOptions.MAX_RETRIES;
+import static
org.apache.seatunnel.connectors.seatunnel.amazondynamodb.config.AmazonDynamoDBSinkOptions.REGION;
+import static
org.apache.seatunnel.connectors.seatunnel.amazondynamodb.config.AmazonDynamoDBSinkOptions.SECRET_ACCESS_KEY;
+import static
org.apache.seatunnel.connectors.seatunnel.amazondynamodb.config.AmazonDynamoDBSinkOptions.TABLE;
+import static
org.apache.seatunnel.connectors.seatunnel.amazondynamodb.config.AmazonDynamoDBSinkOptions.URL;
+
+/** Tests declarative sink validation exposed by {@link
AmazonDynamoDBSinkFactory}. */
+public class AmazonDynamoDBSinkFactoryTest {
+
+ private OptionRule sinkRule;
+
+ @BeforeEach
+ void setUp() {
+ sinkRule = new AmazonDynamoDBSinkFactory().optionRule();
+ }
+
+ @Test
+ void testOmittedMaxRetriesUsesDefault() {
+ Map<String, Object> config = requiredOptions();
+
+ Assertions.assertDoesNotThrow(() -> validate(config));
+ Assertions.assertEquals(10,
ReadonlyConfig.fromMap(config).get(MAX_RETRIES));
+ }
+
+ @Test
+ void testNonNegativeMaxRetries() {
+ Assertions.assertDoesNotThrow(() -> validateMaxRetries(0));
+ Assertions.assertDoesNotThrow(() -> validateMaxRetries(3));
+ }
+
+ @Test
+ void testNegativeMaxRetriesFails() {
+ OptionValidationException exception =
+ Assertions.assertThrows(
+ OptionValidationException.class, () ->
validateMaxRetries(-1));
+
+
Assertions.assertTrue(exception.getMessage().contains(MAX_RETRIES.key()));
+ Assertions.assertTrue(exception.getMessage().contains(">= 0"));
+ }
+
+ @Test
+ void testMaxRetriesIsADeclaredOption() {
+ Map<String, Object> config = requiredOptions();
+ config.put(MAX_RETRIES.key(), 3);
+
+ Assertions.assertDoesNotThrow(
+ () ->
+ ConfigValidator.validateUnknownKeys(
+ ReadonlyConfig.fromMap(config), sinkRule,
"AmazonDynamoDB"));
+ }
+
+ private void validateMaxRetries(int maxRetries) {
+ Map<String, Object> config = requiredOptions();
+ config.put(MAX_RETRIES.key(), maxRetries);
+ validate(config);
+ }
+
+ private Map<String, Object> requiredOptions() {
+ Map<String, Object> config = new HashMap<>();
+ config.put(URL.key(), "http://localhost:8000");
+ config.put(REGION.key(), "us-east-1");
+ config.put(ACCESS_KEY_ID.key(), "access-key");
+ config.put(SECRET_ACCESS_KEY.key(), "secret-key");
+ config.put(TABLE.key(), "orders");
+ return config;
+ }
+
+ private void validate(Map<String, Object> config) {
+ ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(sinkRule);
+ }
+}