This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11821-fa001a4e474841c35869d821584dbb0fbc78ff6c in repository https://gitbox.apache.org/repos/asf/seatunnel.git
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); + } +}
