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-12390-9f243e4eb7e56f6e9f71b6844e77d006c85b40c9 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit f80794289d565e002a63844e6cf5f59d3d48076a Author: chovygo <[email protected]> AuthorDate: Sat Oct 3 11:10:56 2026 +0000 [Feature][Connector-V2] Expose MaxCompute client timeout and retry options (#12390) Co-authored-by: zhaowei <[email protected]> --- docs/en/connectors/sink/Maxcompute.md | 38 ++++ docs/en/connectors/source/Maxcompute.md | 38 ++++ docs/zh/connectors/sink/Maxcompute.md | 36 +++ docs/zh/connectors/source/Maxcompute.md | 36 +++ .../maxcompute/catalog/MaxComputeCatalog.java | 9 +- .../catalog/MaxComputeCatalogFactory.java | 3 + .../maxcompute/config/MaxcomputeBaseOptions.java | 52 +++++ .../maxcompute/sink/MaxcomputeSinkFactory.java | 6 + .../maxcompute/source/MaxcomputeSourceFactory.java | 8 +- .../seatunnel/maxcompute/util/MaxcomputeUtil.java | 65 ++++++ .../maxcompute/MaxcomputeSourceFactoryTest.java | 38 ++++ .../maxcompute/catalog/MaxComputeCatalogTest.java | 79 +++++++ .../maxcompute/util/MaxcomputeUtilTest.java | 249 ++++++++++++++++++++- .../test/resources/maxcompute_to_maxcompute.conf | 9 + 14 files changed, 661 insertions(+), 5 deletions(-) diff --git a/docs/en/connectors/sink/Maxcompute.md b/docs/en/connectors/sink/Maxcompute.md index 44e625b3f7..10b9ace223 100644 --- a/docs/en/connectors/sink/Maxcompute.md +++ b/docs/en/connectors/sink/Maxcompute.md @@ -45,6 +45,12 @@ or upsert session selected by `insert_strategy`. | datetime_format | string | no | yyyy-MM-dd HH:mm:ss | Format string used to convert `LocalDateTime` fields to strings. | | tunnel_endpoint | string | no | - | Custom endpoint URL for the MaxCompute Tunnel service. When not set, the endpoint is auto-inferred from the region. | | tunnel_name | string | no | - | Tunnel Quota name used for exclusive resource groups. Requires both `endpoint` and `tunnel_endpoint` to be VPC endpoints. | +| connect_timeout_ms | long | no | 10000 | HTTP connect timeout for the ODPS REST client (metadata/catalog calls) in ms. Default 10000 (10s). | +| read_timeout_ms | long | no | 120000 | HTTP read timeout for the ODPS REST client (metadata/catalog calls) in ms. Default 120000 (120s). | +| retry_times | int | no | 4 | Max retry times for the ODPS REST client. Default 4. | +| tunnel_connect_timeout_ms | long | no | 180000 | HTTP connect timeout for the Tunnel client (data upload/download) in ms. Default 180000 (180s). | +| tunnel_read_timeout_ms | long | no | 300000 | HTTP read timeout for the Tunnel client (data upload/download) in ms. Default 300000 (300s). | +| tunnel_retry_times | int | no | 4 | Max retry times for the Tunnel client. Default 4. | | insert_strategy | string | no | upload | Insert session strategy: `upload` uses an upload session, `upsert` uses an upsert session and requires a primary key. | | multi_table_sink_replica | int | no | 1 | Number of sink writer replicas for each table in a multi-table job. | | common-options | | no | - | Sink plugin common parameters, such as `plugin_input`. | @@ -206,6 +212,38 @@ Example values: Default: Not set (use default quota) +> **Client timeout & retry** +> MaxCompute has two HTTP clients. The **ODPS REST client** handles the control plane +> (table/schema lookup, catalog listing); tune it with `connect_timeout_ms`, +> `read_timeout_ms`, `retry_times`. The **Tunnel client** handles the data plane +> (bulk row upload/download); tune it with the `tunnel_*` options. Setting the REST +> options alone does **not** change the Tunnel client's timeouts. +> Millisecond timeout values are converted to whole seconds; the minimum is `1000`. + +### connect_timeout_ms [long] + +`connect_timeout_ms` HTTP connect timeout for the MaxCompute ODPS REST client, which handles metadata and catalog calls (table/schema lookup, table listing). In milliseconds. Default `10000` (10 seconds). + +### read_timeout_ms [long] + +`read_timeout_ms` HTTP read timeout for the ODPS REST client (metadata/catalog calls) in milliseconds. Default `120000` (120 seconds). Raise this if listing a project with many tables or fetching very wide schemas times out. + +### retry_times [int] + +`retry_times` Maximum retry times for the ODPS REST client on transient failures. Default `4`. + +### tunnel_connect_timeout_ms [long] + +`tunnel_connect_timeout_ms` HTTP connect timeout for the Tunnel client, which performs bulk data upload/download. In milliseconds. Default `180000` (180 seconds). + +### tunnel_read_timeout_ms [long] + +`tunnel_read_timeout_ms` HTTP read timeout for the Tunnel client (bulk data upload/download) in milliseconds. Default `300000` (300 seconds). Raise this when uploading large partitions or upserting large batches whose single write requests exceed 5 minutes. + +### tunnel_retry_times [int] + +`tunnel_retry_times` Maximum retry times for the Tunnel client on transient failures. Default `4`. + ### insert_strategy [string] If `insert_strategy` is set to `upload`, insert operations use an upload session. diff --git a/docs/en/connectors/source/Maxcompute.md b/docs/en/connectors/source/Maxcompute.md index 9f45708ada..472f602366 100644 --- a/docs/en/connectors/source/Maxcompute.md +++ b/docs/en/connectors/source/Maxcompute.md @@ -42,6 +42,12 @@ for parallel reads. | table_list | Array | no | - | List of tables to read. Use this instead of `table_name` to read multiple MaxCompute tables in one job. | | tunnel_endpoint | string | no | - | Custom endpoint URL for the MaxCompute Tunnel service. When not set, the endpoint is auto-inferred from the region. | | tunnel_name | string | no | - | Tunnel Quota name used for exclusive resource groups. Requires both `endpoint` and `tunnel_endpoint` to be VPC endpoints. | +| connect_timeout_ms | long | no | 10000 | HTTP connect timeout for the ODPS REST client (metadata/catalog calls) in ms. Default 10000 (10s). | +| read_timeout_ms | long | no | 120000 | HTTP read timeout for the ODPS REST client (metadata/catalog calls) in ms. Default 120000 (120s). | +| retry_times | int | no | 4 | Max retry times for the ODPS REST client. Default 4. | +| tunnel_connect_timeout_ms | long | no | 180000 | HTTP connect timeout for the Tunnel client (data download) in ms. Default 180000 (180s). | +| tunnel_read_timeout_ms | long | no | 300000 | HTTP read timeout for the Tunnel client (data download) in ms. Default 300000 (300s). | +| tunnel_retry_times | int | no | 4 | Max retry times for the Tunnel client. Default 4. | | schema | config | no | - | Schema of the source table. When `read_columns` is not set, all schema fields are read. | | common-options | | no | - | Source plugin common parameters, such as `plugin_output`. | @@ -138,6 +144,38 @@ Example values: Default: Not set (use default quota) +> **Client timeout & retry** +> MaxCompute has two HTTP clients. The **ODPS REST client** handles the control plane +> (table/schema lookup, catalog listing); tune it with `connect_timeout_ms`, +> `read_timeout_ms`, `retry_times`. The **Tunnel client** handles the data plane +> (bulk row download); tune it with the `tunnel_*` options. Setting the REST options +> alone does **not** change the Tunnel client's timeouts. +> Millisecond timeout values are converted to whole seconds; the minimum is `1000`. + +### connect_timeout_ms [long] + +`connect_timeout_ms` HTTP connect timeout for the MaxCompute ODPS REST client, which handles metadata and catalog calls (table/schema lookup, table listing). In milliseconds. Default `10000` (10 seconds). + +### read_timeout_ms [long] + +`read_timeout_ms` HTTP read timeout for the ODPS REST client (metadata/catalog calls) in milliseconds. Default `120000` (120 seconds). Raise this if listing a project with many tables or fetching very wide schemas times out. + +### retry_times [int] + +`retry_times` Maximum retry times for the ODPS REST client on transient failures. Default `4`. + +### tunnel_connect_timeout_ms [long] + +`tunnel_connect_timeout_ms` HTTP connect timeout for the Tunnel client, which performs bulk data download. In milliseconds. Default `180000` (180 seconds). + +### tunnel_read_timeout_ms [long] + +`tunnel_read_timeout_ms` HTTP read timeout for the Tunnel client (bulk data download) in milliseconds. Default `300000` (300 seconds). Raise this when reading large tables or large partitions whose single read requests exceed 5 minutes. + +### tunnel_retry_times [int] + +`tunnel_retry_times` Maximum retry times for the Tunnel client on transient failures. Default `4`. + ### common options Source plugin common parameters, please refer to [Source Common Options](../common-options/source-common-options.md) for details. diff --git a/docs/zh/connectors/sink/Maxcompute.md b/docs/zh/connectors/sink/Maxcompute.md index c9ccab389f..7d9c7d1996 100644 --- a/docs/zh/connectors/sink/Maxcompute.md +++ b/docs/zh/connectors/sink/Maxcompute.md @@ -41,6 +41,12 @@ import ChangeLog from '../changelog/connector-maxcompute.md'; | datetime_format | string | 否 | yyyy-MM-dd HH:mm:ss | 将 `LocalDateTime` 字段序列化为字符串时使用的格式。 | | tunnel_endpoint | string | 否 | - | MaxCompute Tunnel 服务的自定义端点;未配置时根据区域自动推断。 | | tunnel_name | string | 否 | - | Tunnel Quota 名称;需同时将 `endpoint` 与 `tunnel_endpoint` 配置为 VPC 端点。 | +| connect_timeout_ms | long | 否 | 10000 | ODPS REST 客户端(元数据/catalog 调用)的 HTTP 连接超时,单位毫秒。默认 10000(10 秒)。 | +| read_timeout_ms | long | 否 | 120000 | ODPS REST 客户端(元数据/catalog 调用)的 HTTP 读超时,单位毫秒。默认 120000(120 秒)。 | +| retry_times | int | 否 | 4 | ODPS REST 客户端最大重试次数。默认 4。 | +| tunnel_connect_timeout_ms | long | 否 | 180000 | Tunnel 客户端(批量数据上传/下载)的 HTTP 连接超时,单位毫秒。默认 180000(180 秒)。 | +| tunnel_read_timeout_ms | long | 否 | 300000 | Tunnel 客户端(批量数据上传/下载)的 HTTP 读超时,单位毫秒。默认 300000(300 秒)。 | +| tunnel_retry_times | int | 否 | 4 | Tunnel 客户端最大重试次数。默认 4。 | | insert_strategy | string | 否 | upload | 插入会话类型:`upload` 使用 upload 会话,`upsert` 使用 upsert 会话并要求目标表存在主键。 | | multi_table_sink_replica | int | 否 | 1 | 多表写入时每张表对应的 Sink Writer 副本数。 | | common-options | | 否 | - | Sink 插件通用参数,例如 `plugin_input`。 | @@ -200,6 +206,36 @@ Tunnel Quota 允许您使用专用的计算资源进行 MaxCompute Tunnel 数据 默认值:未设置(使用默认 quota) +> **客户端超时与重试** +> MaxCompute 有两个 HTTP 客户端。**ODPS REST 客户端**负责控制面(表/schema 查询、catalog 列表), +> 用 `connect_timeout_ms`、`read_timeout_ms`、`retry_times` 调整;**Tunnel 客户端**负责数据面(批量行上传/下载), +> 用 `tunnel_*` 系列调整。单独设置 REST 参数**不会**改变 Tunnel 客户端的超时。 +> 毫秒值会被转换为整秒,最小值为 `1000`。 + +### connect_timeout_ms [long] + +`connect_timeout_ms` MaxCompute ODPS REST 客户端的 HTTP 连接超时,该客户端处理元数据与 catalog 调用(表/schema 查询、表列表)。单位毫秒。默认 `10000`(10 秒)。 + +### read_timeout_ms [long] + +`read_timeout_ms` ODPS REST 客户端(元数据/catalog 调用)的 HTTP 读超时,单位毫秒。默认 `120000`(120 秒)。当项目下表非常多或 schema 极宽导致拉取超时时调大。 + +### retry_times [int] + +`retry_times` ODPS REST 客户端在瞬态失败时的最大重试次数。默认 `4`。 + +### tunnel_connect_timeout_ms [long] + +`tunnel_connect_timeout_ms` Tunnel 客户端(执行批量数据上传/下载)的 HTTP 连接超时,单位毫秒。默认 `180000`(180 秒)。 + +### tunnel_read_timeout_ms [long] + +`tunnel_read_timeout_ms` Tunnel 客户端(批量数据上传/下载)的 HTTP 读超时,单位毫秒。默认 `300000`(300 秒)。当单次写请求超过 5 分钟(上传大分区或大批量 upsert)时调大。 + +### tunnel_retry_times [int] + +`tunnel_retry_times` Tunnel 客户端在瞬态失败时的最大重试次数。默认 `4`。 + ### insert_strategy [string] 如果将 `insert_strategy` 设置为 `upload`,插入操作将使用 upload 会话。 diff --git a/docs/zh/connectors/source/Maxcompute.md b/docs/zh/connectors/source/Maxcompute.md index 46e24a3f9f..1336f50b42 100644 --- a/docs/zh/connectors/source/Maxcompute.md +++ b/docs/zh/connectors/source/Maxcompute.md @@ -39,6 +39,12 @@ import ChangeLog from '../changelog/connector-maxcompute.md'; | table_list | Array | 否 | - | 要读取的表列表;可替代 `table_name` 一次读取多张 MaxCompute 表。 | | tunnel_endpoint | string | 否 | - | MaxCompute Tunnel 服务的自定义端点;未配置时根据区域自动推断。 | | tunnel_name | string | 否 | - | Tunnel Quota 名称;需同时将 `endpoint` 与 `tunnel_endpoint` 配置为 VPC 端点。 | +| connect_timeout_ms | long | 否 | 10000 | ODPS REST 客户端(元数据/catalog 调用)的 HTTP 连接超时,单位毫秒。默认 10000(10 秒)。 | +| read_timeout_ms | long | 否 | 120000 | ODPS REST 客户端(元数据/catalog 调用)的 HTTP 读超时,单位毫秒。默认 120000(120 秒)。 | +| retry_times | int | 否 | 4 | ODPS REST 客户端最大重试次数。默认 4。 | +| tunnel_connect_timeout_ms | long | 否 | 180000 | Tunnel 客户端(批量数据下载)的 HTTP 连接超时,单位毫秒。默认 180000(180 秒)。 | +| tunnel_read_timeout_ms | long | 否 | 300000 | Tunnel 客户端(批量数据下载)的 HTTP 读超时,单位毫秒。默认 300000(300 秒)。 | +| tunnel_retry_times | int | 否 | 4 | Tunnel 客户端最大重试次数。默认 4。 | | schema | config | 否 | - | 源表结构;未配置 `read_columns` 时按 schema 中字段读取。 | | common-options | | 否 | - | Source 插件通用参数,例如 `plugin_output`。 | @@ -132,6 +138,36 @@ Tunnel Quota 允许您使用专用的计算资源进行 MaxCompute Tunnel 数据 默认值:未设置(使用默认 quota) +> **客户端超时与重试** +> MaxCompute 有两个 HTTP 客户端。**ODPS REST 客户端**负责控制面(表/schema 查询、catalog 列表), +> 用 `connect_timeout_ms`、`read_timeout_ms`、`retry_times` 调整;**Tunnel 客户端**负责数据面(批量行下载), +> 用 `tunnel_*` 系列调整。单独设置 REST 参数**不会**改变 Tunnel 客户端的超时。 +> 毫秒值会被转换为整秒,最小值为 `1000`。 + +### connect_timeout_ms [long] + +`connect_timeout_ms` MaxCompute ODPS REST 客户端的 HTTP 连接超时,该客户端处理元数据与 catalog 调用(表/schema 查询、表列表)。单位毫秒。默认 `10000`(10 秒)。 + +### read_timeout_ms [long] + +`read_timeout_ms` ODPS REST 客户端(元数据/catalog 调用)的 HTTP 读超时,单位毫秒。默认 `120000`(120 秒)。当项目下表非常多或 schema 极宽导致拉取超时时调大。 + +### retry_times [int] + +`retry_times` ODPS REST 客户端在瞬态失败时的最大重试次数。默认 `4`。 + +### tunnel_connect_timeout_ms [long] + +`tunnel_connect_timeout_ms` Tunnel 客户端(执行批量数据下载)的 HTTP 连接超时,单位毫秒。默认 `180000`(180 秒)。 + +### tunnel_read_timeout_ms [long] + +`tunnel_read_timeout_ms` Tunnel 客户端(批量数据下载)的 HTTP 读超时,单位毫秒。默认 `300000`(300 秒)。当单次读请求超过 5 分钟(大表或大分区)时调大。 + +### tunnel_retry_times [int] + +`tunnel_retry_times` Tunnel 客户端在瞬态失败时的最大重试次数。默认 `4`。 + ### common options 源插件常用参数, 详见 [源通用选项](../common-options/source-common-options.md) . diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalog.java b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalog.java index 9c7fc5c7d9..04fb30d2a3 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalog.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalog.java @@ -17,6 +17,7 @@ package org.apache.seatunnel.connectors.seatunnel.maxcompute.catalog; +import org.apache.seatunnel.shade.com.google.common.annotations.VisibleForTesting; import org.apache.seatunnel.shade.com.google.common.collect.Lists; import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils; @@ -323,19 +324,23 @@ public class MaxComputeCatalog implements Catalog { /** * Creates an ODPS client for the given project. When {@code schemaName} is non-blank, sets it * as the current schema so that subsequent ODPS API calls resolve tables within that MaxCompute - * Schema namespace. + * Schema namespace. REST client timeout/retry options are applied via {@link + * MaxcomputeUtil#applyRestClientOptions} so that metadata and DDL calls honor + * connect_timeout_ms / read_timeout_ms / retry_times. */ private Odps getOdps(String project) { return getOdps(project, null); } - private Odps getOdps(String project, String schemaName) { + @VisibleForTesting + Odps getOdps(String project, String schemaName) { Odps odps = new Odps(account); odps.setEndpoint(readonlyConfig.get(MaxcomputeBaseOptions.ENDPOINT)); odps.setDefaultProject(project); if (StringUtils.isNotEmpty(schemaName)) { odps.setCurrentSchema(schemaName); } + MaxcomputeUtil.applyRestClientOptions(odps, readonlyConfig); return odps; } diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalogFactory.java b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalogFactory.java index b40c1d4ae0..7acdb43af7 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalogFactory.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalogFactory.java @@ -54,6 +54,9 @@ public class MaxComputeCatalogFactory implements CatalogFactory { MaxcomputeBaseOptions.PARTITION_SPEC, MaxcomputeBaseOptions.SPLIT_ROW, MaxcomputeBaseOptions.SCHEMA_NAME, + MaxcomputeBaseOptions.CONNECT_TIMEOUT_MS, + MaxcomputeBaseOptions.READ_TIMEOUT_MS, + MaxcomputeBaseOptions.RETRY_TIMES, ConnectorCommonOptions.SCHEMA) .build(); } diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/config/MaxcomputeBaseOptions.java b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/config/MaxcomputeBaseOptions.java index 8581b5542d..5c3ebbed8d 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/config/MaxcomputeBaseOptions.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/config/MaxcomputeBaseOptions.java @@ -93,4 +93,56 @@ public class MaxcomputeBaseOptions implements Serializable { .stringType() .noDefaultValue() .withDescription("Tunnel quota name for exclusive resource groups"); + + // ---- Odps REST client (control plane: metadata / schema / catalog) ---- + public static final Option<Long> CONNECT_TIMEOUT_MS = + Options.key("connect_timeout_ms") + .longType() + .defaultValue(10000L) + .withDescription( + "HTTP connect timeout for the MaxCompute (ODPS) REST client " + + "(metadata/catalog calls) in milliseconds. " + + "Millisecond values are converted to whole seconds; " + + "minimum 1000. Default 10000 (10s)."); + public static final Option<Long> READ_TIMEOUT_MS = + Options.key("read_timeout_ms") + .longType() + .defaultValue(120000L) + .withDescription( + "HTTP read timeout for the MaxCompute (ODPS) REST client " + + "(metadata/catalog calls) in milliseconds. " + + "Millisecond values are converted to whole seconds; " + + "minimum 1000. Default 120000 (120s)."); + public static final Option<Integer> RETRY_TIMES = + Options.key("retry_times") + .intType() + .defaultValue(4) + .withDescription( + "Max retry times for the MaxCompute (ODPS) REST client. Default 4."); + + // ---- Tunnel client (data plane: bulk row read / write / upsert) ---- + public static final Option<Long> TUNNEL_CONNECT_TIMEOUT_MS = + Options.key("tunnel_connect_timeout_ms") + .longType() + .defaultValue(180000L) + .withDescription( + "HTTP connect timeout for the MaxCompute Tunnel client " + + "(data upload/download) in milliseconds. " + + "Millisecond values are converted to whole seconds; " + + "minimum 1000. Default 180000 (180s)."); + public static final Option<Long> TUNNEL_READ_TIMEOUT_MS = + Options.key("tunnel_read_timeout_ms") + .longType() + .defaultValue(300000L) + .withDescription( + "HTTP read timeout for the MaxCompute Tunnel client " + + "(data upload/download) in milliseconds. " + + "Millisecond values are converted to whole seconds; " + + "minimum 1000. Default 300000 (300s)."); + public static final Option<Integer> TUNNEL_RETRY_TIMES = + Options.key("tunnel_retry_times") + .intType() + .defaultValue(4) + .withDescription( + "Max retry times for the MaxCompute Tunnel client. Default 4."); } diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/sink/MaxcomputeSinkFactory.java b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/sink/MaxcomputeSinkFactory.java index 89b73cab16..96539e6f77 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/sink/MaxcomputeSinkFactory.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/sink/MaxcomputeSinkFactory.java @@ -58,6 +58,12 @@ public class MaxcomputeSinkFactory implements TableSinkFactory { FormatOptions.DATETIME_FORMAT, MaxcomputeSinkOptions.TUNNEL_ENDPOINT, MaxcomputeSinkOptions.TUNNEL_NAME, + MaxcomputeSinkOptions.CONNECT_TIMEOUT_MS, + MaxcomputeSinkOptions.READ_TIMEOUT_MS, + MaxcomputeSinkOptions.RETRY_TIMES, + MaxcomputeSinkOptions.TUNNEL_CONNECT_TIMEOUT_MS, + MaxcomputeSinkOptions.TUNNEL_READ_TIMEOUT_MS, + MaxcomputeSinkOptions.TUNNEL_RETRY_TIMES, SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA) .build(); } diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/source/MaxcomputeSourceFactory.java b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/source/MaxcomputeSourceFactory.java index bcdfb085a5..8859a9b42f 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/source/MaxcomputeSourceFactory.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/source/MaxcomputeSourceFactory.java @@ -54,7 +54,13 @@ public class MaxcomputeSourceFactory implements TableSourceFactory { MaxcomputeSourceOptions.PROJECT, MaxcomputeSourceOptions.READ_COLUMNS, MaxcomputeSourceOptions.TUNNEL_ENDPOINT, - MaxcomputeSourceOptions.TUNNEL_NAME) + MaxcomputeSourceOptions.TUNNEL_NAME, + MaxcomputeSourceOptions.CONNECT_TIMEOUT_MS, + MaxcomputeSourceOptions.READ_TIMEOUT_MS, + MaxcomputeSourceOptions.RETRY_TIMES, + MaxcomputeSourceOptions.TUNNEL_CONNECT_TIMEOUT_MS, + MaxcomputeSourceOptions.TUNNEL_READ_TIMEOUT_MS, + MaxcomputeSourceOptions.TUNNEL_RETRY_TIMES) .exclusive(CatalogOptions.TABLE_LIST, MaxcomputeSourceOptions.TABLE_NAME) .build(); } diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtil.java b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtil.java index 96fcc4b319..b3918da9bf 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtil.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/main/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtil.java @@ -33,10 +33,14 @@ import com.aliyun.odps.account.Account; import com.aliyun.odps.account.AklessAccount; import com.aliyun.odps.account.AliyunAccount; import com.aliyun.odps.account.StsAccount; +import com.aliyun.odps.rest.RestClient; +import com.aliyun.odps.tunnel.Configuration; import com.aliyun.odps.tunnel.TableTunnel; import com.aliyun.odps.tunnel.TunnelException; import lombok.extern.slf4j.Slf4j; +import static org.apache.seatunnel.shade.com.google.common.base.Preconditions.checkArgument; + @Slf4j public class MaxcomputeUtil { public static Table getTable(ReadonlyConfig readonlyConfig) { @@ -47,6 +51,22 @@ public class MaxcomputeUtil { public static TableTunnel getTableTunnel(ReadonlyConfig readonlyConfig) { Odps odps = getOdps(readonlyConfig); TableTunnel tableTunnel = new TableTunnel(odps); + // Tunnel client timeouts / retry (data plane). The tunnel RestClient is built + // lazily by Configuration.newRestClient() when a session is created, reading + // these socket fields, so setting them before any session is built takes effect. + Configuration tunnelConfig = tableTunnel.getConfig(); + tunnelConfig.setSocketConnectTimeout( + toTimeoutSeconds( + readonlyConfig.get(MaxcomputeBaseOptions.TUNNEL_CONNECT_TIMEOUT_MS), + "tunnel_connect_timeout_ms")); + tunnelConfig.setSocketTimeout( + toTimeoutSeconds( + readonlyConfig.get(MaxcomputeBaseOptions.TUNNEL_READ_TIMEOUT_MS), + "tunnel_read_timeout_ms")); + tunnelConfig.setSocketRetryTimes( + toRetryTimes( + readonlyConfig.get(MaxcomputeBaseOptions.TUNNEL_RETRY_TIMES), + "tunnel_retry_times")); if (StringUtils.isNotEmpty(readonlyConfig.get(MaxcomputeBaseOptions.TUNNEL_ENDPOINT))) { tableTunnel.setEndpoint(readonlyConfig.get(MaxcomputeBaseOptions.TUNNEL_ENDPOINT)); } @@ -84,9 +104,54 @@ public class MaxcomputeUtil { odps.setDefaultProject(readonlyConfig.get(MaxcomputeBaseOptions.PROJECT)); odps.setCurrentSchema( readonlyConfig.getOptional(MaxcomputeBaseOptions.SCHEMA_NAME).orElse(null)); + applyRestClientOptions(odps, readonlyConfig); return odps; } + /** + * Applies the ODPS REST client timeout/retry options (control plane) to an Odps instance. + * Shared by {@link #getOdps(ReadonlyConfig)} and {@code MaxComputeCatalog.getOdps} so that + * metadata/catalog calls honor {@code connect_timeout_ms} / {@code read_timeout_ms} / {@code + * retry_times}. Values default to the SDK defaults, so omitting them preserves existing + * behavior. + */ + public static void applyRestClientOptions(Odps odps, ReadonlyConfig readonlyConfig) { + // RestClient stores connect/read timeout in seconds internally. + RestClient restClient = odps.getRestClient(); + restClient.setConnectTimeout( + toTimeoutSeconds( + readonlyConfig.get(MaxcomputeBaseOptions.CONNECT_TIMEOUT_MS), + "connect_timeout_ms")); + restClient.setReadTimeout( + toTimeoutSeconds( + readonlyConfig.get(MaxcomputeBaseOptions.READ_TIMEOUT_MS), + "read_timeout_ms")); + restClient.setRetryTimes( + toRetryTimes(readonlyConfig.get(MaxcomputeBaseOptions.RETRY_TIMES), "retry_times")); + } + + /** + * Converts a millisecond timeout to whole seconds, rejecting sub-second / negative values and + * values that would overflow the SDK's int-seconds field. The ODPS RestClient and Tunnel + * Configuration store timeouts in seconds internally, so millisecond input is divided by 1000; + * a value below 1000 would silently collapse to 0 or be clamped to 1 second, which is almost + * never what the user intended. + */ + private static int toTimeoutSeconds(long ms, String optionName) { + checkArgument( + ms >= 1000 && ms <= Integer.MAX_VALUE * 1000L, + "%s must be between 1000 and %d (millisecond values are converted to whole seconds). Got: %s", + optionName, + Integer.MAX_VALUE * 1000L, + ms); + return (int) (ms / 1000); + } + + private static int toRetryTimes(int value, String optionName) { + checkArgument(value >= 0, "%s must be non-negative. Got: %s", optionName, value); + return value; + } + public static TableTunnel.DownloadSession getDownloadSession(ReadonlyConfig readonlyConfig) { TableTunnel tunnel = getTableTunnel(readonlyConfig); TableTunnel.DownloadSession session; diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/MaxcomputeSourceFactoryTest.java b/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/MaxcomputeSourceFactoryTest.java index 926ba28834..ec0c5ea932 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/MaxcomputeSourceFactoryTest.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/MaxcomputeSourceFactoryTest.java @@ -17,16 +17,54 @@ package org.apache.seatunnel.connectors.seatunnel.maxcompute; +import org.apache.seatunnel.api.configuration.Option; +import org.apache.seatunnel.api.configuration.util.OptionRule; import org.apache.seatunnel.connectors.seatunnel.maxcompute.sink.MaxcomputeSinkFactory; import org.apache.seatunnel.connectors.seatunnel.maxcompute.source.MaxcomputeSourceFactory; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.util.HashSet; +import java.util.Set; + public class MaxcomputeSourceFactoryTest { @Test void optionRule() { Assertions.assertNotNull((new MaxcomputeSourceFactory()).optionRule()); Assertions.assertNotNull((new MaxcomputeSinkFactory()).optionRule()); } + + /** + * The six client timeout / retry options (3 for the ODPS REST client, 3 for the Tunnel client) + * must be declared as optional in both the source and sink factory option rules, otherwise + * users cannot override them from job configs. + */ + @Test + void optionRuleRegistersTimeoutAndRetryOptions() { + assertTimeoutAndRetryOptionsRegistered(new MaxcomputeSourceFactory().optionRule()); + assertTimeoutAndRetryOptionsRegistered(new MaxcomputeSinkFactory().optionRule()); + } + + private void assertTimeoutAndRetryOptionsRegistered(OptionRule rule) { + Set<String> optionalKeys = new HashSet<>(); + for (Option<?> option : rule.getOptionalOptions()) { + optionalKeys.add(option.key()); + } + // REST client (control plane) + Assertions.assertTrue( + optionalKeys.contains("connect_timeout_ms"), "connect_timeout_ms not registered"); + Assertions.assertTrue( + optionalKeys.contains("read_timeout_ms"), "read_timeout_ms not registered"); + Assertions.assertTrue(optionalKeys.contains("retry_times"), "retry_times not registered"); + // Tunnel client (data plane) + Assertions.assertTrue( + optionalKeys.contains("tunnel_connect_timeout_ms"), + "tunnel_connect_timeout_ms not registered"); + Assertions.assertTrue( + optionalKeys.contains("tunnel_read_timeout_ms"), + "tunnel_read_timeout_ms not registered"); + Assertions.assertTrue( + optionalKeys.contains("tunnel_retry_times"), "tunnel_retry_times not registered"); + } } diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalogTest.java b/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalogTest.java new file mode 100644 index 0000000000..8766d5ba16 --- /dev/null +++ b/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/catalog/MaxComputeCatalogTest.java @@ -0,0 +1,79 @@ +/* + * 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.maxcompute.catalog; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import com.aliyun.odps.Odps; + +import java.util.HashMap; +import java.util.Map; + +public class MaxComputeCatalogTest { + + /** Minimal config that lets the catalog build an Odps client without network calls. */ + private static Map<String, Object> baseConfig() { + Map<String, Object> config = new HashMap<>(); + config.put("accessId", "my-id"); + config.put("accesskey", "my-key"); + config.put("endpoint", "http://service.odps.aliyun.com/api"); + config.put("project", "my_project"); + return config; + } + + /** + * The REST client timeout/retry options must be applied to the Odps instance built by + * MaxComputeCatalog.getOdps, so that metadata and DDL calls honor connect_timeout_ms / + * read_timeout_ms / retry_times. This guards against the catalog path silently ignoring the + * options (the most common metadata path: source schema discovery and sink save-mode DDL). + */ + @Test + void testCatalogAppliesRestClientOptions() { + Map<String, Object> config = baseConfig(); + config.put("connect_timeout_ms", 30000L); + config.put("read_timeout_ms", 60000L); + config.put("retry_times", 7); + + MaxComputeCatalog catalog = new MaxComputeCatalog("test", ReadonlyConfig.fromMap(config)); + catalog.open(); + + Odps odps = catalog.getOdps("my_project", null); + Assertions.assertEquals(30, odps.getRestClient().getConnectTimeout()); + Assertions.assertEquals(60, odps.getRestClient().getReadTimeout()); + Assertions.assertEquals(7, odps.getRestClient().getRetryTimes()); + } + + /** + * When no REST timeout/retry options are supplied, the catalog must fall back to the option + * defaults (connect 10s, read 120s, retry 4) — i.e. the original SDK behavior is preserved. + */ + @Test + void testCatalogUsesRestClientDefaultsWhenOptionsAbsent() { + MaxComputeCatalog catalog = + new MaxComputeCatalog("test", ReadonlyConfig.fromMap(baseConfig())); + catalog.open(); + + Odps odps = catalog.getOdps("my_project", null); + Assertions.assertEquals(10, odps.getRestClient().getConnectTimeout()); + Assertions.assertEquals(120, odps.getRestClient().getReadTimeout()); + Assertions.assertEquals(4, odps.getRestClient().getRetryTimes()); + } +} diff --git a/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtilTest.java b/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtilTest.java index e3aa017bdd..d35c6943bf 100644 --- a/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtilTest.java +++ b/seatunnel-connectors-v2/connector-maxcompute/src/test/java/org/apache/seatunnel/connectors/seatunnel/maxcompute/util/MaxcomputeUtilTest.java @@ -18,14 +18,19 @@ package org.apache.seatunnel.connectors.seatunnel.maxcompute.util; import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.maxcompute.config.MaxcomputeBaseOptions; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import com.aliyun.odps.Odps; import com.aliyun.odps.account.Account; import com.aliyun.odps.account.AklessAccount; import com.aliyun.odps.account.AliyunAccount; import com.aliyun.odps.account.StsAccount; +import com.aliyun.odps.commons.GeneralConfiguration; +import com.aliyun.odps.rest.RestClient; +import com.aliyun.odps.tunnel.TableTunnel; import java.util.HashMap; import java.util.Map; @@ -133,7 +138,7 @@ public class MaxcomputeUtilTest { config.put("project", "my_project"); config.put("schema_name", "my_schema"); - com.aliyun.odps.Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config)); + Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config)); Assertions.assertEquals("my_schema", odps.getCurrentSchema()); } @@ -150,8 +155,248 @@ public class MaxcomputeUtilTest { config.put("endpoint", "http://service.odps.aliyun.com/api"); config.put("project", "my_project"); - com.aliyun.odps.Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config)); + Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config)); Assertions.assertNull(odps.getCurrentSchema()); } + + // --- getOdps() / getTableTunnel() timeout & retry wiring tests --- + // + // The ODPS SDK exposes two independent HTTP clients: + // * REST client (Odps.getRestClient()) -> control plane: metadata / catalog calls + // * Tunnel client (TableTunnel.getConfig()) -> data plane: bulk row read / write / upsert + // Each is configured from its own set of options, so the tests below cover per-option + // application, the option defaults that flow through when options are omitted, the + // milliseconds->seconds validation guard, and that the two clients never cross-contaminate. + + /** Minimal config that lets getOdps()/getTableTunnel() build a client without network calls. */ + private static Map<String, Object> baseConfig() { + Map<String, Object> config = new HashMap<>(); + config.put("accessId", "my-id"); + config.put("accesskey", "my-key"); + config.put("endpoint", "http://service.odps.aliyun.com/api"); + config.put("project", "my_project"); + return config; + } + + /** + * A user-supplied connect_timeout_ms must be applied to the ODPS REST client (converted to + * seconds, since the SDK stores connect/read timeout in seconds internally). This verifies the + * control-plane client timeout override wiring in getOdps(). + */ + @Test + void testGetOdpsAppliesRestClientConnectTimeout() { + Map<String, Object> config = baseConfig(); + config.put("connect_timeout_ms", 30000L); + + Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config)); + + Assertions.assertEquals(30, odps.getRestClient().getConnectTimeout()); + } + + /** read_timeout_ms must be applied to the REST client's read timeout (seconds). */ + @Test + void testGetOdpsAppliesRestClientReadTimeout() { + Map<String, Object> config = baseConfig(); + config.put("read_timeout_ms", 60000L); + + Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config)); + + Assertions.assertEquals(60, odps.getRestClient().getReadTimeout()); + } + + /** retry_times must be applied to the REST client's retry count. */ + @Test + void testGetOdpsAppliesRestClientRetryTimes() { + Map<String, Object> config = baseConfig(); + config.put("retry_times", 7); + + Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config)); + + Assertions.assertEquals(7, odps.getRestClient().getRetryTimes()); + } + + /** + * When no REST timeout/retry options are supplied, the connector must fall back to the option + * defaults (connect 10s, read 120s, retry 4) — i.e. the original SDK behavior is preserved. + * This guards against accidental regressions in backward compatibility. + */ + @Test + void testGetOdpsUsesRestClientDefaultsWhenTimeoutOptionsAbsent() { + Odps odps = MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(baseConfig())); + + Assertions.assertEquals(10, odps.getRestClient().getConnectTimeout()); + Assertions.assertEquals(120, odps.getRestClient().getReadTimeout()); + Assertions.assertEquals(4, odps.getRestClient().getRetryTimes()); + } + + /** tunnel_connect_timeout_ms must be applied to the Tunnel client's socket connect timeout. */ + @Test + void testGetTableTunnelAppliesSocketConnectTimeout() { + Map<String, Object> config = baseConfig(); + config.put("tunnel_connect_timeout_ms", 240000L); + + TableTunnel tableTunnel = MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(config)); + + Assertions.assertEquals(240, tableTunnel.getConfig().getSocketConnectTimeout()); + } + + /** + * A user-supplied tunnel_read_timeout_ms must be applied to the Tunnel client's socket read + * timeout (converted to seconds), independent of the ODPS REST client. This verifies the + * data-plane client timeout override wiring in getTableTunnel(). + */ + @Test + void testGetTableTunnelAppliesSocketReadTimeout() { + Map<String, Object> config = baseConfig(); + config.put("tunnel_read_timeout_ms", 600000L); + + TableTunnel tableTunnel = MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(config)); + + Assertions.assertEquals(600, tableTunnel.getConfig().getSocketTimeout()); + } + + /** tunnel_retry_times must be applied to the Tunnel client's socket retry count. */ + @Test + void testGetTableTunnelAppliesSocketRetryTimes() { + Map<String, Object> config = baseConfig(); + config.put("tunnel_retry_times", 8); + + TableTunnel tableTunnel = MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(config)); + + Assertions.assertEquals(8, tableTunnel.getConfig().getSocketRetryTimes()); + } + + /** + * When no Tunnel timeout/retry options are supplied, the connector must fall back to the option + * defaults (connect 180s, read 300s, retry 4) — i.e. the original SDK behavior is preserved. + */ + @Test + void testGetTableTunnelUsesSocketDefaultsWhenTimeoutOptionsAbsent() { + TableTunnel tableTunnel = + MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(baseConfig())); + + Assertions.assertEquals(180, tableTunnel.getConfig().getSocketConnectTimeout()); + Assertions.assertEquals(300, tableTunnel.getConfig().getSocketTimeout()); + Assertions.assertEquals(4, tableTunnel.getConfig().getSocketRetryTimes()); + } + + /** + * REST client options must not leak onto the Tunnel client and vice versa: the two clients are + * configured from disjoint option sets. Setting connect_timeout_ms (REST) must not change the + * Tunnel socket connect timeout, and setting tunnel_read_timeout_ms must not change the REST + * read timeout. + */ + @Test + void testRestClientAndTunnelClientTimeoutsAreIndependent() { + Map<String, Object> config = baseConfig(); + config.put("connect_timeout_ms", 30000L); // REST-only option + config.put("tunnel_read_timeout_ms", 600000L); // Tunnel-only option + + TableTunnel tableTunnel = MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(config)); + Odps odps = tableTunnel.getConfig().getOdps(); + + // REST client picks up the REST option but is untouched by the Tunnel option. + Assertions.assertEquals(30, odps.getRestClient().getConnectTimeout()); + Assertions.assertEquals(120, odps.getRestClient().getReadTimeout()); + // Tunnel client picks up the Tunnel option but is untouched by the REST option. + Assertions.assertEquals(180, tableTunnel.getConfig().getSocketConnectTimeout()); + Assertions.assertEquals(600, tableTunnel.getConfig().getSocketTimeout()); + } + + // --- validation tests --- + // + // Sub-second / negative timeout values and negative retry counts must be rejected at + // config-application time rather than silently clamped or passed through to the SDK, + // because a 0 turning into a 1-second timeout is a nasty surprise that is hard to diagnose. + + /** + * A sub-second timeout value (e.g. 500ms) must be rejected with a clear error, not silently + * clamped to 1 second, because the SDK stores timeouts in whole seconds and the user almost + * certainly did not intend a 1-second timeout. + */ + @Test + void testSubSecondTimeoutIsRejected() { + Map<String, Object> config = baseConfig(); + config.put("connect_timeout_ms", 500L); + config.put("tunnel_connect_timeout_ms", 999L); + + Assertions.assertThrows( + IllegalArgumentException.class, + () -> MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config))); + Assertions.assertThrows( + IllegalArgumentException.class, + () -> MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(config))); + } + + /** A negative retry_times must be rejected, not passed through to the SDK. */ + @Test + void testNegativeRetryTimesIsRejected() { + Map<String, Object> config = baseConfig(); + config.put("retry_times", -1); + + Assertions.assertThrows( + IllegalArgumentException.class, + () -> MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config))); + } + + /** A negative tunnel_retry_times must be rejected, not passed through to the SDK. */ + @Test + void testNegativeTunnelRetryTimesIsRejected() { + Map<String, Object> config = baseConfig(); + config.put("tunnel_retry_times", -1); + + Assertions.assertThrows( + IllegalArgumentException.class, + () -> MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(config))); + } + + /** + * A millisecond value that would overflow the SDK's int-seconds field (above ~2.1 trillion ms) + * must be rejected with a clear error naming the option, not a bare ArithmeticException. + */ + @Test + void testOverflowTimeoutIsRejected() { + Map<String, Object> config = baseConfig(); + config.put("connect_timeout_ms", Long.MAX_VALUE); + config.put("tunnel_read_timeout_ms", ((long) Integer.MAX_VALUE) * 1000L + 1); + + Assertions.assertThrows( + IllegalArgumentException.class, + () -> MaxcomputeUtil.getOdps(ReadonlyConfig.fromMap(config))); + Assertions.assertThrows( + IllegalArgumentException.class, + () -> MaxcomputeUtil.getTableTunnel(ReadonlyConfig.fromMap(config))); + } + + // --- SDK defaults consistency test --- + + /** + * The option defaults must match the pinned ODPS SDK (0.51.2) defaults so that omitting the + * options preserves the original SDK behavior. If a future SDK bump changes its defaults, this + * test fails loudly and the option defaults can be revisited. + */ + @Test + void testOptionDefaultsMatchSdkDefaults() { + // REST client (control plane) + Assertions.assertEquals( + RestClient.DEFAULT_CONNECT_TIMEOUT, + MaxcomputeBaseOptions.CONNECT_TIMEOUT_MS.defaultValue() / 1000); + Assertions.assertEquals( + RestClient.DEFAULT_READ_TIMEOUT, + MaxcomputeBaseOptions.READ_TIMEOUT_MS.defaultValue() / 1000); + Assertions.assertEquals( + RestClient.DEFAULT_CONNECT_RETRYTIMES, + MaxcomputeBaseOptions.RETRY_TIMES.defaultValue()); + // Tunnel client (data plane) + Assertions.assertEquals( + GeneralConfiguration.DEFAULT_SOCKET_CONNECT_TIMEOUT, + MaxcomputeBaseOptions.TUNNEL_CONNECT_TIMEOUT_MS.defaultValue() / 1000); + Assertions.assertEquals( + GeneralConfiguration.DEFAULT_SOCKET_TIMEOUT, + MaxcomputeBaseOptions.TUNNEL_READ_TIMEOUT_MS.defaultValue() / 1000); + Assertions.assertEquals( + GeneralConfiguration.DEFAULT_SOCKET_RETRY_TIMES, + MaxcomputeBaseOptions.TUNNEL_RETRY_TIMES.defaultValue()); + } } diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-maxcompute-e2e/src/test/resources/maxcompute_to_maxcompute.conf b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-maxcompute-e2e/src/test/resources/maxcompute_to_maxcompute.conf index 8b5ad5bd3a..f64d595027 100644 --- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-maxcompute-e2e/src/test/resources/maxcompute_to_maxcompute.conf +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-maxcompute-e2e/src/test/resources/maxcompute_to_maxcompute.conf @@ -71,5 +71,14 @@ sink { project = "mocked_mc" table_name = "test_table_sink" insert_strategy = "upsert" + + # Non-default but safe values to prove the options are applied in a real job, + # not just accepted by OptionRule validation. + connect_timeout_ms = 15000 + read_timeout_ms = 121000 + retry_times = 5 + tunnel_connect_timeout_ms = 181000 + tunnel_read_timeout_ms = 301000 + tunnel_retry_times = 5 } } \ No newline at end of file
