This is an automated email from the ASF dual-hosted git repository.
diqiu50 pushed a commit to branch trino-irc-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/trino-irc-1.3 by this push:
new 087a95bff8 fix(trino): require Iceberg REST routing by default
087a95bff8 is described below
commit 087a95bff8e43b441bfd09a291234c206091faf1
Author: yuhui <[email protected]>
AuthorDate: Wed Aug 26 14:06:18 2026 +0000
fix(trino): require Iceberg REST routing by default
---
docs/trino-connector/catalog-iceberg.md | 34 +++++-----
docs/trino-connector/configuration.md | 11 ++--
.../gravitino/trino/connector/GravitinoConfig.java | 22 +++++++
.../connector/catalog/CatalogConnectorManager.java | 32 +++++----
.../catalog/iceberg/IcebergConnectorAdapter.java | 21 +++++-
.../trino/connector/TestGravitinoConfig.java | 25 +++++++
.../catalog/TestCatalogConnectorManager.java | 76 ++++++++++++++++++++++
.../TestIcebergCatalogPropertyConverter.java | 59 ++++++++---------
8 files changed, 209 insertions(+), 71 deletions(-)
diff --git a/docs/trino-connector/catalog-iceberg.md
b/docs/trino-connector/catalog-iceberg.md
index 5a214fc624..f457c5b3b0 100644
--- a/docs/trino-connector/catalog-iceberg.md
+++ b/docs/trino-connector/catalog-iceberg.md
@@ -35,19 +35,18 @@ protocol.
The connector already connects to the Gravitino server (it is how catalogs are
discovered in the
first place), so it also asks that server whether it has an Iceberg REST
server running as an
-[auxiliary service](../iceberg-rest-service.md) for the connector's metalake,
and uses the endpoint
-the server reports. Nothing needs to be configured for this:
`etc/catalog/gravitino.properties` only
-needs the usual `gravitino.uri` and `gravitino.metalake`. If the server
reports no endpoint — the
-IRC is not running as an auxiliary service, or it serves a different metalake
— `lakehouse-iceberg`
-catalogs fall back to translating `catalog-backend` as before, and credential
vending does not work.
-
-Only the coordinator polls the Gravitino server, so the coordinator resolves
the endpoint once, when
-a catalog is first registered or reloaded, and hands it to every node
(coordinator and workers alike)
+[auxiliary service](../iceberg-rest-service.md) for the connector's metalake.
By default, a
+non-REST `lakehouse-iceberg` catalog is not registered until an endpoint is
discovered or configured
+explicitly. It is retried during every metadata refresh rather than silently
falling back and
+disabling credential vending. To retain the behavior from older connector
versions, set
+`gravitino.iceberg.rest-routing-enabled=false`; this skips discovery and
translates the catalog's
+`catalog-backend` into the corresponding native Trino Iceberg configuration.
+
+Only the coordinator polls the Gravitino server, so the coordinator resolves
the endpoint when a
+catalog is registered or refreshed and hands it to every node (coordinator and
workers alike)
as part of that catalog's own definition — the same way Trino replicates any
other catalog property
-cluster-wide. A practical consequence: a catalog registered before the IRC
started keeps its existing
-routing until Gravitino reports a change *and* the catalog itself is reloaded
(its metadata changes,
-or Trino restarts) — starting the IRC alone does not retroactively re-route an
already-registered
-catalog.
+cluster-wide. A catalog that could not be registered before the IRC started is
registered
+automatically after a later discovery poll succeeds; no Trino restart is
required.
Set `gravitino.iceberg.rest-uri` to override the discovered endpoint, and it
is required — not just
an override — for a standalone IRC (its own process, not the Gravitino
server's auxiliary service):
@@ -122,13 +121,14 @@ keeping per-user credential vending and per-user
authorization intact. Set
- One IRC serves exactly one metalake, fixed at startup by
`gravitino.iceberg-rest.gravitino-metalake`. The Gravitino server only
reports the IRC's endpoint
- for that metalake; in multi-metalake mode
(`gravitino.use-single-metalake=false`), catalogs in any
- other metalake fall back to translating `catalog-backend` instead of failing
outright.
+ for that metalake. In multi-metalake mode
(`gravitino.use-single-metalake=false`), a non-REST
+ Iceberg catalog in another metalake therefore requires a metalake-scoped
manual URI or remains
+ unregistered while REST routing is enabled.
- 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.
-- A deployment that does not run the IRC needs no configuration at all: the
server reports no
- endpoint, and `lakehouse-iceberg` catalogs are translated into Trino's
`jdbc` or `hive_metastore`
- catalog types as before.
+- A deployment that does not run the IRC must set
+ `gravitino.iceberg.rest-routing-enabled=false` to translate non-REST
`lakehouse-iceberg` catalogs
+ into Trino's `jdbc` or `hive_metastore` catalog types as before.
- Discovery only works for an IRC running as a Gravitino auxiliary service
(`gravitino.auxService.names=iceberg-rest`), embedded in the same process as
the Gravitino server.
A standalone IRC — its own process, started with
`GravitinoIcebergRESTServer` and its own
diff --git a/docs/trino-connector/configuration.md
b/docs/trino-connector/configuration.md
index 7e6afd07f7..0c4a2ed96a 100644
--- a/docs/trino-connector/configuration.md
+++ b/docs/trino-connector/configuration.md
@@ -29,16 +29,19 @@ 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-uri | string | (none)
| The endpoint of the Gravitino Iceberg REST server (IRC). Discovered
automatically from the Gravitino server, which is asked whether it has an IRC
running for this connector's metalake; set this only to override the discovered
value. When an endpoint is available (discovered or configured),
`lakehouse-iceberg` catalogs whose `catalog-backend` is not `rest` are loaded
through it instead of being translated in [...]
+| gravitino.iceberg.rest-routing-enabled | boolean | true
| Whether non-REST `lakehouse-iceberg` catalogs must be routed through the
Gravitino Iceberg REST server. When enabled, a catalog is not registered until
discovery succeeds or `gravitino.iceberg.rest-uri` is configured. Set this to
`false` to retain legacy `catalog-backend` translation and skip discovery. | No
|
+| gravitino.iceberg.rest-uri | string | (none)
| An explicitly configured endpoint for the Gravitino Iceberg REST server
(IRC). It takes precedence over automatic discovery and is required when the
server cannot advertise the endpoint, such as with a standalone IRC. Non-REST
`lakehouse-iceberg` catalogs are loaded through this endpoint instead of being
translated into Trino's `jdbc` or `hive_metastore` Iceberg catalog type,
enabling credential vending. | [...]
| 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`. The `uri`, `warehouse` and `prefix`
keys are reserved and always derived by the connector. | 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.
-Upgrading an existing deployment needs no action: a Gravitino server without
the Iceberg REST server
-running reports no endpoint, so `lakehouse-iceberg` catalogs keep translating
`catalog-backend` as
-before. See [Iceberg
catalog](./catalog-iceberg.md#how-trino-reaches-the-catalog).
+When upgrading a deployment that does not provide an Iceberg REST endpoint,
either configure
+`gravitino.iceberg.rest-uri` or set
`gravitino.iceberg.rest-routing-enabled=false` to retain the
+legacy `catalog-backend` translation. Otherwise, non-REST `lakehouse-iceberg`
catalogs remain
+unregistered until discovery succeeds. 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 440-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.
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 cb9d71a3ff..6f3f213418 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
@@ -287,6 +287,14 @@ public class GravitinoConfig {
"",
false);
+ private static final ConfigEntry GRAVITINO_ICEBERG_REST_ROUTING_ENABLED =
+ new ConfigEntry(
+ "gravitino.iceberg.rest-routing-enabled",
+ "Whether non-REST Iceberg catalogs must be routed through the
Gravitino Iceberg REST "
+ + "server. Disable this only to retain the legacy
catalog-backend translation.",
+ "true",
+ false);
+
private static final ConfigEntry
GRAVITINO_ICEBERG_REST_CATALOG_CONFIG_PREFIX =
new ConfigEntry(
"gravitino.iceberg.rest-catalog.",
@@ -770,6 +778,20 @@ public class GravitinoConfig {
return discoveredIcebergRestUriByMetalake.getOrDefault(metalake, "");
}
+ /**
+ * Returns whether non-REST Iceberg catalogs must be routed through the
Gravitino Iceberg REST
+ * server.
+ *
+ * @return {@code true} when Iceberg REST routing is enabled
+ */
+ public boolean isIcebergRestRoutingEnabled() {
+ String value =
+ config.getOrDefault(
+ GRAVITINO_ICEBERG_REST_ROUTING_ENABLED.key,
+ GRAVITINO_ICEBERG_REST_ROUTING_ENABLED.defaultValue);
+ return parseBooleanConfig(GRAVITINO_ICEBERG_REST_ROUTING_ENABLED.key,
value);
+ }
+
/**
* Retrieves the manually configured Iceberg REST server endpoint for the
given metalake, if any.
* Unlike the discovered endpoint, this is plain local file configuration
and is therefore
diff --git
a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java
index b1a809ca7b..380c5214f8 100644
---
a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java
+++
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java
@@ -201,7 +201,10 @@ public class CatalogConnectorManager {
try {
GravitinoMetalake metalake = metalakes.get(usedMetalake);
LOG.debug("Load metalake: {}", usedMetalake);
- refreshIcebergRestUri(usedMetalake);
+ if (config.isIcebergRestRoutingEnabled()
+ &&
StringUtils.isBlank(config.getManualIcebergRestUri(usedMetalake))) {
+ refreshIcebergRestUri(usedMetalake);
+ }
loadCatalogs(metalake);
} catch (Exception e) {
LOG.error("Load Metalake {} failed.", usedMetalake, e);
@@ -218,10 +221,8 @@ public class CatalogConnectorManager {
* read on the next catalog load. Failures — including talking to a
Gravitino server older than
* this endpoint — must not interrupt catalog loading, so they are swallowed
here; Iceberg
* catalogs simply keep their last known routing decision until the next
successful poll. A
- * failure is logged at WARN only on the transition into that state (not on
every poll), so a
- * persistently failing/misconfigured discovery is still visible at the
default log level without
- * spamming the log for the common case of an older Gravitino server that
lacks this endpoint
- * entirely.
+ * failure is logged at ERROR on every poll because routing through Iceberg
REST is required when
+ * enabled. Catalog loading continues so that unrelated catalogs remain
available.
*/
private void refreshIcebergRestUri(String metalakeName) {
try {
@@ -231,18 +232,15 @@ public class CatalogConnectorManager {
LOG.info("Iceberg REST service discovery for metalake {} recovered.",
metalakeName);
}
} catch (Exception e) {
- if (icebergRestDiscoveryFailing.add(metalakeName)) {
- LOG.warn(
- "Failed to query the Iceberg REST service endpoint for metalake
{}; keeping the "
- + "last known routing decision. This is expected when talking
to a Gravitino "
- + "server that predates this endpoint, but is otherwise worth
investigating. "
- + "Further failures for this metalake are logged at DEBUG
until it recovers.",
- metalakeName,
- e);
- } else {
- LOG.debug(
- "Failed to query the Iceberg REST service endpoint for metalake
{}.", metalakeName, e);
- }
+ icebergRestDiscoveryFailing.add(metalakeName);
+ LOG.error(
+ "Failed to query the Iceberg REST service endpoint for metalake {};
Iceberg catalogs "
+ + "without a configured REST endpoint cannot be registered until
discovery "
+ + "recovers. Set gravitino.iceberg.rest-uri explicitly, upgrade
the Gravitino "
+ + "server to one that supports discovery, or disable Iceberg
REST routing with "
+ + "gravitino.iceberg.rest-routing-enabled=false to use legacy
backend translation.",
+ metalakeName,
+ e);
}
}
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 44c06503f8..29536a8d0e 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
@@ -21,6 +21,7 @@ package org.apache.gravitino.trino.connector.catalog.iceberg;
import static java.util.Collections.emptyList;
import com.google.common.collect.ImmutableMap;
+import io.trino.spi.TrinoException;
import io.trino.spi.session.PropertyMetadata;
import java.util.HashMap;
import java.util.List;
@@ -29,6 +30,7 @@ import org.apache.commons.lang3.StringUtils;
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.GravitinoErrorCode;
import org.apache.gravitino.trino.connector.catalog.CatalogConnectorAdapter;
import
org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadataAdapter;
import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog;
@@ -81,8 +83,8 @@ public class IcebergConnectorAdapter implements
CatalogConnectorAdapter {
// Trino reaches the data. Whenever an Iceberg REST server endpoint is
available for this
// catalog's metalake, the catalog is loaded through it, the only path
that supports temporary
// credentials. A catalog that already has a REST backend keeps pointing
at its own configured
- // endpoint. If no endpoint is available, this falls back to translating
catalog-backend as
- // before — nothing to configure either way.
+ // endpoint. If no endpoint is available, routing is required unless the
compatibility switch
+ // explicitly enables the legacy catalog-backend translation.
//
// The manual override is plain local config, so it is valid on every node
as-is. The
// discovered endpoint is coordinator-only knowledge, so it is read from
the catalog's own
@@ -106,6 +108,21 @@ public class IcebergConnectorAdapter implements
CatalogConnectorAdapter {
return catalogConverter.buildIcebergRestProperties(catalog, config,
restUri);
}
+ String catalogBackend =
catalog.getProperty(IcebergConstants.CATALOG_BACKEND, null);
+ if (config.isIcebergRestRoutingEnabled()
+ && StringUtils.isBlank(restUri)
+ && !REST_CATALOG_BACKEND.equalsIgnoreCase(catalogBackend)) {
+ throw new TrinoException(
+ GravitinoErrorCode.GRAVITINO_RUNTIME_ERROR,
+ String.format(
+ "Cannot register Iceberg catalog '%s' in metalake '%s': no
Iceberg REST endpoint "
+ + "is available. Configure gravitino.iceberg.rest-uri, use a
Gravitino server "
+ + "that supports /api/system/iceberg-rest, or set "
+ + "gravitino.iceberg.rest-routing-enabled=false to use
legacy backend "
+ + "translation.",
+ catalog.getName(), catalog.getMetalake()));
+ }
+
Map<String, String> connectorConfig =
new
HashMap<>(catalogConverter.gravitinoToEngineProperties(catalog.getProperties()));
IcebergCatalogPropertyConverter.applyCredentials(credentials,
connectorConfig);
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 1f67109deb..b7ec6b8ff9 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
@@ -307,9 +307,34 @@ public class TestGravitinoConfig {
// Nothing configured, and nothing discovered yet.
assertEquals("", config.getManualIcebergRestUri("user_001"));
assertEquals("", config.getDiscoveredIcebergRestUri("user_001"));
+ assertTrue(config.isIcebergRestRoutingEnabled());
assertTrue(config.getIcebergRestCatalogConfig().isEmpty());
}
+ @Test
+ public void testIcebergRestRoutingCanBeDisabled() {
+ GravitinoConfig config =
+ new GravitinoConfig(
+ ImmutableMap.of(
+ "gravitino.metalake", "user_001",
+ "gravitino.iceberg.rest-routing-enabled", "false"));
+
+ assertFalse(config.isIcebergRestRoutingEnabled());
+ assertTrue(
+
config.toCatalogConfig().contains("\"gravitino.iceberg.rest-routing-enabled\"='false'"));
+ }
+
+ @Test
+ public void testIcebergRestRoutingRejectsInvalidBoolean() {
+ GravitinoConfig config =
+ new GravitinoConfig(
+ ImmutableMap.of(
+ "gravitino.metalake", "user_001",
+ "gravitino.iceberg.rest-routing-enabled", "yes"));
+
+ assertThrows(TrinoException.class, config::isIcebergRestRoutingEnabled);
+ }
+
@Test
public void testIcebergRestConfig() {
ImmutableMap<String, String> configMap =
diff --git
a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java
b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java
index 5ac9efb5e7..995aedc29d 100644
---
a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java
+++
b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java
@@ -25,6 +25,9 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import com.google.common.collect.ImmutableMap;
@@ -193,6 +196,79 @@ public class TestCatalogConnectorManager {
assertEquals("", config.getDiscoveredIcebergRestUri("test"));
}
+ @Test
+ public void testIcebergRestRoutingDisabledSkipsDiscovery() throws Exception {
+ GravitinoAdminClient client = mock(GravitinoAdminClient.class);
+ CatalogRegister catalogRegister = mock(CatalogRegister.class);
+ when(catalogRegister.isTrinoStarted()).thenReturn(true);
+
when(client.loadMetalake("test")).thenReturn(mock(GravitinoMetalake.class));
+
+ CatalogConnectorManager manager =
+ new CatalogConnectorManager(catalogRegister,
createCatalogConnectorFactory(), null);
+ GravitinoConfig config =
+ new GravitinoConfig(
+ ImmutableMap.of(
+ "gravitino.uri", "http://127.0.0.1:8090",
+ "gravitino.metalake", "test",
+ "gravitino.use-single-metalake", "true",
+ "gravitino.iceberg.rest-routing-enabled", "false"));
+ manager.config(config, client);
+
+ manager.loadMetalakeSync();
+
+ verify(client, never()).icebergRestServiceUri("test");
+ }
+
+ @Test
+ public void testConfiguredIcebergRestUriSkipsDiscovery() throws Exception {
+ GravitinoAdminClient client = mock(GravitinoAdminClient.class);
+ CatalogRegister catalogRegister = mock(CatalogRegister.class);
+ when(catalogRegister.isTrinoStarted()).thenReturn(true);
+
when(client.loadMetalake("test")).thenReturn(mock(GravitinoMetalake.class));
+
+ CatalogConnectorManager manager =
+ new CatalogConnectorManager(catalogRegister,
createCatalogConnectorFactory(), null);
+ GravitinoConfig config =
+ new GravitinoConfig(
+ ImmutableMap.of(
+ "gravitino.uri", "http://127.0.0.1:8090",
+ "gravitino.metalake", "test",
+ "gravitino.use-single-metalake", "true",
+ "gravitino.iceberg.rest-uri", "http://irc-host:9001/iceberg"));
+ manager.config(config, client);
+
+ manager.loadMetalakeSync();
+
+ verify(client, never()).icebergRestServiceUri("test");
+ }
+
+ @Test
+ public void testIcebergRestDiscoveryRetriesAndRecovers() throws Exception {
+ GravitinoAdminClient client = mock(GravitinoAdminClient.class);
+ CatalogRegister catalogRegister = mock(CatalogRegister.class);
+ when(catalogRegister.isTrinoStarted()).thenReturn(true);
+
when(client.loadMetalake("test")).thenReturn(mock(GravitinoMetalake.class));
+ when(client.icebergRestServiceUri("test"))
+ .thenThrow(new RESTException("simulated discovery failure"))
+ .thenReturn(Optional.of("http://irc-host:9001/iceberg"));
+
+ CatalogConnectorManager manager =
+ new CatalogConnectorManager(catalogRegister,
createCatalogConnectorFactory(), null);
+ GravitinoConfig config =
+ new GravitinoConfig(
+ ImmutableMap.of(
+ "gravitino.uri", "http://127.0.0.1:8090",
+ "gravitino.metalake", "test",
+ "gravitino.use-single-metalake", "true"));
+ manager.config(config, client);
+
+ manager.loadMetalakeSync();
+ manager.loadMetalakeSync();
+
+ verify(client, times(2)).icebergRestServiceUri("test");
+ assertEquals("http://irc-host:9001/iceberg",
config.getDiscoveredIcebergRestUri("test"));
+ }
+
private CatalogConnectorManager createManager(ImmutableMap<String, String>
configMap)
throws Exception {
return createManager(createCatalogConnectorFactory(), configMap);
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 1951003a19..74bea3b161 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
@@ -318,8 +318,7 @@ public class TestIcebergCatalogPropertyConverter {
}
@Test
- public void
testBuildConnectorPropertiesFallsBackWhenNoIcebergRestUriIsAvailable()
- throws Exception {
+ public void
testBuildConnectorPropertiesFailsWhenNoIcebergRestUriIsAvailable() {
Map<String, String> properties =
ImmutableMap.<String, String>builder()
.put("uri", "jdbc:postgresql://localhost:5432/iceberg")
@@ -327,13 +326,12 @@ public class TestIcebergCatalogPropertyConverter {
.put("jdbc-driver", "org.postgresql.Driver")
.build();
- // Neither a manual gravitino.iceberg.rest-uri nor a discovered one for
this metalake: the
- // routing gate never fires, so this is the ordinary backend-translation
path, not an error.
- Map<String, String> config =
- buildConnectorConfig("catalog1", properties,
icebergRestUnavailableConfig());
+ GravitinoConfig config = new
GravitinoConfig(ImmutableMap.of("gravitino.metalake", "test"));
- Assertions.assertEquals("jdbc", config.get("iceberg.catalog.type"));
- Assertions.assertNull(config.get("iceberg.rest-catalog.uri"));
+ TrinoException exception =
+ Assertions.assertThrows(
+ TrinoException.class, () -> buildConnectorConfig("catalog1",
properties, config));
+ Assertions.assertTrue(exception.getMessage().contains("no Iceberg REST
endpoint is available"));
}
@Test
@@ -382,7 +380,7 @@ public class TestIcebergCatalogPropertyConverter {
}
@Test
- public void testBuildConnectorPropertiesIgnoresDiscoveryForOtherMetalakes()
throws Exception {
+ public void testBuildConnectorPropertiesRejectsDiscoveryForOtherMetalakes() {
Map<String, String> properties =
ImmutableMap.<String, String>builder()
.put("catalog-backend", "jdbc")
@@ -392,19 +390,18 @@ public class TestIcebergCatalogPropertyConverter {
// Discovery reported an endpoint for a different metalake than the one
this catalog belongs
// to (the mock catalog's metalake is "test");
embedDiscoveredIcebergRestUri must not embed it.
- Map<String, String> config =
- buildConnectorConfig(
- "catalog1",
- properties,
- icebergRestDiscoveredConfig("other_metalake",
"http://other:9001/iceberg"),
- /* embedDiscovery= */ true);
-
- Assertions.assertEquals("jdbc", config.get("iceberg.catalog.type"));
- Assertions.assertNull(config.get("iceberg.rest-catalog.uri"));
+ Assertions.assertThrows(
+ TrinoException.class,
+ () ->
+ buildConnectorConfig(
+ "catalog1",
+ properties,
+ icebergRestDiscoveredConfig("other_metalake",
"http://other:9001/iceberg"),
+ /* embedDiscovery= */ true));
}
@Test
- public void testDiscoveredUriUnusedWithoutEmbedding() throws Exception {
+ public void testDiscoveredUriUnusedWithoutEmbedding() {
// A worker never runs discovery, so its own GravitinoConfig never has a
discovered value
// populated — this asserts that IcebergConnectorAdapter really does read
the routing signal
// from the catalog, not from GravitinoConfig's discovered map directly.
@@ -415,15 +412,14 @@ public class TestIcebergCatalogPropertyConverter {
.put("jdbc-driver", "org.postgresql.Driver")
.build();
- Map<String, String> config =
- buildConnectorConfig(
- "catalog1",
- properties,
- icebergRestDiscoveredConfig("test",
"http://discovered:9001/iceberg"),
- /* embedDiscovery= */ false);
-
- Assertions.assertEquals("jdbc", config.get("iceberg.catalog.type"));
- Assertions.assertNull(config.get("iceberg.rest-catalog.uri"));
+ Assertions.assertThrows(
+ TrinoException.class,
+ () ->
+ buildConnectorConfig(
+ "catalog1",
+ properties,
+ icebergRestDiscoveredConfig("test",
"http://discovered:9001/iceberg"),
+ /* embedDiscovery= */ false));
}
@Test
@@ -784,9 +780,10 @@ public class TestIcebergCatalogPropertyConverter {
}
private static GravitinoConfig icebergRestUnavailableConfig() {
- // No manual gravitino.iceberg.rest-uri, and nothing discovered for this
metalake: the
- // connector has no endpoint to route through, so catalogs fall back to
their own backend.
- return new GravitinoConfig(ImmutableMap.of("gravitino.metalake", "test"));
+ return new GravitinoConfig(
+ ImmutableMap.of(
+ "gravitino.metalake", "test",
+ "gravitino.iceberg.rest-routing-enabled", "false"));
}
private static GravitinoConfig icebergRestConfiguredConfig(Map<String,
String> extraConfig) {