This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 5c9ac457a4d branch-4.1: [feature](catalog) Support Alibaba Cloud OSS
Tables REST catalog (#66823)
5c9ac457a4d is described below
commit 5c9ac457a4d35e69869f5e00af5e4c8ea00ff17c
Author: Chenjunwei <[email protected]>
AuthorDate: Mon Aug 24 09:46:20 2026 +0800
branch-4.1: [feature](catalog) Support Alibaba Cloud OSS Tables REST
catalog (#66823)
### What problem does this PR solve?
Issue Number: None
Related PR: #66567
Problem Summary:
Backport Alibaba Cloud OSS Tables REST Catalog support from #66567 to
branch-4.1.
The branch-4.1 Iceberg implementation predates the connector-plugin
refactor, so this backport applies the same behavior to
`IcebergRestProperties`:
- recognize `osstables` as a managed SigV4 signing name;
- require a signing region and `sigv4-enabled=true`;
- reuse `oss.access_key`, `oss.secret_key`, and `oss.session_token` for
REST signing;
- keep the OSS S3-compatible endpoint and credentials for S3FileIO.
### Release note
Support Alibaba Cloud OSS Tables through the Iceberg REST Catalog on
branch-4.1.
### Check List (For Author)
- Test
- [x] Unit Test
- `./run-fe-ut.sh --run
org.apache.doris.datasource.property.metastore.IcebergRestPropertiesTest`
- Behavior changed:
- [x] Yes. Iceberg REST catalogs with `signing-name=osstables` reuse OSS
credentials and validate required SigV4 settings.
- Does this need documentation?
- [x] Yes. apache/doris-website#4054
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label
---
.../apache/doris/datasource/ExternalCatalog.java | 1 +
.../CatalogConnectivityTestCoordinator.java | 6 +-
.../IcebergRestConnectivityTester.java | 14 ++--
.../doris/datasource/iceberg/IcebergUtils.java | 54 +++++++++++++++
.../datasource/iceberg/source/IcebergScanNode.java | 1 +
.../metastore/AbstractIcebergProperties.java | 21 ++----
.../property/metastore/IcebergRestProperties.java | 33 +++++++--
.../org/apache/doris/planner/IcebergTableSink.java | 18 ++---
.../IcebergRestConnectivityTesterTest.java | 57 ++++++++++++++++
.../doris/datasource/iceberg/IcebergUtilsTest.java | 30 ++++++++
.../metastore/IcebergRestPropertiesTest.java | 79 ++++++++++++++++++++++
.../property/storage/StoragePropertiesTest.java | 18 +++++
.../apache/doris/planner/IcebergTableSinkTest.java | 48 +++++++++++++
13 files changed, 345 insertions(+), 35 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/ExternalCatalog.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/ExternalCatalog.java
index 481c1fcb3ac..d72dadb554f 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/datasource/ExternalCatalog.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/ExternalCatalog.java
@@ -304,6 +304,7 @@ public abstract class ExternalCatalog
CatalogConnectivityTestCoordinator testCoordinator = new
CatalogConnectivityTestCoordinator(
name,
catalogProperty.getMetastoreProperties(),
+ catalogProperty.getOrderedStoragePropertiesList(),
catalogProperty.getStoragePropertiesMap()
);
testCoordinator.runTests();
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/CatalogConnectivityTestCoordinator.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/CatalogConnectivityTestCoordinator.java
index cf8c308849a..9c3baf45b99 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/CatalogConnectivityTestCoordinator.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/CatalogConnectivityTestCoordinator.java
@@ -35,6 +35,7 @@ import org.apache.commons.lang3.StringUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
+import java.util.List;
import java.util.Map;
/**
@@ -47,6 +48,7 @@ public class CatalogConnectivityTestCoordinator {
private final String catalogName;
private final MetastoreProperties metastoreProperties;
+ private final List<StorageProperties> storagePropertiesList;
private final Map<StorageProperties.Type, StorageProperties>
storagePropertiesMap;
private String warehouseLocation;
@@ -54,9 +56,11 @@ public class CatalogConnectivityTestCoordinator {
public CatalogConnectivityTestCoordinator(
String catalogName,
MetastoreProperties metastoreProperties,
+ List<StorageProperties> storagePropertiesList,
Map<StorageProperties.Type, StorageProperties>
storagePropertiesMap) {
this.catalogName = catalogName;
this.metastoreProperties = metastoreProperties;
+ this.storagePropertiesList = storagePropertiesList;
this.storagePropertiesMap = storagePropertiesMap;
}
@@ -295,7 +299,7 @@ public class CatalogConnectivityTestCoordinator {
// Iceberg REST
if (props instanceof IcebergRestProperties) {
- return new IcebergRestConnectivityTester((IcebergRestProperties)
props);
+ return new IcebergRestConnectivityTester((IcebergRestProperties)
props, storagePropertiesList);
}
// Iceberg S3Table
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/IcebergRestConnectivityTester.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/IcebergRestConnectivityTester.java
index def265fea82..630c6c87073 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/IcebergRestConnectivityTester.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/connectivity/IcebergRestConnectivityTester.java
@@ -18,11 +18,12 @@
package org.apache.doris.datasource.connectivity;
import
org.apache.doris.datasource.property.metastore.AbstractIcebergProperties;
-import org.apache.doris.datasource.property.metastore.IcebergRestProperties;
+import org.apache.doris.datasource.property.storage.StorageProperties;
import org.apache.iceberg.CatalogProperties;
import org.apache.iceberg.rest.RESTCatalog;
+import java.util.List;
import java.util.Map;
public class IcebergRestConnectivityTester extends
AbstractIcebergConnectivityTester {
@@ -30,9 +31,12 @@ public class IcebergRestConnectivityTester extends
AbstractIcebergConnectivityTe
private static final String DEFAULT_BASE_LOCATION =
"default-base-location";
private String warehouseLocation;
+ private final List<StorageProperties> storagePropertiesList;
- public IcebergRestConnectivityTester(AbstractIcebergProperties properties)
{
+ public IcebergRestConnectivityTester(AbstractIcebergProperties properties,
+ List<StorageProperties> storagePropertiesList) {
super(properties);
+ this.storagePropertiesList = storagePropertiesList;
}
@Override
@@ -48,10 +52,8 @@ public class IcebergRestConnectivityTester extends
AbstractIcebergConnectivityTe
@Override
public void testConnection() throws Exception {
- Map<String, String> restProps = ((IcebergRestProperties)
properties).getIcebergRestCatalogProperties();
-
- try (RESTCatalog catalog = new RESTCatalog()) {
- catalog.initialize("connectivity-test", restProps);
+ try (RESTCatalog catalog = (RESTCatalog) properties.initializeCatalog(
+ "connectivity-test", storagePropertiesList)) {
// Validate connection by listing namespaces.
// This verifies authentication and warehouse configuration.
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
index 0b375c70d67..891183adbd2 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java
@@ -59,6 +59,9 @@ import org.apache.doris.datasource.metacache.CacheSpec;
import org.apache.doris.datasource.mvcc.MvccSnapshot;
import org.apache.doris.datasource.mvcc.MvccUtil;
import org.apache.doris.datasource.property.metastore.HMSBaseProperties;
+import
org.apache.doris.datasource.property.storage.AbstractS3CompatibleProperties;
+import org.apache.doris.datasource.property.storage.S3Properties;
+import org.apache.doris.datasource.property.storage.StorageProperties;
import org.apache.doris.nereids.exceptions.NotSupportedException;
import org.apache.doris.nereids.trees.expressions.literal.Result;
import org.apache.doris.nereids.types.VarBinaryType;
@@ -145,6 +148,7 @@ import java.time.temporal.TemporalAccessor;
import java.util.ArrayList;
import java.util.Base64;
import java.util.Comparator;
+import java.util.EnumMap;
import java.util.HashMap;
import java.util.LinkedHashSet;
import java.util.List;
@@ -177,6 +181,56 @@ public class IcebergUtils {
public static final String TOTAL_POSITION_DELETES =
"total-position-deletes";
public static final String TOTAL_EQUALITY_DELETES =
"total-equality-deletes";
+ /**
+ * Selects the storage bindings Iceberg should consume together. Iceberg
can configure only one
+ * S3-compatible data plane, so a concrete provider such as OSS takes
precedence over the generic
+ * S3 fallback while unrelated storage bindings are preserved.
+ */
+ public static List<StorageProperties> selectEffectiveStorageProperties(
+ List<StorageProperties> storagePropertiesList) {
+ StorageProperties chosenS3 =
chooseS3CompatibleStorage(storagePropertiesList);
+ List<StorageProperties> selected = new ArrayList<>();
+ for (StorageProperties storageProperties : storagePropertiesList) {
+ if (!(storageProperties instanceof AbstractS3CompatibleProperties)
+ || storageProperties == chosenS3) {
+ selected.add(storageProperties);
+ }
+ }
+ return selected;
+ }
+
+ public static Map<StorageProperties.Type, StorageProperties>
selectEffectiveStorageProperties(
+ Map<StorageProperties.Type, StorageProperties>
storagePropertiesMap) {
+ List<StorageProperties> ordered = new ArrayList<>();
+ for (StorageProperties.Type type : StorageProperties.Type.values()) {
+ StorageProperties storageProperties =
storagePropertiesMap.get(type);
+ if (storageProperties != null) {
+ ordered.add(storageProperties);
+ }
+ }
+
+ Map<StorageProperties.Type, StorageProperties> selected = new
EnumMap<>(StorageProperties.Type.class);
+ for (StorageProperties storageProperties :
selectEffectiveStorageProperties(ordered)) {
+ selected.put(storageProperties.getType(), storageProperties);
+ }
+ return selected;
+ }
+
+ private static StorageProperties
chooseS3CompatibleStorage(List<StorageProperties> storagePropertiesList) {
+ StorageProperties fallback = null;
+ for (StorageProperties storageProperties : storagePropertiesList) {
+ if (storageProperties instanceof AbstractS3CompatibleProperties) {
+ if (fallback == null) {
+ fallback = storageProperties;
+ }
+ if (!(storageProperties instanceof S3Properties)) {
+ return storageProperties;
+ }
+ }
+ }
+ return fallback;
+ }
+
// nickname in flink and spark
public static final String WRITE_FORMAT = "write-format";
public static final String COMPRESSION_CODEC = "compression-codec";
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
index 3e2d560f60e..561d2557364 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
@@ -284,6 +284,7 @@ public class IcebergScanNode extends FileQueryScanNode {
source.getCatalog().getCatalogProperty().getStoragePropertiesMap(),
icebergTable
);
+ storagePropertiesMap =
IcebergUtils.selectEffectiveStorageProperties(storagePropertiesMap);
backendStorageProperties =
CredentialUtils.getBackendPropertiesFromStorageMap(storagePropertiesMap);
} finally {
if (getSummaryProfile() != null) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractIcebergProperties.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractIcebergProperties.java
index 9a3a5ef5d2a..ef225fae44c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractIcebergProperties.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractIcebergProperties.java
@@ -19,6 +19,7 @@ package org.apache.doris.datasource.property.metastore;
import org.apache.doris.common.security.authentication.ExecutionAuthenticator;
import org.apache.doris.datasource.iceberg.IcebergExternalCatalog;
+import org.apache.doris.datasource.iceberg.IcebergUtils;
import org.apache.doris.datasource.metacache.CacheSpec;
import
org.apache.doris.datasource.property.common.IcebergAwsAssumeRoleProperties;
import
org.apache.doris.datasource.property.storage.AbstractS3CompatibleProperties;
@@ -210,28 +211,18 @@ public abstract class AbstractIcebergProperties extends
MetastoreProperties {
*/
public void toFileIOProperties(List<StorageProperties>
storagePropertiesList,
Map<String, String> fileIOProperties, Configuration conf) {
- // We only support one S3-compatible storage property for FileIO
configuration.
- // When multiple AbstractS3CompatibleProperties exist, prefer the
first non-S3Properties one,
- // because a non-S3 type (e.g. OSSProperties, COSProperties) indicates
the user has explicitly
- // specified a concrete S3-compatible storage, which should take
priority over the generic S3Properties.
- AbstractS3CompatibleProperties s3Fallback = null;
AbstractS3CompatibleProperties s3Target = null;
- for (StorageProperties storageProperties : storagePropertiesList) {
+ for (StorageProperties storageProperties :
+
IcebergUtils.selectEffectiveStorageProperties(storagePropertiesList)) {
if (conf != null && storageProperties.getHadoopStorageConfig() !=
null) {
conf.addResource(storageProperties.getHadoopStorageConfig());
}
if (storageProperties instanceof AbstractS3CompatibleProperties) {
- if (s3Fallback == null) {
- s3Fallback = (AbstractS3CompatibleProperties)
storageProperties;
- }
- if (s3Target == null && !(storageProperties instanceof
S3Properties)) {
- s3Target = (AbstractS3CompatibleProperties)
storageProperties;
- }
+ s3Target = (AbstractS3CompatibleProperties) storageProperties;
}
}
- AbstractS3CompatibleProperties chosen = s3Target != null ? s3Target :
s3Fallback;
- if (chosen != null) {
- toS3FileIOProperties(chosen, fileIOProperties);
+ if (s3Target != null) {
+ toS3FileIOProperties(s3Target, fileIOProperties);
} else {
String region =
AbstractS3CompatibleProperties.getRegionFromProperties(fileIOProperties);
if (!Strings.isNullOrEmpty(region)) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/IcebergRestProperties.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/IcebergRestProperties.java
index 3c3da3fce5f..25648fb3551 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/IcebergRestProperties.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/IcebergRestProperties.java
@@ -20,6 +20,7 @@ package org.apache.doris.datasource.property.metastore;
import org.apache.doris.datasource.iceberg.IcebergExternalCatalog;
import org.apache.doris.datasource.property.common.AwsCredentialsProviderMode;
import
org.apache.doris.datasource.property.common.IcebergAwsClientCredentialsProperties;
+import org.apache.doris.datasource.property.storage.OSSProperties;
import org.apache.doris.datasource.property.storage.S3Properties;
import org.apache.doris.datasource.property.storage.StorageProperties;
import org.apache.doris.foundation.property.ConnectorProperty;
@@ -49,6 +50,7 @@ public class IcebergRestProperties extends
AbstractIcebergProperties {
private Map<String, String> icebergRestCatalogProperties;
private S3Properties s3Properties;
+ private OSSProperties ossProperties;
@Getter
@ConnectorProperty(names = {"iceberg.rest.uri", "uri"},
@@ -214,7 +216,11 @@ public class IcebergRestProperties extends
AbstractIcebergProperties {
AwsCredentialsProviderMode.fromString(icebergRestCredentialsProviderType);
buildRules().validate();
if (shouldUseS3PropertiesForRestCredentials()) {
- s3Properties = S3Properties.of(origProps);
+ if (isOssTables()) {
+ ossProperties = OSSProperties.of(origProps);
+ } else {
+ s3Properties = S3Properties.of(origProps);
+ }
}
initIcebergRestCatalogProperties();
}
@@ -251,13 +257,19 @@ public class IcebergRestProperties extends
AbstractIcebergProperties {
}
}
- // When signing-name is glue or s3tables: require signing-region and
sigv4-enabled
+ // SigV4-backed REST catalogs require a signing region and SigV4 to be
enabled.
rules.requireIf(icebergRestSigningName, "glue",
new String[] {icebergRestSigningRegion,
icebergRestSigV4Enabled},
"Rest Catalog requires signing-region and sigv4-enabled set to
true when signing-name is glue");
rules.requireIf(icebergRestSigningName, "s3tables",
new String[] {icebergRestSigningRegion,
icebergRestSigV4Enabled},
"Rest Catalog requires signing-region and sigv4-enabled set to
true when signing-name is s3tables");
+ rules.requireIf(icebergRestSigningName, "osstables",
+ new String[] {icebergRestSigningRegion,
icebergRestSigV4Enabled},
+ "Rest Catalog requires signing-region and sigv4-enabled set to
true when signing-name is osstables");
+ rules.check(() -> shouldUseS3PropertiesForRestCredentials()
+ && !"true".equalsIgnoreCase(icebergRestSigV4Enabled),
+ "Rest Catalog requires sigv4-enabled set to true when
signing-name is " + icebergRestSigningName);
rejectUnsupportedAwsAssumeRoleProperty(ICEBERG_REST_ROLE_ARN);
rejectUnsupportedAwsAssumeRoleProperty(ICEBERG_REST_EXTERNAL_ID);
@@ -350,8 +362,14 @@ public class IcebergRestProperties extends
AbstractIcebergProperties {
icebergRestCatalogProperties.put("rest.signing-region",
icebergRestSigningRegion);
if (shouldUseS3PropertiesForRestCredentials()) {
-
IcebergAwsClientCredentialsProperties.putCredentialProviderProperties(
- icebergRestCatalogProperties, s3Properties);
+ if (isOssTables()) {
+
IcebergAwsClientCredentialsProperties.putCredentialProviderProperties(
+ icebergRestCatalogProperties,
ossProperties.getAccessKey(), ossProperties.getSecretKey(),
+ ossProperties.getSessionToken(),
icebergRestCredentialsProviderMode);
+ } else {
+
IcebergAwsClientCredentialsProperties.putCredentialProviderProperties(
+ icebergRestCatalogProperties, s3Properties);
+ }
} else {
IcebergAwsClientCredentialsProperties.putCredentialProviderProperties(
icebergRestCatalogProperties, icebergRestAccessKeyId,
@@ -362,7 +380,12 @@ public class IcebergRestProperties extends
AbstractIcebergProperties {
private boolean shouldUseS3PropertiesForRestCredentials() {
return "glue".equals(icebergRestSigningName)
- || "s3tables".equals(icebergRestSigningName);
+ || "s3tables".equals(icebergRestSigningName)
+ || "osstables".equals(icebergRestSigningName);
+ }
+
+ private boolean isOssTables() {
+ return "osstables".equals(icebergRestSigningName);
}
public Map<String, String> getIcebergRestCatalogProperties() {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java
b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java
index 4f643f61748..fea8a84c27b 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java
@@ -79,10 +79,11 @@ public class IcebergTableSink extends
BaseExternalTableDataSink {
this.targetTable = targetTable;
this.icebergTable = targetTable.getIcebergTable();
IcebergExternalCatalog catalog = (IcebergExternalCatalog)
targetTable.getCatalog();
- storagePropertiesMap =
VendedCredentialsFactory.getStoragePropertiesMapWithVendedCredentials(
- catalog.getCatalogProperty().getMetastoreProperties(),
- catalog.getCatalogProperty().getStoragePropertiesMap(),
- icebergTable);
+ storagePropertiesMap = IcebergUtils.selectEffectiveStorageProperties(
+
VendedCredentialsFactory.getStoragePropertiesMapWithVendedCredentials(
+ catalog.getCatalogProperty().getMetastoreProperties(),
+ catalog.getCatalogProperty().getStoragePropertiesMap(),
+ icebergTable));
}
public IcebergTableSink(IcebergExternalTable targetTable, Table
icebergTable) {
@@ -94,10 +95,11 @@ public class IcebergTableSink extends
BaseExternalTableDataSink {
// Keep credentials and every writer option on the metadata generation
pinned during analysis.
this.icebergTable = Objects.requireNonNull(icebergTable, "icebergTable
is not null");
IcebergExternalCatalog catalog = (IcebergExternalCatalog)
targetTable.getCatalog();
- storagePropertiesMap =
VendedCredentialsFactory.getStoragePropertiesMapWithVendedCredentials(
- catalog.getCatalogProperty().getMetastoreProperties(),
- catalog.getCatalogProperty().getStoragePropertiesMap(),
- icebergTable);
+ storagePropertiesMap = IcebergUtils.selectEffectiveStorageProperties(
+
VendedCredentialsFactory.getStoragePropertiesMapWithVendedCredentials(
+ catalog.getCatalogProperty().getMetastoreProperties(),
+ catalog.getCatalogProperty().getStoragePropertiesMap(),
+ icebergTable));
}
@Override
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/connectivity/IcebergRestConnectivityTesterTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/connectivity/IcebergRestConnectivityTesterTest.java
new file mode 100644
index 00000000000..4d0236b821c
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/connectivity/IcebergRestConnectivityTesterTest.java
@@ -0,0 +1,57 @@
+// 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.doris.datasource.connectivity;
+
+import org.apache.doris.datasource.property.metastore.IcebergRestProperties;
+import org.apache.doris.datasource.property.storage.StorageProperties;
+
+import org.apache.iceberg.CatalogProperties;
+import org.apache.iceberg.rest.RESTCatalog;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+public class IcebergRestConnectivityTesterTest {
+
+ @Test
+ public void testUsesNormalCatalogInitializationPath() throws Exception {
+ IcebergRestProperties properties =
Mockito.mock(IcebergRestProperties.class);
+ RESTCatalog catalog = Mockito.mock(RESTCatalog.class);
+ StorageProperties storageProperties =
Mockito.mock(StorageProperties.class);
+ List<StorageProperties> storagePropertiesList =
Collections.singletonList(storageProperties);
+ Map<String, String> catalogProperties = new HashMap<>();
+ catalogProperties.put(CatalogProperties.WAREHOUSE_LOCATION,
"s3://warehouse/path");
+
+ Mockito.when(properties.initializeCatalog("connectivity-test",
storagePropertiesList)).thenReturn(catalog);
+ Mockito.when(catalog.properties()).thenReturn(catalogProperties);
+
+ IcebergRestConnectivityTester tester = new
IcebergRestConnectivityTester(
+ properties, storagePropertiesList);
+ tester.testConnection();
+
+ Mockito.verify(properties).initializeCatalog("connectivity-test",
storagePropertiesList);
+ Mockito.verify(catalog).listNamespaces();
+ Mockito.verify(catalog).close();
+ Assertions.assertEquals("s3://warehouse/path",
tester.getTestLocation());
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
index 45619851cdc..21e918a8439 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
@@ -24,7 +24,11 @@ import org.apache.doris.catalog.StructField;
import org.apache.doris.catalog.Type;
import org.apache.doris.common.UserException;
import org.apache.doris.common.security.authentication.ExecutionAuthenticator;
+import org.apache.doris.common.util.LocationPath;
import org.apache.doris.datasource.iceberg.source.IcebergTableQueryInfo;
+import org.apache.doris.datasource.property.storage.OSSProperties;
+import org.apache.doris.datasource.property.storage.S3Properties;
+import org.apache.doris.datasource.property.storage.StorageProperties;
import org.apache.doris.nereids.exceptions.AnalysisException;
import org.apache.doris.system.Backend;
@@ -80,6 +84,32 @@ import java.util.UUID;
import java.util.concurrent.atomic.AtomicReference;
public class IcebergUtilsTest {
+ @Test
+ public void testSelectEffectiveStoragePropertiesPrefersOssOverGenericS3()
throws UserException {
+ Map<String, String> properties = new HashMap<>();
+ properties.put("iceberg.rest.signing-name", "osstables");
+ properties.put("iceberg.rest.signing-region", "cn-beijing");
+ properties.put("oss.endpoint", "https://oss-cn-beijing.aliyuncs.com");
+ properties.put("oss.region", "cn-beijing");
+ properties.put("oss.access_key", "ak");
+ properties.put("oss.secret_key", "sk");
+
+ Map<StorageProperties.Type, StorageProperties> detected = new
HashMap<>();
+ for (StorageProperties storageProperties :
StorageProperties.createAll(properties)) {
+ detected.put(storageProperties.getType(), storageProperties);
+ }
+ Assert.assertTrue(detected.get(StorageProperties.Type.S3) instanceof
S3Properties);
+ Assert.assertTrue(detected.get(StorageProperties.Type.OSS) instanceof
OSSProperties);
+
+ Map<StorageProperties.Type, StorageProperties> selected =
+ IcebergUtils.selectEffectiveStorageProperties(detected);
+
+ Assert.assertFalse(selected.containsKey(StorageProperties.Type.S3));
+ Assert.assertTrue(selected.get(StorageProperties.Type.OSS) instanceof
OSSProperties);
+ Assert.assertSame(selected.get(StorageProperties.Type.OSS),
+ LocationPath.of("s3://bucket/data.parquet",
selected).getStorageProperties());
+ }
+
@Test
public void testSnapshotCacheFreezesSharedTableOperations() {
Schema originalSchema = new Schema(
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/IcebergRestPropertiesTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/IcebergRestPropertiesTest.java
index cd9820c2722..9449e21f467 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/IcebergRestPropertiesTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/IcebergRestPropertiesTest.java
@@ -568,6 +568,85 @@ public class IcebergRestPropertiesTest {
Assertions.assertTrue(e.getMessage().contains("signing-region") &&
e.getMessage().contains("s3tables"));
}
+ @Test
+ public void testOssTablesRestCatalogUsesSharedS3Credentials() throws
Exception {
+ Map<String, String> props = new HashMap<>();
+ props.put("iceberg.rest.uri",
"https://cn-hangzhou.oss-tables.aliyuncs.com/iceberg");
+ props.put("warehouse",
"acs:osstables:cn-hangzhou:1234567890:bucket/my-table-bucket");
+ props.put("iceberg.rest.signing-name", "osstables");
+ props.put("iceberg.rest.signing-region", "cn-hangzhou");
+ props.put("iceberg.rest.sigv4-enabled", "true");
+ props.put("iceberg.rest.view-enabled", "false");
+ props.put("io-impl", "org.apache.iceberg.aws.s3.S3FileIO");
+ props.put("oss.endpoint", "https://oss-cn-hangzhou.aliyuncs.com");
+ props.put("oss.region", "cn-hangzhou");
+ props.put("oss.access_key", "oss-access-key");
+ props.put("oss.secret_key", "oss-secret-key");
+ props.put("oss.session_token", "oss-session-token");
+ props.put("oss.use_path_style", "true");
+
+ IcebergRestProperties restProps = new IcebergRestProperties(props);
+ restProps.initNormalizeAndCheckProps();
+ Assertions.assertFalse(restProps.isIcebergRestViewEnabled());
+
+ Map<String, String> catalogProps = new HashMap<>(props);
+ catalogProps.putAll(restProps.getIcebergRestCatalogProperties());
+ List<StorageProperties> storageProperties =
StorageProperties.createAll(props);
+ restProps.toFileIOProperties(storageProperties, catalogProps, new
Configuration());
+
+
Assertions.assertEquals("https://cn-hangzhou.oss-tables.aliyuncs.com/iceberg",
+ catalogProps.get(CatalogProperties.URI));
+
Assertions.assertEquals("acs:osstables:cn-hangzhou:1234567890:bucket/my-table-bucket",
+ catalogProps.get(CatalogProperties.WAREHOUSE_LOCATION));
+ Assertions.assertEquals("org.apache.iceberg.aws.s3.S3FileIO",
+ catalogProps.get(CatalogProperties.FILE_IO_IMPL));
+ Assertions.assertEquals("osstables",
catalogProps.get("rest.signing-name"));
+ Assertions.assertEquals("cn-hangzhou",
catalogProps.get("rest.signing-region"));
+ Assertions.assertEquals("true",
catalogProps.get("rest.sigv4-enabled"));
+ Assertions.assertEquals("oss-access-key",
catalogProps.get("rest.access-key-id"));
+ Assertions.assertEquals("oss-secret-key",
catalogProps.get("rest.secret-access-key"));
+ Assertions.assertEquals("oss-session-token",
catalogProps.get("rest.session-token"));
+
Assertions.assertTrue(storageProperties.stream().anyMatch(OSSProperties.class::isInstance));
+ Assertions.assertEquals("https://oss-cn-hangzhou.aliyuncs.com",
+ catalogProps.get(S3FileIOProperties.ENDPOINT));
+ Assertions.assertEquals("cn-hangzhou",
catalogProps.get(AwsClientProperties.CLIENT_REGION));
+ Assertions.assertEquals("oss-access-key",
catalogProps.get(S3FileIOProperties.ACCESS_KEY_ID));
+ Assertions.assertEquals("oss-secret-key",
catalogProps.get(S3FileIOProperties.SECRET_ACCESS_KEY));
+ Assertions.assertEquals("oss-session-token",
catalogProps.get(S3FileIOProperties.SESSION_TOKEN));
+ Assertions.assertEquals("true",
catalogProps.get(S3FileIOProperties.PATH_STYLE_ACCESS));
+ }
+
+ @Test
+ public void testOssTablesSigningNameMissingSigningRegionFails() {
+ Map<String, String> props = new HashMap<>();
+ props.put("iceberg.rest.uri",
"https://cn-hangzhou.oss-tables.aliyuncs.com/iceberg");
+ props.put("iceberg.rest.signing-name", "osstables");
+ props.put("iceberg.rest.sigv4-enabled", "true");
+ props.put("s3.access_key", "oss-access-key");
+ props.put("s3.secret_key", "oss-secret-key");
+
+ IcebergRestProperties restProps = new IcebergRestProperties(props);
+ IllegalArgumentException e =
Assertions.assertThrows(IllegalArgumentException.class,
+ restProps::initNormalizeAndCheckProps);
+ Assertions.assertTrue(e.getMessage().contains("signing-region") &&
e.getMessage().contains("osstables"));
+ }
+
+ @Test
+ public void testOssTablesSigningNameWithSigV4DisabledFails() {
+ Map<String, String> props = new HashMap<>();
+ props.put("iceberg.rest.uri",
"https://cn-hangzhou.oss-tables.aliyuncs.com/iceberg");
+ props.put("iceberg.rest.signing-name", "osstables");
+ props.put("iceberg.rest.signing-region", "cn-hangzhou");
+ props.put("iceberg.rest.sigv4-enabled", "false");
+ props.put("s3.access_key", "oss-access-key");
+ props.put("s3.secret_key", "oss-secret-key");
+
+ IcebergRestProperties restProps = new IcebergRestProperties(props);
+ IllegalArgumentException e =
Assertions.assertThrows(IllegalArgumentException.class,
+ restProps::initNormalizeAndCheckProps);
+ Assertions.assertTrue(e.getMessage().contains("sigv4-enabled") &&
e.getMessage().contains("osstables"));
+ }
+
@Test
public void testAccessKeyAndSecretKeyMustBeSetTogether() {
Map<String, String> props1 = new HashMap<>();
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/property/storage/StoragePropertiesTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/property/storage/StoragePropertiesTest.java
index 022c1e9b890..77315099a22 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/property/storage/StoragePropertiesTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/property/storage/StoragePropertiesTest.java
@@ -70,6 +70,24 @@ public class StoragePropertiesTest {
"OSS should be detected via guessIsMe when no explicit
fs.xx.support is set");
}
+ @Test
+ public void testOssTablesPropertiesKeepS3AndOssProvidersAvailable() throws
UserException {
+ Map<String, String> props = new HashMap<>();
+ props.put("iceberg.rest.signing-name", "osstables");
+ props.put("iceberg.rest.signing-region", "cn-beijing");
+ props.put("oss.endpoint", "https://oss-cn-beijing.aliyuncs.com");
+ props.put("oss.region", "cn-beijing");
+ props.put("oss.access_key", "ak");
+ props.put("oss.secret_key", "sk");
+
+ List<StorageProperties> all = StorageProperties.createAll(props);
+ List<Class<?>> types = toTypeList(all);
+
+ Assertions.assertTrue(types.contains(OSSProperties.class));
+ Assertions.assertTrue(types.contains(S3Properties.class),
+ "Global storage detection must keep S3 available when OSS is
also configured");
+ }
+
/**
* When no {@code fs.xx.support} flag is set, an S3 endpoint containing
* "amazonaws.com" should be detected as S3 via guessIsMe.
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergTableSinkTest.java
b/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergTableSinkTest.java
index 334c70880e2..7b93261a1ba 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergTableSinkTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergTableSinkTest.java
@@ -20,6 +20,9 @@ package org.apache.doris.planner;
import org.apache.doris.datasource.CatalogProperty;
import org.apache.doris.datasource.iceberg.IcebergExternalCatalog;
import org.apache.doris.datasource.iceberg.IcebergExternalTable;
+import org.apache.doris.datasource.property.storage.OSSProperties;
+import org.apache.doris.datasource.property.storage.S3Properties;
+import org.apache.doris.datasource.property.storage.StorageProperties;
import
org.apache.doris.nereids.trees.plans.commands.insert.IcebergInsertCommandContext;
import org.apache.doris.thrift.TIcebergTableSink;
@@ -34,10 +37,55 @@ import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
import java.util.Collections;
+import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
public class IcebergTableSinkTest {
+ @Test
+ public void testBindPrefersOssDataPlanePropertiesOverGenericS3() throws
Exception {
+ Map<String, String> properties = new HashMap<>();
+ properties.put("iceberg.rest.signing-name", "osstables");
+ properties.put("iceberg.rest.signing-region", "cn-beijing");
+ properties.put("oss.endpoint", "https://oss-cn-beijing.aliyuncs.com");
+ properties.put("oss.region", "cn-beijing");
+ properties.put("oss.access_key", "oss-ak");
+ properties.put("oss.secret_key", "oss-sk");
+
+ Map<StorageProperties.Type, StorageProperties> storagePropertiesMap =
new HashMap<>();
+ for (StorageProperties storageProperties :
StorageProperties.createAll(properties)) {
+ storagePropertiesMap.put(storageProperties.getType(),
storageProperties);
+ }
+
Assertions.assertTrue(storagePropertiesMap.get(StorageProperties.Type.S3)
instanceof S3Properties);
+
Assertions.assertTrue(storagePropertiesMap.get(StorageProperties.Type.OSS)
instanceof OSSProperties);
+
+ IcebergExternalCatalog catalog =
Mockito.mock(IcebergExternalCatalog.class);
+ CatalogProperty catalogProperty = Mockito.mock(CatalogProperty.class);
+ Mockito.when(catalog.getCatalogProperty()).thenReturn(catalogProperty);
+
Mockito.when(catalogProperty.getMetastoreProperties()).thenReturn(null);
+
Mockito.when(catalogProperty.getStoragePropertiesMap()).thenReturn(storagePropertiesMap);
+
+ IcebergExternalTable targetTable =
Mockito.mock(IcebergExternalTable.class);
+ Mockito.when(targetTable.isView()).thenReturn(false);
+ Mockito.when(targetTable.getCatalog()).thenReturn(catalog);
+ Mockito.when(targetTable.getDbName()).thenReturn("db");
+ Mockito.when(targetTable.getName()).thenReturn("table");
+
+ Schema schema = new Schema(1,
+ Types.NestedField.required(1, "id", Types.IntegerType.get()));
+ Table table = mockTable(schema);
+ Mockito.when(table.location()).thenReturn("s3://bucket/table");
+
+ IcebergTableSink sink = new IcebergTableSink(targetTable, table);
+ sink.bindDataSink(Optional.empty());
+
+ TIcebergTableSink thriftSink = sink.tDataSink.getIcebergTableSink();
+ Assertions.assertEquals("https://oss-cn-beijing.aliyuncs.com",
+ thriftSink.getHadoopConfig().get("AWS_ENDPOINT"));
+ Assertions.assertEquals("oss-ak",
thriftSink.getHadoopConfig().get("AWS_ACCESS_KEY"));
+ Assertions.assertEquals("s3://bucket/table/data",
thriftSink.getOriginalOutputPath());
+ }
+
@Test
public void testBindUsesPinnedIcebergTableMetadata() throws Exception {
IcebergExternalCatalog catalog =
Mockito.mock(IcebergExternalCatalog.class);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]