This is an automated email from the ASF dual-hosted git repository. diqiu50 pushed a commit to branch dell-1.3 in repository https://gitbox.apache.org/repos/asf/gravitino.git
commit f1beaf7b97c3bf29fe10cc839e1890a80ab33180 Author: diqiu50 <[email protected]> AuthorDate: Fri Aug 21 18:16:08 2026 +0800 [Cherry-pick to branch-1.3] [#12554] improvement(trino-connector): Route lakehouse-iceberg catalogs through the Iceberg REST server (cherry picked from commit 859ad4c82049c29099c90298f3f51085042d4b83) --- docs/trino-connector/authentication.md | 20 +- docs/trino-connector/catalog-iceberg.md | 79 +++++++ docs/trino-connector/configuration.md | 11 + .../docker-script/docker-compose.yaml | 2 + .../init/trino/config/catalog/gravitino.properties | 1 + .../docker-script/init/trino/init.sh | 3 + .../integration/test/container/ContainerSuite.java | 2 + .../test/container/TrinoITContainers.java | 6 +- .../gravitino/integration/test/util/BaseIT.java | 22 ++ trino-connector/integration-test/build.gradle.kts | 1 + .../integration/test/TrinoConnectorIT.java | 19 ++ .../integration/test/TrinoQueryITBase.java | 13 ++ .../gravitino/trino/connector/GravitinoConfig.java | 76 ++++++- .../catalog/DefaultCatalogConnectorFactory.java | 2 +- .../iceberg/IcebergCatalogPropertyConverter.java | 120 +++++++++- .../catalog/iceberg/IcebergConnectorAdapter.java | 29 ++- .../trino/connector/TestGravitinoConfig.java | 53 +++++ .../TestIcebergCatalogPropertyConverter.java | 253 ++++++++++++++++++++- 18 files changed, 696 insertions(+), 16 deletions(-) diff --git a/docs/trino-connector/authentication.md b/docs/trino-connector/authentication.md index 9793e7749f..322eef9f5d 100644 --- a/docs/trino-connector/authentication.md +++ b/docs/trino-connector/authentication.md @@ -193,7 +193,25 @@ Whether the coordinator can populate this extra-credential depends on the Trino The `gravitino.client.oauth2.*` properties above still configure the shared bootstrap/admin client used for catalog discovery — they are unrelated to the per-user forwarded token. -For an Iceberg catalog with `catalog-backend=rest` (backed by an Iceberg REST Catalog), the connector does not set `iceberg.rest-catalog.security`/`iceberg.rest-catalog.session` on its own — that catalog's own `gravitino.client.*` config is unrelated to how its underlying Iceberg REST catalog authenticates. To also forward the end user's token to the REST catalog itself, set `trino.bypass.iceberg.rest-catalog.security=OAUTH2` and `trino.bypass.iceberg.rest-catalog.session=USER` explicitl [...] +For an Iceberg catalog reached through the Gravitino Iceberg REST server (IRC) — every +`lakehouse-iceberg` catalog, unless `gravitino.iceberg.rest-enabled=false`; see [Iceberg +catalog](./catalog-iceberg.md#how-trino-reaches-the-catalog) — the IRC's own authentication is +configured once per Trino cluster with the `gravitino.iceberg.rest-catalog.` prefix, and +`iceberg.rest-catalog.session=USER` is set automatically when `forwardUser=true`: + +```properties +gravitino.iceberg.rest-uri=http://gravitino-host:9001/iceberg +gravitino.iceberg.rest-catalog.security=OAUTH2 +gravitino.iceberg.rest-catalog.oauth2.credential=service-account-id:service-account-secret +gravitino.iceberg.rest-catalog.oauth2.server-uri=http://your-idp/realms/gravitino/protocol/openid-connect/token +gravitino.iceberg.rest-catalog.oauth2.scope=email +``` + +This is a completely separate credential from the one the connector uses against the main Gravitino +server; it is not reused automatically. + +For an Iceberg catalog with `catalog-backend=rest` (pointing at an Iceberg REST Catalog of its own, +which the connector does not re-route), the connector does not set `iceberg.rest-catalog.security`/`iceberg.rest-catalog.session` on its own — that catalog's own `gravitino.client.*` config is unrelated to how its underlying Iceberg REST catalog authenticates. To also forward the end user's token to the REST catalog itself, set `trino.bypass.iceberg.rest-catalog.security=OAUTH2` and `trino.bypass.iceberg.rest-catalog.session=USER` explicitly on that catalog's properties, alongside its bo [...] **Configuration properties:** diff --git a/docs/trino-connector/catalog-iceberg.md b/docs/trino-connector/catalog-iceberg.md index fbecbcb7d0..223b597fe9 100644 --- a/docs/trino-connector/catalog-iceberg.md +++ b/docs/trino-connector/catalog-iceberg.md @@ -20,6 +20,85 @@ To use Iceberg, you need: - ORC - Parquet (default) +## How Trino Reaches the Catalog + +The Gravitino Trino connector loads every `lakehouse-iceberg` catalog through the Gravitino Iceberg +REST server (IRC), regardless of the catalog's `catalog-backend`. `catalog-backend` describes how +Gravitino stores the catalog's metadata; it does not decide how the query engine reaches the data. + +This is what makes [credential vending](../security/credential-vending.md) work. Trino only consumes +vended credentials in its `rest` Iceberg catalog type — the `jdbc` and `hive_metastore` types have +nowhere to put the session token of an STS temporary credential — so a catalog with +`credential-providers=s3-token` produces no usable credential on those paths. Routing through the +IRC means every table access gets a freshly issued temporary credential over the Iceberg REST +protocol. + +Configure the IRC endpoint once per Trino cluster, in `etc/catalog/gravitino.properties`: + +```properties +connector.name=gravitino +gravitino.metalake=test +gravitino.uri=http://gravitino-host:8090 + +gravitino.iceberg.rest-uri=http://gravitino-host:9001/iceberg +``` + +The endpoint is deployment topology, not a property of the data source, so it does not belong on the +catalog: it would have to be repeated on every new catalog, moving the IRC would mean editing all of +them, and one catalog could not serve two Trino clusters that reach the IRC by different hostnames. + +The connector derives everything else from the catalog itself. The Gravitino catalog name is passed +as both `iceberg.rest-catalog.warehouse` and `iceberg.rest-catalog.prefix` — the Iceberg client +selects the catalog twice over, first as the query parameter of the `GET /v1/config` call that +discovers it, then as the path segment of every request after that. The Trino native file system +(`fs.native-s3.enabled` and `s3.region`, or the GCS and Azure equivalents) is derived from the +catalog's `warehouse` scheme and storage properties. + +On the Gravitino side, the IRC must run with the dynamic config provider so it can serve catalogs +defined in Gravitino: + +```properties +gravitino.auxService.names = iceberg-rest +gravitino.iceberg-rest.catalog-config-provider = dynamic-config-provider +gravitino.iceberg-rest.gravitino-metalake = test +``` + +### Authenticating to the Iceberg REST Server + +If the IRC has authentication enabled, the internal Iceberg REST catalog that the connector builds +must authenticate to it on its own. This is a separate credential from the one the connector uses +against the main Gravitino server, and it is not reused automatically. Any property prefixed with +`gravitino.iceberg.rest-catalog.` is passed through to the internal catalog with the prefix +rewritten to `iceberg.rest-catalog.`: + +```properties +gravitino.iceberg.rest-catalog.security=OAUTH2 +gravitino.iceberg.rest-catalog.oauth2.credential=client_id:client_secret +gravitino.iceberg.rest-catalog.oauth2.server-uri=http://your-idp/token +gravitino.iceberg.rest-catalog.oauth2.scope=email +``` + +Omitting this block against an authenticated IRC surfaces as an authentication error rather than a +missing-configuration error, which is easy to misdiagnose. + +When `gravitino.client.session.forwardUser=true`, the connector also sets +`iceberg.rest-catalog.session=USER` so that each query carries the end user's identity to the IRC, +keeping per-user credential vending and per-user authorization intact. Set +`gravitino.iceberg.rest-catalog.session` explicitly to override it. See +[Authentication](./authentication.md) for the full setup. + +### Limitations + +- One IRC serves exactly one metalake, fixed at startup by + `gravitino.iceberg-rest.gravitino-metalake`. In multi-metalake mode + (`gravitino.use-single-metalake=false`), only catalogs in that metalake are reachable; the others + fail on the IRC side with `NoSuchCatalogException`. +- A catalog created with `catalog-backend=rest` keeps pointing at its own configured `uri` and is + not re-routed, since it already reaches an Iceberg REST catalog directly. +- To disable this routing entirely — for a deployment that does not run the IRC — set + `gravitino.iceberg.rest-enabled=false`. Iceberg catalogs are then translated into Trino's `jdbc` + or `hive_metastore` catalog types as before, and credential vending does not work. + ## Schema Operations ### Create a Schema diff --git a/docs/trino-connector/configuration.md b/docs/trino-connector/configuration.md index 5bcb98636e..0436de81c5 100644 --- a/docs/trino-connector/configuration.md +++ b/docs/trino-connector/configuration.md @@ -29,11 +29,22 @@ license: "This software is licensed under the Apache License version 2." | gravitino.client. | string | (none) | The configuration key prefix for the Gravitino client config. | No | | gravitino.trino.skip-catalog-patterns | string | (none) | The `gravitino.trino.skip-catalog-patterns` defines a comma-separated list of catalog name regex patterns that should be excluded from loading. For example, `test_.*, .*_tmp` excludes all catalogs starting with `test_` or ending with `_tmp`. | No | | gravitino.use-single-metalake | boolean | true | If `true`, only one metalake is used and catalogs are identified by `<catalog_name>`. If `false`, multi-metalake mode is enabled and catalogs are identified by `<metalake_name>.<catalog_name>`. | No | +| gravitino.iceberg.rest-enabled | boolean | true | If `true`, `lakehouse-iceberg` catalogs are loaded through the Gravitino Iceberg REST server (IRC) instead of being translated into Trino's `jdbc` or `hive_metastore` Iceberg catalog type. This is what makes credential vending work. Requires `gravitino.iceberg.rest-uri`. Set to `false` for a deployment that does not run the IRC. | No | +| gravitino.iceberg.rest-uri | string | (none) | The endpoint of the Gravitino Iceberg REST server, for example `http://gravitino-host:9001/iceberg`. Required when `gravitino.iceberg.rest-enabled` is `true` and the metalake contains `lakehouse-iceberg` catalogs. | No | +| gravitino.iceberg.rest-catalog. | string | (none) | The configuration key prefix for properties passed through to the internal Trino Iceberg REST catalog. The prefix is rewritten to `iceberg.rest-catalog.`, so `gravitino.iceberg.rest-catalog.security=OAUTH2` becomes `iceberg.rest-catalog.security=OAUTH2`. | No | To configure the Gravitino client, use properties prefixed with `gravitino.client.`. These properties will directly passed to the Gravitino client. **Note:** Invalid configuration properties will result in exceptions. Please see [Gravitino Java client configurations](../how-to-use-gravitino-client.md#java-client-configuration) for more support client configuration. +:::caution +`gravitino.iceberg.rest-enabled` defaults to `true`. When upgrading an existing deployment whose +metalake contains `lakehouse-iceberg` catalogs with `catalog-backend=jdbc` or `hive`, those catalogs +fail to load unless you either set `gravitino.iceberg.rest-uri` to your Iceberg REST server endpoint, +or set `gravitino.iceberg.rest-enabled=false` to keep the previous behavior. See +[Iceberg catalog](./catalog-iceberg.md#how-trino-reaches-the-catalog). +::: + Multi-metalake mode (`gravitino.use-single-metalake=false`) is supported on Trino connector versions 435-445 and 469-478. On versions 446-468, a warning is logged and the connector initializes, but the mode is not fully supported and some operations may fail. ## Connecting to a TLS-enabled coordinator diff --git a/integration-test-common/docker-script/docker-compose.yaml b/integration-test-common/docker-script/docker-compose.yaml index 3ed57d07ca..1dfb8f5a7d 100644 --- a/integration-test-common/docker-script/docker-compose.yaml +++ b/integration-test-common/docker-script/docker-compose.yaml @@ -87,6 +87,7 @@ services: - HADOOP_USER_NAME=anonymous - GRAVITINO_HOST_IP=host.docker.internal - GRAVITINO_HOST_PORT=${GRAVITINO_SERVER_PORT:-8090} + - GRAVITINO_ICEBERG_REST_PORT=${GRAVITINO_ICEBERG_REST_PORT:-9001} - GRAVITINO_METALAKE_NAME=test - HIVE_HOST_IP=hive - TRINO_WORKER_NUM=${TRINO_WORKER_NUM:-0} @@ -122,6 +123,7 @@ services: - HADOOP_USER_NAME=anonymous - GRAVITINO_HOST_IP=host.docker.internal - GRAVITINO_HOST_PORT=${GRAVITINO_SERVER_PORT:-8090} + - GRAVITINO_ICEBERG_REST_PORT=${GRAVITINO_ICEBERG_REST_PORT:-9001} - GRAVITINO_METALAKE_NAME=test - HIVE_HOST_IP=hive - TRINO_ROLE=worker diff --git a/integration-test-common/docker-script/init/trino/config/catalog/gravitino.properties b/integration-test-common/docker-script/init/trino/config/catalog/gravitino.properties index 08789c84c0..d38c18c458 100644 --- a/integration-test-common/docker-script/init/trino/config/catalog/gravitino.properties +++ b/integration-test-common/docker-script/init/trino/config/catalog/gravitino.properties @@ -20,6 +20,7 @@ connector.name = gravitino gravitino.uri = http://GRAVITINO_HOST_IP:GRAVITINO_HOST_PORT gravitino.metalake = GRAVITINO_METALAKE_NAME +gravitino.iceberg.rest-uri = http://GRAVITINO_ICEBERG_REST_HOST:GRAVITINO_ICEBERG_REST_PORT/iceberg gravitino.trino.skip-version-validation=true gravitino.client.authType = simple gravitino.client.session.forwardUser = true diff --git a/integration-test-common/docker-script/init/trino/init.sh b/integration-test-common/docker-script/init/trino/init.sh index 077c6f77d0..d713708e52 100644 --- a/integration-test-common/docker-script/init/trino/init.sh +++ b/integration-test-common/docker-script/init/trino/init.sh @@ -34,6 +34,9 @@ cp /usr/lib/trino/plugin/mysql/*mysql-connector-j-*.jar /usr/lib/trino/plugin/ic sed -i "s/GRAVITINO_HOST_IP:GRAVITINO_HOST_PORT/${GRAVITINO_HOST_IP}:${GRAVITINO_HOST_PORT}/g" /etc/trino/catalog/gravitino.properties # Update `gravitino.metalake = GRAVITINO_METALAKE_NAME` in the `conf/catalog/gravitino.properties` sed -i "s/GRAVITINO_METALAKE_NAME/${GRAVITINO_METALAKE_NAME}/g" /etc/trino/catalog/gravitino.properties +# Update `gravitino.iceberg.rest-uri` in the `conf/catalog/gravitino.properties`. It uses its own +# placeholders so that the substitution above cannot leave a half-replaced host behind. +sed -i "s/GRAVITINO_ICEBERG_REST_HOST:GRAVITINO_ICEBERG_REST_PORT/${GRAVITINO_HOST_IP}:${GRAVITINO_ICEBERG_REST_PORT}/g" /etc/trino/catalog/gravitino.properties # Update `node.id=NODE_ID` in the `/conf/node.properties` sed -i "s/NODE_ID/${RANDOM}-${RANDOM}-${RANDOM}-${RANDOM}-${RANDOM}/g" /etc/trino/node.properties diff --git a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java index b2b0ae50da..bc29a24e06 100644 --- a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java +++ b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java @@ -288,6 +288,7 @@ public class ContainerSuite implements Closeable { String trinoConfDir, String trinoConnectorLibDir, int gravitinoServerPort, + int icebergRestServerPort, String metalakeName) { ITUtils.cleanDisk(); if (trinoContainer == null) { @@ -303,6 +304,7 @@ public class ContainerSuite implements Closeable { .put("HADOOP_USER_NAME", "anonymous") .put("GRAVITINO_HOST_IP", "host.docker.internal") .put("GRAVITINO_HOST_PORT", String.valueOf(gravitinoServerPort)) + .put("GRAVITINO_ICEBERG_REST_PORT", String.valueOf(icebergRestServerPort)) .put("GRAVITINO_METALAKE_NAME", metalakeName) .build()) .withNetwork(getNetwork()) diff --git a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/TrinoITContainers.java b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/TrinoITContainers.java index a02e12de6a..14fe768c0b 100644 --- a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/TrinoITContainers.java +++ b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/TrinoITContainers.java @@ -50,11 +50,12 @@ public class TrinoITContainers implements AutoCloseable { } public void launch(int gravitinoServerPort) throws Exception { - launch(gravitinoServerPort, "hive2", false, null, null, null); + launch(gravitinoServerPort, 0, "hive2", false, null, null, null); } public void launch( int gravitinoServerPort, + int icebergRestServerPort, String hiveRuntimeVersion, boolean isTrinoConnectorTest, Integer trinoWorkerNum, @@ -81,6 +82,9 @@ public class TrinoITContainers implements AutoCloseable { } env.put("GRAVITINO_SERVER_PORT", String.valueOf(gravitinoServerPort)); env.put("HIVE_RUNTIME_VERSION", hiveRuntimeVersion); + if (icebergRestServerPort > 0) { + env.put("GRAVITINO_ICEBERG_REST_PORT", String.valueOf(icebergRestServerPort)); + } env.put("TRINO_CONNECTOR_TEST", String.valueOf(isTrinoConnectorTest)); if (System.getProperty("gravitino.log.path") != null) { env.put("GRAVITINO_LOG_PATH", System.getProperty("gravitino.log.path")); diff --git a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java index 8c9cbfa510..f3a0724ef0 100644 --- a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java +++ b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java @@ -30,6 +30,7 @@ import com.google.common.collect.ImmutableMap; import java.io.File; import java.io.IOException; import java.lang.reflect.Field; +import java.net.URI; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.nio.file.Path; @@ -674,6 +675,27 @@ public class BaseIT { return null; } + /** + * Enables the Iceberg REST auxiliary service before the server starts. Callers that hold a {@link + * BaseIT} instance rather than extending it use this instead of the protected fields. + * + * @param icebergRestConfigs extra server configs, typically the dynamic config provider and the + * metalake the Iceberg REST service should serve + */ + public void enableIcebergAuxRestService(Map<String, String> icebergRestConfigs) { + this.ignoreIcebergAuxRestService = false; + this.customConfigs.putAll(icebergRestConfigs); + } + + /** + * Returns the port the Iceberg REST auxiliary service listens on. + * + * @return the Iceberg REST service port + */ + public int getIcebergRestServicePort() { + return URI.create(getIcebergRestServiceUri()).getPort(); + } + protected String getIcebergRestServiceUri() { JettyServerConfig jettyServerConfig = JettyServerConfig.fromConfig(serverConfig, String.format("gravitino.iceberg-rest.")); diff --git a/trino-connector/integration-test/build.gradle.kts b/trino-connector/integration-test/build.gradle.kts index f66c9e3649..c48e9719bf 100644 --- a/trino-connector/integration-test/build.gradle.kts +++ b/trino-connector/integration-test/build.gradle.kts @@ -34,6 +34,7 @@ dependencies { } testImplementation(project(":api")) + testImplementation(project(":catalogs:catalog-common")) testImplementation(project(":clients:client-java")) testImplementation(project(":common")) testImplementation(project(":core")) diff --git a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java index 10eb37db74..d0796e044f 100644 --- a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java +++ b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java @@ -23,6 +23,7 @@ import static org.testcontainers.shaded.org.awaitility.Awaitility.await; import com.google.common.collect.ImmutableMap; import com.google.common.collect.Maps; import com.google.common.collect.Sets; +import java.net.URI; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -33,6 +34,7 @@ import java.util.concurrent.TimeUnit; import org.apache.gravitino.Catalog; import org.apache.gravitino.NameIdentifier; import org.apache.gravitino.Schema; +import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants; import org.apache.gravitino.client.GravitinoMetalake; import org.apache.gravitino.integration.test.container.ContainerSuite; import org.apache.gravitino.integration.test.container.HiveContainer; @@ -76,6 +78,8 @@ public class TrinoConnectorIT extends BaseIT { private static final ContainerSuite containerSuite = ContainerSuite.getInstance(); + private static final String GRAVITINO_ICEBERG_REST_PREFIX = "gravitino.iceberg-rest."; + public static String metalakeName = GravitinoITUtils.genRandomName("TrinoIT_metalake").toLowerCase(); public static String catalogName = @@ -90,6 +94,20 @@ public class TrinoConnectorIT extends BaseIT { private static GravitinoMetalake metalake; private static Catalog catalog; + @BeforeAll + @Override + public void startIntegrationTest() throws Exception { + // The Trino connector loads every lakehouse-iceberg catalog through the Iceberg REST server, + // so the auxiliary service has to run and serve this test's metalake. + ignoreIcebergAuxRestService = false; + customConfigs.put( + GRAVITINO_ICEBERG_REST_PREFIX + IcebergConstants.ICEBERG_REST_CATALOG_CONFIG_PROVIDER, + IcebergConstants.DYNAMIC_ICEBERG_CATALOG_CONFIG_PROVIDER_NAME); + customConfigs.put( + GRAVITINO_ICEBERG_REST_PREFIX + IcebergConstants.GRAVITINO_METALAKE, metalakeName); + super.startIntegrationTest(); + } + @BeforeAll public void startDockerContainer() throws TException, InterruptedException { String trinoConfDir = System.getenv("TRINO_CONF_DIR"); @@ -113,6 +131,7 @@ public class TrinoConnectorIT extends BaseIT { trinoConfDir, System.getenv("GRAVITINO_ROOT_DIR") + "/trino-connector/build/libs", getGravitinoServerPort(), + URI.create(getIcebergRestServiceUri()).getPort(), metalakeName); Assertions.assertTrue( containerSuite.getTrinoContainer().checkSyncCatalogFromGravitino(5, catalogName), diff --git a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java index 5d955bca7e..24d33cbf2e 100644 --- a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java +++ b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java @@ -20,6 +20,7 @@ package org.apache.gravitino.trino.connector.integration.test; import static java.lang.Thread.sleep; +import com.google.common.collect.ImmutableMap; import java.io.File; import java.io.IOException; import java.nio.charset.StandardCharsets; @@ -31,6 +32,7 @@ import org.apache.gravitino.Catalog; import org.apache.gravitino.NameIdentifier; import org.apache.gravitino.Namespace; import org.apache.gravitino.SupportsSchemas; +import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants; import org.apache.gravitino.client.GravitinoAdminClient; import org.apache.gravitino.client.GravitinoMetalake; import org.apache.gravitino.exceptions.RESTException; @@ -45,6 +47,8 @@ import org.slf4j.LoggerFactory; public class TrinoQueryITBase { private static final Logger LOG = LoggerFactory.getLogger(TrinoQueryITBase.class); + private static final String GRAVITINO_ICEBERG_REST_PREFIX = "gravitino.iceberg-rest."; + // Auto start docker containers and Gravitino server protected static boolean autoStart = true; @@ -87,6 +91,14 @@ public class TrinoQueryITBase { private void setEnv() throws Exception { baseIT = new BaseIT(); if (autoStart) { + // The Trino connector loads every lakehouse-iceberg catalog through the Iceberg REST server, + // so the auxiliary service has to run and serve this test's metalake. + baseIT.enableIcebergAuxRestService( + ImmutableMap.of( + GRAVITINO_ICEBERG_REST_PREFIX + IcebergConstants.ICEBERG_REST_CATALOG_CONFIG_PROVIDER, + IcebergConstants.DYNAMIC_ICEBERG_CATALOG_CONFIG_PROVIDER_NAME, + GRAVITINO_ICEBERG_REST_PREFIX + IcebergConstants.GRAVITINO_METALAKE, + metalakeName)); baseIT.startIntegrationTest(); gravitinoClient = baseIT.getGravitinoClient(); gravitinoUri = String.format("http://127.0.0.1:%d", baseIT.getGravitinoServerPort()); @@ -108,6 +120,7 @@ public class TrinoQueryITBase { trinoITContainers.launch( baseIT.getGravitinoServerPort(), + baseIT.getIcebergRestServicePort(), hiveRuntimeVersion, isTrinoConnectorTest, trinoWorkerNum, diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java index 817d9f6306..e79645d0e6 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java @@ -76,6 +76,9 @@ public class GravitinoConfig { public static final String GRAVITINO_DYNAMIC_CONNECTOR_CATALOG_CONFIG = "__gravitino.dynamic.connector.catalog.config"; + /** The Trino Iceberg REST catalog property prefix. */ + private static final String TRINO_ICEBERG_REST_CATALOG_PREFIX = "iceberg.rest-catalog."; + private static final Map<String, ConfigEntry> CONFIG_DEFINITIONS = new HashMap<>(); private final Map<String, String> config; private final List<Pattern> skipCatalogPatternList; @@ -263,6 +266,32 @@ public class GravitinoConfig { "3600", false); + private static final ConfigEntry GRAVITINO_ICEBERG_REST_ENABLED = + new ConfigEntry( + "gravitino.iceberg.rest-enabled", + "When true, lakehouse-iceberg catalogs are loaded through the Gravitino Iceberg REST " + + "server instead of being translated into a Trino JDBC or Hive metastore Iceberg " + + "catalog. Requires gravitino.iceberg.rest-uri.", + "true", + false); + + private static final ConfigEntry GRAVITINO_ICEBERG_REST_URI = + new ConfigEntry( + "gravitino.iceberg.rest-uri", + "The endpoint of the Gravitino Iceberg REST server, for example " + + "http://localhost:9001/iceberg.", + "", + false); + + private static final ConfigEntry GRAVITINO_ICEBERG_REST_CATALOG_CONFIG_PREFIX = + new ConfigEntry( + "gravitino.iceberg.rest-catalog.", + "Prefix for properties passed through to the internal Trino Iceberg REST catalog. Any " + + "property beginning with this prefix is rewritten to iceberg.rest-catalog. and " + + "passed through (e.g., gravitino.iceberg.rest-catalog.security=OAUTH2).", + "", + false); + /** * Constructs a new GravitinoConfig with the specified configuration. * @@ -600,9 +629,13 @@ public class GravitinoConfig { stringList.add(String.format("\"%s\"='%s'", entry.getKey(), value)); } } - // copy the configuration by the prefix of GRAVITINO_CLIENT_CONFIG_PREFIX + // copy the configuration by the prefix of GRAVITINO_CLIENT_CONFIG_PREFIX and + // GRAVITINO_ICEBERG_REST_CATALOG_CONFIG_PREFIX config.entrySet().stream() - .filter(entry -> entry.getKey().startsWith(GRAVITINO_CLIENT_CONFIG_PREFIX.key)) + .filter( + entry -> + entry.getKey().startsWith(GRAVITINO_CLIENT_CONFIG_PREFIX.key) + || entry.getKey().startsWith(GRAVITINO_ICEBERG_REST_CATALOG_CONFIG_PREFIX.key)) .forEach( entry -> stringList.add(String.format("\"%s\"='%s'", entry.getKey(), entry.getValue()))); @@ -678,6 +711,45 @@ public class GravitinoConfig { return parseLongConfigEntry(GRAVITINO_SESSION_CACHE_EXPIRE_AFTER_ACCESS_SECONDS); } + /** + * Returns whether lakehouse-iceberg catalogs are routed through the Gravitino Iceberg REST + * server. + * + * @return true if the Iceberg REST routing is enabled + */ + public boolean isIcebergRestEnabled() { + return Boolean.parseBoolean( + config.getOrDefault( + GRAVITINO_ICEBERG_REST_ENABLED.key, GRAVITINO_ICEBERG_REST_ENABLED.defaultValue)); + } + + /** + * Retrieves the endpoint of the Gravitino Iceberg REST server. + * + * @return the Iceberg REST server endpoint, or an empty string if not configured + */ + public String getIcebergRestUri() { + return config.getOrDefault( + GRAVITINO_ICEBERG_REST_URI.key, GRAVITINO_ICEBERG_REST_URI.defaultValue); + } + + /** + * Retrieves the properties passed through to the internal Trino Iceberg REST catalog, with the + * {@code gravitino.iceberg.rest-catalog.} prefix rewritten to {@code iceberg.rest-catalog.}. + * + * @return the Trino Iceberg REST catalog properties + */ + public Map<String, String> getIcebergRestCatalogConfig() { + String prefix = GRAVITINO_ICEBERG_REST_CATALOG_CONFIG_PREFIX.key; + return config.entrySet().stream() + .filter(entry -> entry.getKey().startsWith(prefix)) + .collect( + Collectors.toMap( + entry -> + TRINO_ICEBERG_REST_CATALOG_PREFIX + entry.getKey().substring(prefix.length()), + Map.Entry::getValue)); + } + private long parseLongConfigEntry(ConfigEntry entry) { String value = config.getOrDefault(entry.key, entry.defaultValue); try { diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/DefaultCatalogConnectorFactory.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/DefaultCatalogConnectorFactory.java index 29b2979da2..66d6937d84 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/DefaultCatalogConnectorFactory.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/DefaultCatalogConnectorFactory.java @@ -72,7 +72,7 @@ public class DefaultCatalogConnectorFactory implements CatalogConnectorFactory { new CatalogConnectorContext.Builder(new MemoryConnectorAdapter())); catalogBuilders.put( ICEBERG_CONNECTOR_PROVIDER_NAME, - new CatalogConnectorContext.Builder(new IcebergConnectorAdapter())); + new CatalogConnectorContext.Builder(new IcebergConnectorAdapter(config))); catalogBuilders.put( MYSQL_CONNECTOR_PROVIDER_NAME, new CatalogConnectorContext.Builder(new MySQLConnectorAdapter())); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java index 7f5460cbf0..8073f1e1c2 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java @@ -22,16 +22,21 @@ package org.apache.gravitino.trino.connector.catalog.iceberg; import com.google.common.collect.Sets; import io.trino.spi.TrinoException; import java.util.HashMap; +import java.util.Locale; import java.util.Map; import java.util.Set; +import org.apache.commons.lang3.StringUtils; import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants; import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergPropertiesUtils; import org.apache.gravitino.credential.AzureAccountKeyCredential; import org.apache.gravitino.credential.Credential; import org.apache.gravitino.credential.JdbcCredential; import org.apache.gravitino.credential.S3SecretKeyCredential; +import org.apache.gravitino.storage.S3Properties; +import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.catalog.CatalogPropertyConverter; +import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; /** * A property converter for Iceberg catalogs that handles the conversion between Trino and Gravitino @@ -46,6 +51,21 @@ public class IcebergCatalogPropertyConverter extends CatalogPropertyConverter { private static final Set<String> REST_BACKEND_REQUIRED_PROPERTIES = Set.of("uri"); + private static final String TRINO_ICEBERG_CATALOG_TYPE = "iceberg.catalog.type"; + private static final String TRINO_ICEBERG_CATALOG_TYPE_REST = "rest"; + private static final String TRINO_ICEBERG_REST_URI = "iceberg.rest-catalog.uri"; + private static final String TRINO_ICEBERG_REST_PREFIX = "iceberg.rest-catalog.prefix"; + private static final String TRINO_ICEBERG_REST_WAREHOUSE = "iceberg.rest-catalog.warehouse"; + private static final String TRINO_ICEBERG_REST_VENDED_CREDENTIALS = + "iceberg.rest-catalog.vended-credentials-enabled"; + private static final String TRINO_ICEBERG_REST_SESSION = "iceberg.rest-catalog.session"; + private static final String TRINO_FS_HADOOP_ENABLED = "fs.hadoop.enabled"; + private static final String TRINO_FS_NATIVE_S3_ENABLED = "fs.native-s3.enabled"; + private static final String TRINO_FS_NATIVE_GCS_ENABLED = "fs.native-gcs.enabled"; + private static final String TRINO_FS_NATIVE_AZURE_ENABLED = "fs.native-azure.enabled"; + private static final String TRINO_S3_REGION = "s3.region"; + private static final String TRINO_S3_ENDPOINT = "s3.endpoint"; + /** * Injects credentials from credential vending into the Iceberg catalog config. Applies JDBC * user/password for the JDBC backend and S3 credentials for S3-backed storage. @@ -147,6 +167,102 @@ public class IcebergCatalogPropertyConverter extends CatalogPropertyConverter { return jdbcProperties; } + /** + * Builds the Trino Iceberg connector config that reaches this catalog through the Gravitino + * Iceberg REST server, regardless of the catalog backend Gravitino uses to store its metadata. + * + * <p>This is the only path on which credential vending works: the Iceberg REST protocol issues a + * fresh temporary credential per table access, while Trino's jdbc and hive_metastore Iceberg + * catalog types have nowhere to put the session token of an STS credential. + * + * @param catalog the Gravitino catalog to load + * @param gravitinoConfig the connector configuration holding the Iceberg REST server endpoint + * @return the Trino Iceberg connector config + */ + public Map<String, String> buildIcebergRestProperties( + GravitinoCatalog catalog, GravitinoConfig gravitinoConfig) { + String restUri = gravitinoConfig.getIcebergRestUri(); + if (StringUtils.isBlank(restUri)) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_MISSING_CONFIG, + "Missing required config 'gravitino.iceberg.rest-uri'. Set it to the Gravitino Iceberg " + + "REST server endpoint, for example http://localhost:9001/iceberg, or set " + + "'gravitino.iceberg.rest-enabled=false' to load Iceberg catalogs through their " + + "catalog backend instead."); + } + + Map<String, String> config = new HashMap<>(); + // The order of put operations determines the priority of parameters. + config.putAll(buildStorageProperties(catalog.getProperties())); + config.put(TRINO_ICEBERG_REST_VENDED_CREDENTIALS, "true"); + if (gravitinoConfig.isForwardUser()) { + config.put(TRINO_ICEBERG_REST_SESSION, "USER"); + } + // The catalog's own trino.bypass properties override the defaults above, so that a Trino + // release renaming one of them can be worked around without a connector change. + config.putAll(super.gravitinoToEngineProperties(catalog.getProperties())); + // The Iceberg REST server endpoint and its authentication are cluster-level operational + // settings, so they take precedence over anything set on a single catalog. + config.putAll(gravitinoConfig.getIcebergRestCatalogConfig()); + config.put(TRINO_ICEBERG_CATALOG_TYPE, TRINO_ICEBERG_CATALOG_TYPE_REST); + config.put(TRINO_ICEBERG_REST_URI, restUri); + // The Iceberg client selects the catalog twice over: `warehouse` is the query parameter of the + // initial GET /v1/config that discovers the catalog, and `prefix` is the path segment of every + // request after it. Setting only the latter leaves the config call without a catalog name. + config.put(TRINO_ICEBERG_REST_WAREHOUSE, catalog.getName()); + config.put(TRINO_ICEBERG_REST_PREFIX, catalog.getName()); + return config; + } + + /** + * Derives the Trino file system config from the catalog's warehouse location. Vended credentials + * are only consumed by Trino's native file systems, so the native implementation matching the + * warehouse scheme has to be enabled. + */ + private Map<String, String> buildStorageProperties(Map<String, String> properties) { + Map<String, String> storageProperties = new HashMap<>(); + storageProperties.put(TRINO_FS_HADOOP_ENABLED, "true"); + + String warehouse = properties.get(IcebergConstants.WAREHOUSE); + if (StringUtils.isBlank(warehouse)) { + return storageProperties; + } + + String scheme = StringUtils.substringBefore(warehouse, "://").toLowerCase(Locale.ROOT); + switch (scheme) { + case "s3": + case "s3a": + case "s3n": + storageProperties.put(TRINO_FS_NATIVE_S3_ENABLED, "true"); + copyProperty( + properties, S3Properties.GRAVITINO_S3_REGION, storageProperties, TRINO_S3_REGION); + copyProperty( + properties, S3Properties.GRAVITINO_S3_ENDPOINT, storageProperties, TRINO_S3_ENDPOINT); + break; + case "gs": + storageProperties.put(TRINO_FS_NATIVE_GCS_ENABLED, "true"); + break; + case "abfs": + case "abfss": + case "wasb": + case "wasbs": + storageProperties.put(TRINO_FS_NATIVE_AZURE_ENABLED, "true"); + break; + default: + // hdfs, file and any scheme without a Trino native file system stay on fs.hadoop.enabled. + break; + } + return storageProperties; + } + + private void copyProperty( + Map<String, String> source, String sourceKey, Map<String, String> target, String targetKey) { + String value = source.get(sourceKey); + if (StringUtils.isNotBlank(value)) { + target.put(targetKey, value); + } + } + private Map<String, String> buildRestBackendProperties(Map<String, String> properties) { Set<String> missingProperty = Sets.difference(REST_BACKEND_REQUIRED_PROPERTIES, properties.keySet()); @@ -157,8 +273,8 @@ public class IcebergCatalogPropertyConverter extends CatalogPropertyConverter { } Map<String, String> restProperties = new HashMap<>(); - restProperties.put("iceberg.catalog.type", "rest"); - restProperties.put("iceberg.rest-catalog.uri", properties.get(IcebergConstants.URI)); + restProperties.put(TRINO_ICEBERG_CATALOG_TYPE, TRINO_ICEBERG_CATALOG_TYPE_REST); + restProperties.put(TRINO_ICEBERG_REST_URI, properties.get(IcebergConstants.URI)); if (properties.containsKey(IcebergConstants.WAREHOUSE)) { restProperties.put( "iceberg.rest-catalog.warehouse", properties.get(IcebergConstants.WAREHOUSE)); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java index 8900775fb7..102c15d0f5 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java @@ -24,8 +24,9 @@ import io.trino.spi.session.PropertyMetadata; import java.util.HashMap; import java.util.List; import java.util.Map; -import org.apache.gravitino.catalog.property.PropertyConverter; +import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants; import org.apache.gravitino.credential.Credential; +import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorAdapter; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadataAdapter; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; @@ -37,25 +38,41 @@ import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; public class IcebergConnectorAdapter implements CatalogConnectorAdapter { private static final String CONNECTOR_ICEBERG = "iceberg"; + private static final String REST_CATALOG_BACKEND = "rest"; + private final IcebergPropertyMeta propertyMetadata; - private final PropertyConverter catalogConverter; + private final IcebergCatalogPropertyConverter catalogConverter; + private final GravitinoConfig config; /** * Constructs a new IcebergConnectorAdapter. Initializes the property metadata and catalog * converter for handling Iceberg-specific configurations. + * + * @param config the Gravitino connector configuration */ - public IcebergConnectorAdapter() { + public IcebergConnectorAdapter(GravitinoConfig config) { this.propertyMetadata = new IcebergPropertyMeta(); this.catalogConverter = new IcebergCatalogPropertyConverter(); + this.config = config; } @Override public Map<String, String> buildInternalConnectorConfig( GravitinoCatalog catalog, Credential[] credentials) throws Exception { - Map<String, String> config = + // The catalog backend describes how Gravitino stores the metadata; it does not decide how + // Trino reaches the data. Unless the routing is disabled, every catalog is loaded through the + // Gravitino Iceberg REST server, which is the only path that supports credential vending. + // A catalog that already has a REST backend keeps pointing at its own configured endpoint. + if (config.isIcebergRestEnabled() + && !REST_CATALOG_BACKEND.equals( + catalog.getProperty(IcebergConstants.CATALOG_BACKEND, null))) { + return catalogConverter.buildIcebergRestProperties(catalog, config); + } + + Map<String, String> connectorConfig = new HashMap<>(catalogConverter.gravitinoToEngineProperties(catalog.getProperties())); - IcebergCatalogPropertyConverter.applyCredentials(credentials, config); - return config; + IcebergCatalogPropertyConverter.applyCredentials(credentials, connectorConfig); + return connectorConfig; } @Override diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java index 2d486b52c0..c31c594be0 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java @@ -300,6 +300,59 @@ public class TestGravitinoConfig { assertTrue(config.isTrinoJdbcSslEnabled()); } + @Test + public void testIcebergRestConfigDefaults() { + GravitinoConfig config = new GravitinoConfig(ImmutableMap.of("gravitino.metalake", "user_001")); + + assertTrue(config.isIcebergRestEnabled()); + assertEquals("", config.getIcebergRestUri()); + assertTrue(config.getIcebergRestCatalogConfig().isEmpty()); + } + + @Test + public void testIcebergRestConfig() { + ImmutableMap<String, String> configMap = + ImmutableMap.of( + "gravitino.metalake", + "user_001", + "gravitino.iceberg.rest-enabled", + "false", + "gravitino.iceberg.rest-uri", + "http://127.0.0.1:9001/iceberg", + "gravitino.iceberg.rest-catalog.security", + "OAUTH2", + "gravitino.iceberg.rest-catalog.oauth2.credential", + "client_id:client_secret"); + GravitinoConfig config = new GravitinoConfig(configMap); + + assertFalse(config.isIcebergRestEnabled()); + assertEquals("http://127.0.0.1:9001/iceberg", config.getIcebergRestUri()); + + Map<String, String> restCatalogConfig = config.getIcebergRestCatalogConfig(); + assertEquals(2, restCatalogConfig.size()); + assertEquals("OAUTH2", restCatalogConfig.get("iceberg.rest-catalog.security")); + assertEquals( + "client_id:client_secret", restCatalogConfig.get("iceberg.rest-catalog.oauth2.credential")); + } + + @Test + public void testToCatalogConfigWithIcebergRestProperties() { + ImmutableMap<String, String> configMap = + ImmutableMap.of( + "gravitino.metalake", + "user_001", + "gravitino.iceberg.rest-uri", + "http://127.0.0.1:9001/iceberg", + "gravitino.iceberg.rest-catalog.security", + "OAUTH2"); + GravitinoConfig config = new GravitinoConfig(configMap); + + String catalogConfig = config.toCatalogConfig(); + assertTrue( + catalogConfig.contains("\"gravitino.iceberg.rest-uri\"='http://127.0.0.1:9001/iceberg'")); + assertTrue(catalogConfig.contains("\"gravitino.iceberg.rest-catalog.security\"='OAUTH2'")); + } + private static boolean skipCatalog(String catalogName, GravitinoConfig config) { for (Pattern pattern : config.getSkipCatalogPatterns()) { if (pattern.matcher(catalogName).matches()) { diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java index e16afd8e6f..fc3571c8e8 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java @@ -27,6 +27,7 @@ import org.apache.gravitino.Catalog; import org.apache.gravitino.catalog.property.PropertyConverter; import org.apache.gravitino.credential.Credential; import org.apache.gravitino.credential.JdbcCredential; +import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; import org.apache.gravitino.trino.connector.metadata.TestGravitinoCatalog; import org.junit.jupiter.api.Assertions; @@ -136,7 +137,7 @@ public class TestIcebergCatalogPropertyConverter { Catalog mockCatalog = TestGravitinoCatalog.mockCatalog( name, "lakehouse-iceberg", "test catalog", Catalog.Type.RELATIONAL, properties); - IcebergConnectorAdapter adapter = new IcebergConnectorAdapter(); + IcebergConnectorAdapter adapter = new IcebergConnectorAdapter(icebergRestDisabledConfig()); Map<String, String> config = adapter.buildInternalConnectorConfig( @@ -176,7 +177,7 @@ public class TestIcebergCatalogPropertyConverter { Catalog mockCatalog = TestGravitinoCatalog.mockCatalog( name, "lakehouse-iceberg", "test catalog", Catalog.Type.RELATIONAL, properties); - IcebergConnectorAdapter adapter = new IcebergConnectorAdapter(); + IcebergConnectorAdapter adapter = new IcebergConnectorAdapter(icebergRestDisabledConfig()); Map<String, String> config = adapter.buildInternalConnectorConfig( @@ -213,7 +214,7 @@ public class TestIcebergCatalogPropertyConverter { Catalog mockCatalog = TestGravitinoCatalog.mockCatalog( name, "lakehouse-iceberg", "test catalog", Catalog.Type.RELATIONAL, properties); - IcebergConnectorAdapter adapter = new IcebergConnectorAdapter(); + IcebergConnectorAdapter adapter = new IcebergConnectorAdapter(icebergRestDisabledConfig()); Map<String, String> config = adapter.buildInternalConnectorConfig( @@ -224,4 +225,250 @@ public class TestIcebergCatalogPropertyConverter { Assertions.assertEquals(config.get("iceberg.jdbc-catalog.connection-user"), "root"); Assertions.assertEquals(config.get("iceberg.jdbc-catalog.connection-password"), "ds123"); } + + @Test + public void testBuildConnectorPropertiesRoutesJdbcBackendThroughIcebergRest() throws Exception { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("uri", "jdbc:postgresql://localhost:5432/iceberg") + .put("catalog-backend", "jdbc") + .put("jdbc-driver", "org.postgresql.Driver") + .put("jdbc-user", "iceberg") + .put("jdbc-password", "secret") + .put("warehouse", "s3://bucket/warehouse/") + .put("s3-region", "us-east-1") + .put("credential-providers", "s3-token") + .build(); + + Map<String, String> config = + buildConnectorConfig("catalog1", properties, icebergRestEnabledConfig(ImmutableMap.of())); + + Assertions.assertEquals("rest", config.get("iceberg.catalog.type")); + Assertions.assertEquals( + "http://localhost:9001/iceberg", config.get("iceberg.rest-catalog.uri")); + Assertions.assertEquals("catalog1", config.get("iceberg.rest-catalog.prefix")); + // The initial GET /v1/config selects the catalog by `warehouse`, not by `prefix`. + Assertions.assertEquals("catalog1", config.get("iceberg.rest-catalog.warehouse")); + Assertions.assertEquals("true", config.get("iceberg.rest-catalog.vended-credentials-enabled")); + + // The JDBC backend is how Gravitino stores the metadata; it must not reach Trino. + Assertions.assertNull(config.get("iceberg.jdbc-catalog.connection-url")); + Assertions.assertNull(config.get("iceberg.jdbc-catalog.connection-user")); + Assertions.assertNull(config.get("iceberg.jdbc-catalog.connection-password")); + } + + @Test + public void testBuildConnectorPropertiesRoutesHiveBackendThroughIcebergRest() throws Exception { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("uri", "thrift://localhost:9083") + .put("catalog-backend", "hive") + .put("warehouse", "s3://bucket/warehouse/") + .build(); + + Map<String, String> config = + buildConnectorConfig("catalog1", properties, icebergRestEnabledConfig(ImmutableMap.of())); + + Assertions.assertEquals("rest", config.get("iceberg.catalog.type")); + Assertions.assertEquals("catalog1", config.get("iceberg.rest-catalog.prefix")); + Assertions.assertNull(config.get("hive.metastore.uri")); + } + + @Test + public void testBuildConnectorPropertiesKeepsRestBackendEndpoint() throws Exception { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("uri", "http://other-irc:9001/iceberg") + .put("catalog-backend", "rest") + .put("warehouse", "gt_iceberg_rest") + .build(); + + Map<String, String> config = + buildConnectorConfig("catalog1", properties, icebergRestEnabledConfig(ImmutableMap.of())); + + Assertions.assertEquals("rest", config.get("iceberg.catalog.type")); + Assertions.assertEquals( + "http://other-irc:9001/iceberg", config.get("iceberg.rest-catalog.uri")); + Assertions.assertEquals("gt_iceberg_rest", config.get("iceberg.rest-catalog.warehouse")); + Assertions.assertNull(config.get("iceberg.rest-catalog.prefix")); + } + + @Test + public void testBuildConnectorPropertiesWithIcebergRestDisabled() throws Exception { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("uri", "jdbc:postgresql://localhost:5432/iceberg") + .put("catalog-backend", "jdbc") + .put("jdbc-driver", "org.postgresql.Driver") + .put("warehouse", "s3://bucket/warehouse/") + .build(); + + Map<String, String> config = + buildConnectorConfig("catalog1", properties, icebergRestDisabledConfig()); + + Assertions.assertEquals("jdbc", config.get("iceberg.catalog.type")); + Assertions.assertEquals( + "jdbc:postgresql://localhost:5432/iceberg", + config.get("iceberg.jdbc-catalog.connection-url")); + Assertions.assertNull(config.get("iceberg.rest-catalog.uri")); + } + + @Test + public void testBuildConnectorPropertiesMissingIcebergRestUri() { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("uri", "jdbc:postgresql://localhost:5432/iceberg") + .put("catalog-backend", "jdbc") + .put("jdbc-driver", "org.postgresql.Driver") + .build(); + GravitinoConfig gravitinoConfig = + new GravitinoConfig(ImmutableMap.of("gravitino.metalake", "test")); + + Assertions.assertThrows( + TrinoException.class, () -> buildConnectorConfig("catalog1", properties, gravitinoConfig)); + } + + @Test + public void testBuildConnectorPropertiesWithIcebergRestAuthentication() throws Exception { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("catalog-backend", "jdbc") + .put("uri", "jdbc:postgresql://localhost:5432/iceberg") + .put("jdbc-driver", "org.postgresql.Driver") + .put("warehouse", "s3://bucket/warehouse/") + .build(); + GravitinoConfig gravitinoConfig = + icebergRestEnabledConfig( + ImmutableMap.of( + "gravitino.iceberg.rest-catalog.security", "OAUTH2", + "gravitino.iceberg.rest-catalog.oauth2.credential", "client_id:client_secret", + "gravitino.iceberg.rest-catalog.uri", "http://ignored:9001/iceberg")); + + Map<String, String> config = buildConnectorConfig("catalog1", properties, gravitinoConfig); + + Assertions.assertEquals("OAUTH2", config.get("iceberg.rest-catalog.security")); + Assertions.assertEquals( + "client_id:client_secret", config.get("iceberg.rest-catalog.oauth2.credential")); + // The endpoint the connector routes to cannot be overridden by the pass-through prefix. + Assertions.assertEquals( + "http://localhost:9001/iceberg", config.get("iceberg.rest-catalog.uri")); + } + + @Test + public void testBuildConnectorPropertiesForwardsSessionUser() throws Exception { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("catalog-backend", "jdbc") + .put("uri", "jdbc:postgresql://localhost:5432/iceberg") + .put("jdbc-driver", "org.postgresql.Driver") + .build(); + + Map<String, String> config = + buildConnectorConfig( + "catalog1", + properties, + icebergRestEnabledConfig( + ImmutableMap.of("gravitino.client.session.forwardUser", "true"))); + Assertions.assertEquals("USER", config.get("iceberg.rest-catalog.session")); + + Map<String, String> explicitConfig = + buildConnectorConfig( + "catalog1", + properties, + icebergRestEnabledConfig( + ImmutableMap.of( + "gravitino.client.session.forwardUser", "true", + "gravitino.iceberg.rest-catalog.session", "NONE"))); + Assertions.assertEquals("NONE", explicitConfig.get("iceberg.rest-catalog.session")); + + Map<String, String> defaultConfig = + buildConnectorConfig("catalog1", properties, icebergRestEnabledConfig(ImmutableMap.of())); + Assertions.assertNull(defaultConfig.get("iceberg.rest-catalog.session")); + } + + @Test + public void testBuildConnectorPropertiesStorageDetection() throws Exception { + Map<String, String> s3Config = + buildConnectorConfig( + "catalog1", + ImmutableMap.of( + "catalog-backend", "jdbc", + "warehouse", "s3://bucket/warehouse/", + "s3-region", "us-east-1", + "s3-endpoint", "http://minio:9000"), + icebergRestEnabledConfig(ImmutableMap.of())); + Assertions.assertEquals("true", s3Config.get("fs.native-s3.enabled")); + Assertions.assertEquals("us-east-1", s3Config.get("s3.region")); + Assertions.assertEquals("http://minio:9000", s3Config.get("s3.endpoint")); + Assertions.assertEquals("true", s3Config.get("fs.hadoop.enabled")); + + Map<String, String> gcsConfig = + buildConnectorConfig( + "catalog1", + ImmutableMap.of("catalog-backend", "jdbc", "warehouse", "gs://bucket/warehouse/"), + icebergRestEnabledConfig(ImmutableMap.of())); + Assertions.assertEquals("true", gcsConfig.get("fs.native-gcs.enabled")); + Assertions.assertNull(gcsConfig.get("fs.native-s3.enabled")); + + Map<String, String> azureConfig = + buildConnectorConfig( + "catalog1", + ImmutableMap.of( + "catalog-backend", "jdbc", "warehouse", "abfss://container@account/warehouse/"), + icebergRestEnabledConfig(ImmutableMap.of())); + Assertions.assertEquals("true", azureConfig.get("fs.native-azure.enabled")); + + Map<String, String> hdfsConfig = + buildConnectorConfig( + "catalog1", + ImmutableMap.of("catalog-backend", "jdbc", "warehouse", "hdfs://namenode:9000/wh"), + icebergRestEnabledConfig(ImmutableMap.of())); + Assertions.assertEquals("true", hdfsConfig.get("fs.hadoop.enabled")); + Assertions.assertNull(hdfsConfig.get("fs.native-s3.enabled")); + Assertions.assertNull(hdfsConfig.get("fs.native-gcs.enabled")); + Assertions.assertNull(hdfsConfig.get("fs.native-azure.enabled")); + } + + @Test + public void testBuildConnectorPropertiesIcebergRestKeepsCatalogBypass() throws Exception { + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("catalog-backend", "jdbc") + .put("warehouse", "s3://bucket/warehouse/") + .put("trino.bypass.iceberg.table-statistics-enabled", "true") + .put("trino.bypass.fs.native-s3.enabled", "false") + .build(); + + Map<String, String> config = + buildConnectorConfig("catalog1", properties, icebergRestEnabledConfig(ImmutableMap.of())); + + Assertions.assertEquals("true", config.get("iceberg.table-statistics-enabled")); + // A catalog can override a derived default, so a Trino release renaming it is not a blocker. + Assertions.assertEquals("false", config.get("fs.native-s3.enabled")); + } + + @SuppressWarnings("unchecked") + private static Map<String, String> buildConnectorConfig( + String catalogName, Map<String, String> properties, GravitinoConfig gravitinoConfig) + throws Exception { + Catalog mockCatalog = + TestGravitinoCatalog.mockCatalog( + catalogName, "lakehouse-iceberg", "test catalog", Catalog.Type.RELATIONAL, properties); + return new IcebergConnectorAdapter(gravitinoConfig) + .buildInternalConnectorConfig(new GravitinoCatalog("test", mockCatalog), new Credential[0]); + } + + private static GravitinoConfig icebergRestDisabledConfig() { + return new GravitinoConfig( + ImmutableMap.of("gravitino.metalake", "test", "gravitino.iceberg.rest-enabled", "false")); + } + + private static GravitinoConfig icebergRestEnabledConfig(Map<String, String> extraConfig) { + return new GravitinoConfig( + ImmutableMap.<String, String>builder() + .put("gravitino.metalake", "test") + .put("gravitino.iceberg.rest-uri", "http://localhost:9001/iceberg") + .putAll(extraConfig) + .build()); + } }
