This is an automated email from the ASF dual-hosted git repository. github-actions[bot] pushed a commit to branch cherry-pick-52b8f534-to-branch-1.3 in repository https://gitbox.apache.org/repos/asf/gravitino.git
commit b98e15f51e7f7704f411201e37f35617670ff798 Author: Yuhui <[email protected]> AuthorDate: Thu Aug 27 14:12:17 2026 +0800 [#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635) ### What changes were proposed in this pull request? Switch trino-connector's logging from SLF4J/Log4j2 to `io.airlift.log.Logger`, matching Trino's own connectors. ### Why are the changes needed? Trino loads the connector plugin in an isolated classloader. Log4j2 is classloader-scoped and the plugin jar ships no config, so it falls back to Log4j2's default (root ERROR, console) — most log output is silently dropped instead of reaching `var/log/server.log`. `io.airlift.log.Logger` wraps `java.util.logging`, a JVM-wide singleton unaffected by classloader isolation. Fix: #12634 ### Does this PR introduce any user-facing change? No. Trino admins will now see Gravitino connector logs in `var/log/server.log`. ### How was this patch tested? - Compiled all trino-connector modules (base + 5 version-range variants) - `./gradlew :trino-connector:trino-connector:test` - `./build.sh sp` # Conflicts: # trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java # trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java # trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java # trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/security/GravitinoAuthProvider.java # trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/util/json/JsonCodec.java --- build.gradle.kts | 16 +- docs/trino-connector/development.md | 23 +-- docs/trino-connector/installation.md | 3 +- gradle/libs.versions.toml | 4 +- .../trino-connector-435-439/build.gradle.kts | 3 +- .../trino-connector-440-445/build.gradle.kts | 3 +- .../trino-connector-446-451/build.gradle.kts | 3 +- .../trino-connector-452-468/build.gradle.kts | 3 +- .../trino-connector-469-472/build.gradle.kts | 3 +- .../trino-connector-473-478/build.gradle.kts | 3 +- trino-connector/trino-connector/build.gradle.kts | 3 +- .../trino/connector/GravitinoConnector.java | 30 ++- .../trino/connector/GravitinoConnectorFactory.java | 19 +- .../connector/GravitinoConnectorPluginManager.java | 18 +- .../trino/connector/GravitinoMetadata.java | 13 +- .../connector/catalog/CatalogConnectorManager.java | 46 ++--- .../catalog/CatalogConnectorMetadata.java | 7 +- .../catalog/CatalogConnectorMetadataAdapter.java | 7 +- .../catalog/CatalogPropertyConverter.java | 9 +- .../trino/connector/catalog/CatalogRegister.java | 220 +++++++++++++++++++-- .../catalog/DefaultCatalogConnectorFactory.java | 5 +- .../catalog/jdbc/mysql/MySQLMetadataAdapter.java | 3 +- .../connector/security/GravitinoAuthProvider.java | 15 +- .../AlterCatalogStoredProcedure.java | 7 +- .../CreateCatalogStoredProcedure.java | 7 +- .../DropCatalogStoredProcedure.java | 12 +- .../trino/connector/util/json/JsonCodec.java | 117 +++++++++++ 27 files changed, 469 insertions(+), 133 deletions(-) diff --git a/build.gradle.kts b/build.gradle.kts index a62a0ea857..5e957ef255 100644 --- a/build.gradle.kts +++ b/build.gradle.kts @@ -157,11 +157,17 @@ allprojects { "import\\s+.*\\.org\\.apache\\.commons\\.io\\.([A-Z][a-zA-Z0-9_]*);", "import org.apache.commons.io.${'$'}1;" ) - replaceRegex( - "Use SLF4J Logger instead of other logging frameworks", - "import\\s+.*\\.(Logger|LoggerFactory);", - "import org.slf4j.${'$'}1;" - ) + // The trino-connector module logs via io.airlift.log.Logger to match Trino's own + // logging so plugin log output routes into Trino's unified log, instead of SLF4J. + // integration-test is intentionally excluded from this carve-out: it runs outside + // Trino's isolated plugin classloader, so it should keep using SLF4J as normal. + if (!project.path.startsWith(":trino-connector:trino-connector")) { + replaceRegex( + "Use SLF4J Logger instead of other logging frameworks", + "import\\s+.*\\.(Logger|LoggerFactory);", + "import org.slf4j.${'$'}1;" + ) + } replaceRegex( "Remove Testcontainers shading", "import\\s+org\\.testcontainers\\.shaded\\.([^;]+);", diff --git a/docs/trino-connector/development.md b/docs/trino-connector/development.md index 5120955a77..cd2b54aa47 100644 --- a/docs/trino-connector/development.md +++ b/docs/trino-connector/development.md @@ -178,27 +178,14 @@ Change `localhost`, `port`, and the names of metalake and catalogs to match your </dependency> <dependency> - <groupId>org.slf4j</groupId> - <artifactId>slf4j-api</artifactId> - <version>2.0.9</version> - </dependency> - - <dependency> - <groupId>org.apache.logging.log4j</groupId> - <artifactId>log4j-slf4j2-impl</artifactId> - <version>2.22.0</version> - </dependency> - - <dependency> - <groupId>org.apache.logging.log4j</groupId> - <artifactId>log4j-api</artifactId> - <version>2.22.0</version> + <groupId>io.airlift</groupId> + <artifactId>log</artifactId> </dependency> <dependency> - <groupId>org.apache.logging.log4j</groupId> - <artifactId>log4j-core</artifactId> - <version>2.22.0</version> + <groupId>org.slf4j</groupId> + <artifactId>slf4j-jdk14</artifactId> + <version>2.0.17</version> </dependency> <dependency> diff --git a/docs/trino-connector/installation.md b/docs/trino-connector/installation.md index d39b200ee3..9543261510 100644 --- a/docs/trino-connector/installation.md +++ b/docs/trino-connector/installation.md @@ -64,8 +64,7 @@ After unpacking, you can see the connector directory: 1. Download and unpack the correct Gravitino Trino connector tarball for your Trino version. 2. Rename the unpacked connector directory to `gravitino`, and then copy it to the Trino plugin directory. Normally, the directory location is `Trino-server-<version>/plugin`, and the directory contains other catalogs used by Trino. -3. Add Trino JVM arguments `-Dlog4j.configurationFile=file:////etc/trino/log4j2.properties` to enable logging for the Gravitino Trino connector. -4. Update Trino coordinator configuration. +3. Update Trino coordinator configuration. You need to set `catalog.management=dynamic`, The config location is `Trino-server-<version>/etc/config.properties`, and the contents like: ```text diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index 6c93c5af40..bc85fd3732 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -38,6 +38,7 @@ jetty = "9.4.58.v20250814" jersey = "2.41" mockito = "4.11.0" airlift-json = "237" +airlift-log = "237" airlift-resolver = "1.6" hive2 = "2.3.9" hive3 = "3.1.3" @@ -195,7 +196,7 @@ junit-jupiter-params = { group = "org.junit.jupiter", name = "junit-jupiter-para junit-jupiter-engine = { group = "org.junit.jupiter", name = "junit-jupiter-engine", version.ref = "junit" } slf4j-api = { group = "org.slf4j", name = "slf4j-api", version.ref = "slf4j" } slf4j-simple = { group = "org.slf4j", name = "slf4j-simple", version.ref = "slf4j" } -slf4j-jdk14 = { group = "org.slf4j", name = "slf4j-jdk14", version = "1.7.30" } +slf4j-jdk14 = { group = "org.slf4j", name = "slf4j-jdk14", version.ref = "slf4j" } log4j-slf4j2-impl = { group = "org.apache.logging.log4j", name = "log4j-slf4j2-impl", version.ref = "log4j" } log4j-api = { group = "org.apache.logging.log4j", name = "log4j-api", version.ref = "log4j" } log4j-core = { group = "org.apache.logging.log4j", name = "log4j-core", version.ref = "log4j" } @@ -244,6 +245,7 @@ hadoop3-shaded-guava = { group = "org.apache.hadoop.thirdparty", name = "hadoop- hadoop3-shaded-protobuf = { group = "org.apache.hadoop.thirdparty", name = "hadoop-shaded-protobuf_3_7", version.ref = "hadoop3-thirdparty" } re2j = { group = "com.google.re2j", name = "re2j", version.ref = "re2j" } airlift-json = { group = "io.airlift", name = "json", version.ref = "airlift-json"} +airlift-log = { group = "io.airlift", name = "log", version.ref = "airlift-log" } airlift-resolver = { group = "io.airlift.resolver", name = "resolver", version.ref = "airlift-resolver"} httpclient = { group = "org.apache.httpcomponents", name = "httpclient", version.ref = "httpclient" } httpclient5 = { group = "org.apache.httpcomponents.client5", name = "httpclient5", version.ref = "httpclient5" } diff --git a/trino-connector/trino-connector-435-439/build.gradle.kts b/trino-connector/trino-connector-435-439/build.gradle.kts index befe0343a7..a3864ffec9 100644 --- a/trino-connector/trino-connector-435-439/build.gradle.kts +++ b/trino-connector/trino-connector-435-439/build.gradle.kts @@ -51,7 +51,8 @@ dependencies { implementation(project(":catalogs:catalog-common")) implementation(project(":clients:client-java-runtime", configuration = "shadow")) implementation(libs.airlift.json) - implementation(libs.bundles.log4j) + implementation(libs.airlift.log) + implementation(libs.slf4j.jdk14) implementation(libs.commons.collections4) implementation(libs.commons.lang3) implementation("io.trino:trino-jdbc:$trinoVersion") diff --git a/trino-connector/trino-connector-440-445/build.gradle.kts b/trino-connector/trino-connector-440-445/build.gradle.kts index d7bddb39e7..af4f5158ea 100644 --- a/trino-connector/trino-connector-440-445/build.gradle.kts +++ b/trino-connector/trino-connector-440-445/build.gradle.kts @@ -52,11 +52,12 @@ dependencies { implementation(project(":catalogs:catalog-common")) implementation(project(":clients:client-java-runtime", configuration = "shadow")) implementation(libs.airlift.json) - implementation(libs.bundles.log4j) implementation(libs.commons.collections4) implementation(libs.commons.lang3) implementation("io.trino:trino-jdbc:$trinoVersion") runtimeOnly("io.opentelemetry.semconv:opentelemetry-semconv:$otelSemconvVersion") + implementation(libs.airlift.log) + implementation(libs.slf4j.jdk14) compileOnly(libs.airlift.resolver) compileOnly("io.trino:trino-spi:$trinoVersion") { exclude("org.apache.logging.log4j") diff --git a/trino-connector/trino-connector-446-451/build.gradle.kts b/trino-connector/trino-connector-446-451/build.gradle.kts index 6ae0102706..e4e75eda71 100644 --- a/trino-connector/trino-connector-446-451/build.gradle.kts +++ b/trino-connector/trino-connector-446-451/build.gradle.kts @@ -52,11 +52,12 @@ dependencies { implementation(project(":catalogs:catalog-common")) implementation(project(":clients:client-java-runtime", configuration = "shadow")) implementation(libs.airlift.json) - implementation(libs.bundles.log4j) implementation(libs.commons.collections4) implementation(libs.commons.lang3) implementation("io.trino:trino-jdbc:$trinoVersion") runtimeOnly("io.opentelemetry.semconv:opentelemetry-semconv-incubating:$otelSemconvVersion") + implementation(libs.airlift.log) + implementation(libs.slf4j.jdk14) compileOnly(libs.airlift.resolver) compileOnly("io.trino:trino-spi:$trinoVersion") { exclude("org.apache.logging.log4j") diff --git a/trino-connector/trino-connector-452-468/build.gradle.kts b/trino-connector/trino-connector-452-468/build.gradle.kts index 6bcdff5e41..767903355f 100644 --- a/trino-connector/trino-connector-452-468/build.gradle.kts +++ b/trino-connector/trino-connector-452-468/build.gradle.kts @@ -52,11 +52,12 @@ dependencies { implementation(project(":catalogs:catalog-common")) implementation(project(":clients:client-java-runtime", configuration = "shadow")) implementation(libs.airlift.json) - implementation(libs.bundles.log4j) implementation(libs.commons.collections4) implementation(libs.commons.lang3) implementation("io.trino:trino-jdbc:$trinoVersion") runtimeOnly("io.opentelemetry.semconv:opentelemetry-semconv:$otelSemconvVersion") + implementation(libs.airlift.log) + implementation(libs.slf4j.jdk14) compileOnly(libs.airlift.resolver) compileOnly("io.trino:trino-spi:$trinoVersion") { exclude("org.apache.logging.log4j") diff --git a/trino-connector/trino-connector-469-472/build.gradle.kts b/trino-connector/trino-connector-469-472/build.gradle.kts index ef3857a23c..e8ca4e2994 100644 --- a/trino-connector/trino-connector-469-472/build.gradle.kts +++ b/trino-connector/trino-connector-469-472/build.gradle.kts @@ -57,11 +57,12 @@ dependencies { compileOnly(project(":common")) implementation(libs.airlift.json) - implementation(libs.bundles.log4j) implementation(libs.commons.collections4) implementation(libs.commons.lang3) implementation("io.trino:trino-jdbc:$trinoVersion") runtimeOnly("io.opentelemetry.semconv:opentelemetry-semconv:$otelSemconvVersion") + implementation(libs.airlift.log) + implementation(libs.slf4j.jdk14) compileOnly(libs.airlift.resolver) compileOnly("io.trino:trino-spi:$trinoVersion") { exclude("org.apache.logging.log4j") diff --git a/trino-connector/trino-connector-473-478/build.gradle.kts b/trino-connector/trino-connector-473-478/build.gradle.kts index d74624bb6d..f41b2aba87 100644 --- a/trino-connector/trino-connector-473-478/build.gradle.kts +++ b/trino-connector/trino-connector-473-478/build.gradle.kts @@ -52,11 +52,12 @@ dependencies { implementation(project(":catalogs:catalog-common")) implementation(project(":clients:client-java-runtime", configuration = "shadow")) implementation(libs.airlift.json) - implementation(libs.bundles.log4j) implementation(libs.commons.collections4) implementation(libs.commons.lang3) implementation("io.trino:trino-jdbc:$trinoVersion") runtimeOnly("io.opentelemetry.semconv:opentelemetry-semconv:$otelSemconvVersion") + implementation(libs.airlift.log) + implementation(libs.slf4j.jdk14) compileOnly(libs.airlift.resolver) compileOnly("io.trino:trino-spi:$trinoVersion") { exclude("org.apache.logging.log4j") diff --git a/trino-connector/trino-connector/build.gradle.kts b/trino-connector/trino-connector/build.gradle.kts index 69ed4e7593..17b87d889c 100644 --- a/trino-connector/trino-connector/build.gradle.kts +++ b/trino-connector/trino-connector/build.gradle.kts @@ -35,10 +35,11 @@ dependencies { implementation(project(":clients:client-java-runtime", configuration = "shadow")) implementation(libs.airlift.json) - implementation(libs.bundles.log4j) implementation(libs.commons.collections4) implementation(libs.commons.lang3) implementation("io.trino:trino-jdbc:$trinoVersion") + implementation(libs.airlift.log) + implementation(libs.slf4j.jdk14) compileOnly(libs.airlift.resolver) compileOnly("io.trino:trino-spi:$trinoVersion") { exclude("org.apache.logging.log4j") diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java index 01d419d1fe..ebca989a4c 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java @@ -25,6 +25,11 @@ import com.google.common.base.Preconditions; import com.google.common.cache.Cache; import com.google.common.cache.CacheBuilder; import com.google.common.cache.RemovalNotification; +<<<<<<< HEAD +======= +import com.google.common.util.concurrent.UncheckedExecutionException; +import io.airlift.log.Logger; +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) import io.trino.spi.TrinoException; import io.trino.spi.connector.Connector; import io.trino.spi.connector.ConnectorAccessControl; @@ -54,8 +59,6 @@ import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadata; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadataAdapter; import org.apache.gravitino.trino.connector.security.GravitinoAuthProvider; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * GravitinoConnector serves as the entry point for operations on the connector managed by Trino and @@ -64,7 +67,7 @@ import org.slf4j.LoggerFactory; */ public class GravitinoConnector implements Connector { - private static final Logger LOG = LoggerFactory.getLogger(GravitinoConnector.class); + private static final Logger LOG = Logger.get(GravitinoConnector.class); private final NameIdentifier catalogIdentifier; protected final CatalogConnectorContext catalogConnectorContext; @@ -243,9 +246,28 @@ public class GravitinoConnector implements Connector { } catch (ExecutionException e) { Throwable cause = e.getCause(); LOG.warn( +<<<<<<< HEAD "Failed to create per-user Gravitino client for user '{}': {}", session.getUser(), cause.getMessage()); +======= + cause, "Failed to create per-user Gravitino client for user '%s'", session.getUser()); + if (cause instanceof TrinoException) { + // Already carries a specific Trino error code (e.g. from buildForSession); re-wrapping + // would swallow it. + throw (TrinoException) cause; + } + if (cause instanceof IllegalArgumentException + || cause instanceof UnsupportedOperationException) { + throw new TrinoException( + PERMISSION_DENIED, + "Failed to authenticate user '" + + session.getUser() + + "' with Gravitino: " + + cause.getMessage(), + cause); + } +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) throw new TrinoException( PERMISSION_DENIED, "Failed to authenticate user '" @@ -304,7 +326,7 @@ public class GravitinoConnector implements Connector { try { client.close(); } catch (Exception e) { - LOG.warn("Failed to close GravitinoAdminClient", e); + LOG.warn(e, "Failed to close GravitinoAdminClient"); } } } diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java index 1420a60c23..31cbdb39bd 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java @@ -24,6 +24,7 @@ import static org.apache.gravitino.trino.connector.GravitinoErrorCode.GRAVITINO_ import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; import com.google.common.base.Strings; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import io.trino.spi.connector.Connector; import io.trino.spi.connector.ConnectorContext; @@ -39,13 +40,11 @@ import org.apache.gravitino.trino.connector.catalog.DefaultCatalogConnectorFacto import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** Gravitino connector factory. */ public class GravitinoConnectorFactory implements ConnectorFactory { - private static final Logger LOG = LoggerFactory.getLogger(GravitinoConnectorFactory.class); + private static final Logger LOG = Logger.get(GravitinoConnectorFactory.class); private static final int MIN_SUPPORT_TRINO_SPI_VERSION = 435; private static final int MAX_SUPPORT_TRINO_SPI_VERSION = Integer.MAX_VALUE; /** The default connector name. */ @@ -116,6 +115,13 @@ public class GravitinoConnectorFactory implements ConnectorFactory { LOG.error(message); throw new TrinoException(GRAVITINO_RUNTIME_ERROR, message, e); } +<<<<<<< HEAD +======= + } catch (Exception e) { + String message = "Initialization of the GravitinoConnector failed " + e.getMessage(); + LOG.error(e, message); + throw new TrinoException(GRAVITINO_RUNTIME_ERROR, message, e); +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) } } @@ -162,10 +168,9 @@ public class GravitinoConnectorFactory implements ConnectorFactory { // check catalog name with metalake are supported in this trino version if (!config.singleMetalakeMode() && !supportCatalogNameWithMetalake()) { LOG.warn( - "The trino-connector-{}-{} does not fully support catalog name with metalake. " + "The trino-connector-%s-%s does not fully support catalog name with metalake. " + "The DROP CATALOG operation may not work correctly in multi-metalake mode.", - getMinSupportTrinoSpiVersion(), - getMaxSupportTrinoSpiVersion()); + getMinSupportTrinoSpiVersion(), getMaxSupportTrinoSpiVersion()); } // skip version validation @@ -174,7 +179,7 @@ public class GravitinoConnectorFactory implements ConnectorFactory { if (trinoVersion < getMinSupportTrinoSpiVersion() || trinoVersion > getMaxSupportTrinoSpiVersion()) { LOG.warn( - "Trino version {} has not been tested with Gravitino and may have compatibility issues", + "Trino version %s has not been tested with Gravitino and may have compatibility issues", trinoVersion); } return; diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorPluginManager.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorPluginManager.java index ccf59b8e59..192f272d55 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorPluginManager.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorPluginManager.java @@ -22,6 +22,7 @@ import static org.apache.gravitino.trino.connector.GravitinoConfig.TRINO_PLUGIN_ import com.google.common.base.Splitter; import com.google.common.collect.ImmutableList; +import io.airlift.log.Logger; import io.airlift.resolver.ArtifactResolver; import io.trino.spi.Plugin; import io.trino.spi.TrinoException; @@ -39,14 +40,12 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.ServiceLoader; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import org.sonatype.aether.artifact.Artifact; /** This class is mange the internal connector plugin and help to create the connector. */ public class GravitinoConnectorPluginManager { - private static final Logger LOG = LoggerFactory.getLogger(GravitinoConnectorPluginManager.class); + private static final Logger LOG = Logger.get(GravitinoConnectorPluginManager.class); /** The app class loader name. */ public static final String APP_CLASS_LOADER_NAME = "app"; @@ -136,7 +135,7 @@ public class GravitinoConnectorPluginManager { .forEach( file -> { loadPlugin(pluginDir, file.getName()); - LOG.info("Load plugin {}/{} successful", pluginDir, file.getName()); + LOG.info("Load plugin %s/%s successful", pluginDir, file.getName()); }); } catch (Exception e) { throw new TrinoException( @@ -149,7 +148,7 @@ public class GravitinoConnectorPluginManager { File directory = new File(dirName); File[] pluginFiles = directory.listFiles(); if (pluginFiles == null || pluginFiles.length == 0) { - LOG.warn("Cannot load plugin {} from empty directory {}", pluginName, dirName); + LOG.warn("Cannot load plugin %s from empty directory %s", pluginName, dirName); return; } List<URL> files = @@ -183,6 +182,7 @@ public class GravitinoConnectorPluginManager { "io.trino.spi.", "com.fasterxml.jackson.annotation.", "io.airlift.slice.", + "io.airlift.log.", "org.openjdk.jol.", "io.opentelemetry.api.", "io.opentelemetry.context.")); @@ -191,13 +191,13 @@ public class GravitinoConnectorPluginManager { ServiceLoader.load(Plugin.class, (ClassLoader) pluginClassLoader); List<Plugin> pluginList = ImmutableList.copyOf(serviceLoader); if (pluginList.isEmpty()) { - LOG.warn("The {} plugin directory does not contain a connector SPI interface", pluginName); + LOG.warn("The %s plugin directory does not contain a connector SPI interface", pluginName); return; } Plugin plugin = pluginList.get(0); if (plugin.getConnectorFactories() == null || !plugin.getConnectorFactories().iterator().hasNext()) { - LOG.warn("The {} plugin does not contain any ConnectorFactories", pluginName); + LOG.warn("The %s plugin does not contain any ConnectorFactories", pluginName); return; } connectorPlugins.put(pluginName, pluginList.get(0)); @@ -235,7 +235,7 @@ public class GravitinoConnectorPluginManager { try { loadPluginByPom(artifactResolver.resolvePom(new File(v)), key); } catch (Exception e) { - LOG.error("Fatal error in load plugin by {}", v, e); + LOG.error(e, "Fatal error in load plugin by %s", v); } }); } @@ -292,7 +292,7 @@ public class GravitinoConnectorPluginManager { new ThreadContextClassLoader(plugin.getClass().getClassLoader())) { ConnectorFactory connectorFactory = plugin.getConnectorFactories().iterator().next(); Connector connector = connectorFactory.create(connectorName, config, context); - LOG.info("create connector {} with config {} successful", connectorName, config); + LOG.info("create connector %s with config %s successful", connectorName, config); return connector; } } catch (Exception e) { diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoMetadata.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoMetadata.java index 04a264e077..f05b557233 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoMetadata.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoMetadata.java @@ -24,6 +24,7 @@ import static org.apache.gravitino.trino.connector.GravitinoErrorCode.GRAVITINO_ import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableSet; +import io.airlift.log.Logger; import io.airlift.slice.Slice; import io.trino.spi.TrinoException; import io.trino.spi.connector.AggregateFunction; @@ -85,8 +86,6 @@ import org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadata; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadataAdapter; import org.apache.gravitino.trino.connector.metadata.GravitinoSchema; import org.apache.gravitino.trino.connector.metadata.GravitinoTable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * The GravitinoMetadata class provides operations for Apache Gravitino metadata on the Gravitino @@ -95,7 +94,7 @@ import org.slf4j.LoggerFactory; */ public abstract class GravitinoMetadata implements ConnectorMetadata { - private static final Logger LOG = LoggerFactory.getLogger(GravitinoMetadata.class); + private static final Logger LOG = Logger.get(GravitinoMetadata.class); // The column handle name that will generate row IDs for the merge operation. public static final String MERGE_ROW_ID = "$row_id"; @@ -285,7 +284,7 @@ public abstract class GravitinoMetadata implements ConnectorMetadata { try { catalogConnectorMetadata.dropTable(tableName); } catch (Exception dropException) { - LOG.warn("Failed to drop table {} during CTAS cleanup", tableName, dropException); + LOG.warn(dropException, "Failed to drop table %s during CTAS cleanup", tableName); } throw e; } @@ -310,7 +309,7 @@ public abstract class GravitinoMetadata implements ConnectorMetadata { // (e.g., Hive 'format' is a String in Gravitino but HiveStorageFormat enum internally). // Returning empty is correct for non-bucketed CTAS. LOG.debug( - "Skipping internal getNewTableLayout due to property type mismatch: {}", e.getMessage()); + "Skipping internal getNewTableLayout due to property type mismatch: %s", e.getMessage()); return Optional.empty(); } } @@ -858,7 +857,7 @@ public abstract class GravitinoMetadata implements ConnectorMetadata { } return toLanguageFunctions(function); } catch (NoSuchFunctionException e) { - LOG.debug("Function {} not found in schema {}", name.getFunctionName(), name.getSchemaName()); + LOG.debug("Function %s not found in schema %s", name.getFunctionName(), name.getSchemaName()); return List.of(); } } @@ -881,7 +880,7 @@ public abstract class GravitinoMetadata implements ConnectorMetadata { String signatureToken = buildSignatureToken(function.name(), definition.parameters()); result.add(new LanguageFunction(signatureToken, sql, List.of(), Optional.empty())); } catch (TrinoException e) { - LOG.warn("Failed to build signature token for function {}", function.name(), e); + LOG.warn(e, "Failed to build signature token for function %s", function.name()); } } } 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 7539a0b439..2bd0fb826a 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 @@ -20,6 +20,7 @@ package org.apache.gravitino.trino.connector.catalog; import com.google.common.base.Preconditions; import com.google.common.util.concurrent.ThreadFactoryBuilder; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import io.trino.spi.connector.ConnectorContext; import java.util.Arrays; @@ -43,8 +44,6 @@ import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; import org.apache.gravitino.trino.connector.security.GravitinoAuthProvider; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * This class has the following main functions: @@ -57,7 +56,7 @@ import org.slf4j.LoggerFactory; * </pre> */ public class CatalogConnectorManager { - private static final Logger LOG = LoggerFactory.getLogger(CatalogConnectorManager.class); + private static final Logger LOG = Logger.get(CatalogConnectorManager.class); private static final int NUMBER_EXECUTOR_THREAD = 1; private static final int LOAD_METALAKE_TIMEOUT = 60; @@ -103,7 +102,7 @@ public class CatalogConnectorManager { .setNameFormat("gravitino-connector-schedule-%d") .setUncaughtExceptionHandler( (thread, throwable) -> - LOG.warn("{} uncaught exception:", thread.getName(), throwable)) + LOG.warn(throwable, "%s uncaught exception:", thread.getName())) .build()); } @@ -119,7 +118,7 @@ public class CatalogConnectorManager { if (client == null) { String authType = config.getClientConfig().getOrDefault(GravitinoAuthProvider.AUTH_TYPE_KEY, "none"); - LOG.info("Building Gravitino client with authType: {}", authType); + LOG.info("Building Gravitino client with authType: %s", authType); try { this.gravitinoClient = GravitinoAuthProvider.build(config); } catch (IllegalArgumentException e) { @@ -183,14 +182,14 @@ public class CatalogConnectorManager { for (String usedMetalake : usedMetalakes) { try { GravitinoMetalake metalake = metalakes.get(usedMetalake); - LOG.debug("Load metalake: {}", usedMetalake); + LOG.debug("Load metalake: %s", usedMetalake); loadCatalogs(metalake); } catch (Exception e) { - LOG.error("Load Metalake {} failed.", usedMetalake, e); + LOG.error(e, "Load Metalake %s failed.", usedMetalake); } } } catch (Exception e) { - LOG.error("Error when loading metalake", e); + LOG.error(e, "Error when loading metalake"); } } @@ -219,11 +218,11 @@ public class CatalogConnectorManager { .filter(id -> !skipCatalog(getTrinoCatalogName(metalake.name(), id))) .collect(Collectors.toList()); } catch (Exception e) { - LOG.error("Failed to list catalogs in metalake {}.", metalake.name(), e); + LOG.error(e, "Failed to list catalogs in metalake %s.", metalake.name()); return; } - LOG.debug("Load metalake {}'s catalogs. catalogs: {}.", metalake.name(), catalogNames); + LOG.debug("Load metalake %s's catalogs. catalogs: %s.", metalake.name(), catalogNames); // Delete those catalogs that have been deleted in Gravitino server Set<String> catalogNameStrings = @@ -239,7 +238,7 @@ public class CatalogConnectorManager { try { unloadCatalog(entry.getValue().getCatalog()); } catch (Exception e) { - LOG.error("Failed to remove catalog {}.", entry.getKey(), e); + LOG.error(e, "Failed to remove catalog %s.", entry.getKey()); } } } @@ -264,13 +263,11 @@ public class CatalogConnectorManager { } } catch (UnsupportedOperationException e) { LOG.warn( - "Unsupported catalog type for catalog {} in metalake {}: {}", - catalogName, - metalake.name(), - e.getMessage()); + "Unsupported catalog type for catalog %s in metalake %s: %s", + catalogName, metalake.name(), e.getMessage()); } catch (Exception e) { LOG.error( - "Failed to load metalake {}'s catalog {}.", metalake.name(), catalogName, e); + e, "Failed to load metalake %s's catalog %s.", metalake.name(), catalogName); } }); } @@ -286,12 +283,12 @@ public class CatalogConnectorManager { catalogConnectors.remove(catalogFullName); loadCatalogImpl(catalog); - LOG.info("Update catalog '{}' in metalake {} successfully.", catalog, catalog.getMetalake()); + LOG.info("Update catalog '%s' in metalake %s successfully.", catalog, catalog.getMetalake()); } private void loadCatalog(GravitinoCatalog catalog) { loadCatalogImpl(catalog); - LOG.info("Load catalog {} in metalake {} successfully.", catalog, catalog.getMetalake()); + LOG.info("Load catalog %s in metalake %s successfully.", catalog, catalog.getMetalake()); } private void loadCatalogImpl(GravitinoCatalog catalog) { @@ -300,7 +297,7 @@ public class CatalogConnectorManager { } catch (Exception e) { String message = String.format("Failed to create internal catalog connector. The catalog is: %s", catalog); - LOG.error(message, e); + LOG.error(e, message); throw new TrinoException( GravitinoErrorCode.GRAVITINO_CREATE_INTERNAL_CONNECTOR_ERROR, message, e); } @@ -311,9 +308,8 @@ public class CatalogConnectorManager { catalogRegister.unregisterCatalog(catalogFullName); catalogConnectors.remove(catalogFullName); LOG.info( - "Remove catalog '{}' in metalake {} successfully.", - catalog.getName(), - catalog.getMetalake()); + "Remove catalog '%s' in metalake %s successfully.", + catalog.getName(), catalog.getMetalake()); } /** @@ -422,10 +418,10 @@ public class CatalogConnectorManager { CatalogConnectorContext connectorContext = builder.build(); String fullCatalogName = getTrinoCatalogName(catalog); catalogConnectors.put(fullCatalogName, connectorContext); - LOG.info("Create connector {} successful", connectorName); + LOG.info("Create connector %s successful", connectorName); return connectorContext; } catch (Exception e) { - LOG.error("Failed to create connector: {}", connectorName, e); + LOG.error(e, "Failed to create connector: %s", connectorName); throw new TrinoException( GravitinoErrorCode.GRAVITINO_OPERATION_FAILED, "Failed to create connector: " + connectorName, @@ -464,7 +460,7 @@ public class CatalogConnectorManager { for (Pattern pattern : config.getSkipCatalogPatterns()) { if (pattern.matcher(catalogName).matches()) { LOG.debug( - "Skip catalog {} with config `gravitino.trino.skip-catalog-patterns`.", catalogName); + "Skip catalog %s with config `gravitino.trino.skip-catalog-patterns`.", catalogName); return true; } } diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadata.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadata.java index 7fd83e7575..4e71662872 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadata.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadata.java @@ -19,6 +19,7 @@ package org.apache.gravitino.trino.connector.catalog; import com.google.common.base.Strings; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import io.trino.spi.connector.SchemaTableName; import java.util.Arrays; @@ -47,13 +48,11 @@ import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.metadata.GravitinoColumn; import org.apache.gravitino.trino.connector.metadata.GravitinoSchema; import org.apache.gravitino.trino.connector.metadata.GravitinoTable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** This class implements Apache Gravitino metadata operators. */ public class CatalogConnectorMetadata { - private static final Logger LOG = LoggerFactory.getLogger(CatalogConnectorMetadata.class); + private static final Logger LOG = Logger.get(CatalogConnectorMetadata.class); private static final String CATALOG_DOES_NOT_EXIST_MSG = "Catalog does not exist"; private static final String SCHEMA_DOES_NOT_EXIST_MSG = "Schema does not exist"; @@ -80,7 +79,7 @@ public class CatalogConnectorMetadata { try { fc = catalog.asFunctionCatalog(); } catch (UnsupportedOperationException e) { - LOG.debug("Catalog {} does not support function operations", catalogName); + LOG.debug("Catalog %s does not support function operations", catalogName); } this.functionCatalog = fc; } catch (NoSuchCatalogException e) { diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadataAdapter.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadataAdapter.java index cb6f54110b..d61f1d79e8 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadataAdapter.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorMetadataAdapter.java @@ -18,6 +18,7 @@ */ package org.apache.gravitino.trino.connector.catalog; +import io.airlift.log.Logger; import io.trino.spi.connector.ColumnMetadata; import io.trino.spi.connector.ConnectorTableMetadata; import io.trino.spi.connector.ConnectorTableProperties; @@ -35,8 +36,6 @@ import org.apache.gravitino.trino.connector.metadata.GravitinoColumn; import org.apache.gravitino.trino.connector.metadata.GravitinoSchema; import org.apache.gravitino.trino.connector.metadata.GravitinoTable; import org.apache.gravitino.trino.connector.util.GeneralDataTypeTransformer; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * This interface is used to handle different parts of catalog metadata from different catalog @@ -44,7 +43,7 @@ import org.slf4j.LoggerFactory; */ public class CatalogConnectorMetadataAdapter { - private static final Logger LOG = LoggerFactory.getLogger(CatalogConnectorMetadataAdapter.class); + private static final Logger LOG = Logger.get(CatalogConnectorMetadataAdapter.class); /** The list of schema properties supported by this catalog connector. */ protected final List<PropertyMetadata<?>> schemaProperties; @@ -217,7 +216,7 @@ public class CatalogConnectorMetadataAdapter { if (properties.containsKey(name)) { validProperties.put(name, properties.get(name)); } else { - LOG.warn("Property {} is not defined in Trino, we will ignore it", name); + LOG.warn("Property %s is not defined in Trino, we will ignore it", name); } } return validProperties; diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogPropertyConverter.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogPropertyConverter.java index adb6d11d43..b42ce2a9c5 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogPropertyConverter.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogPropertyConverter.java @@ -20,14 +20,13 @@ package org.apache.gravitino.trino.connector.catalog; import com.google.common.collect.ImmutableMap; +import io.airlift.log.Logger; import java.util.HashMap; import java.util.Map; import org.apache.gravitino.catalog.property.PropertyConverter; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; public class CatalogPropertyConverter extends PropertyConverter { - private static final Logger LOG = LoggerFactory.getLogger(PropertyConverter.class); + private static final Logger LOG = Logger.get(CatalogPropertyConverter.class); private static final String TRINO_PROPERTIES_PREFIX = "trino.bypass."; @@ -59,7 +58,7 @@ public class CatalogPropertyConverter extends PropertyConverter { if (engineKey != null) { engineProperties.put(engineKey, entry.getValue()); } else { - LOG.info("Property {} is not supported by engine", entry.getKey()); + LOG.info("Property %s is not supported by engine", entry.getKey()); } } // trino.bypass properties will be skipped when the catalog properties is defined by Gravitino @@ -69,7 +68,7 @@ public class CatalogPropertyConverter extends PropertyConverter { if (!engineProperties.containsKey(key)) { engineProperties.put(key, entry.getValue()); } else { - LOG.info("Property {} which with trino.bypass prefix is skipped", key); + LOG.info("Property %s which with trino.bypass prefix is skipped", key); } } } diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java index 09bf1cd658..104ba75d71 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java @@ -21,6 +21,12 @@ package org.apache.gravitino.trino.connector.catalog; import static org.apache.gravitino.trino.connector.GravitinoConfig.GRAVITINO_DYNAMIC_CONNECTOR; import static org.apache.gravitino.trino.connector.GravitinoConfig.GRAVITINO_DYNAMIC_CONNECTOR_CATALOG_CONFIG; +<<<<<<< HEAD +======= +import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.ImmutableSet; +import io.airlift.log.Logger; +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) import io.trino.jdbc.TrinoDriver; import io.trino.spi.TrinoException; import java.io.File; @@ -35,8 +41,6 @@ import java.util.Properties; import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * This class dynamically register the Catalog managed by Apache Gravitino into Trino using Trino @@ -44,7 +48,7 @@ import org.slf4j.LoggerFactory; */ public class CatalogRegister { - private static final Logger LOG = LoggerFactory.getLogger(CatalogRegister.class); + private static final Logger LOG = Logger.get(CatalogRegister.class); private static final int EXECUTE_QUERY_MAX_RETRIES = 6; private static final int EXECUTE_QUERY_BACKOFF_TIME_SECOND = 5; @@ -64,7 +68,7 @@ public class CatalogRegister { isStarted = statement.execute(command); return isStarted; } catch (Exception e) { - LOG.warn("Trino server is not started: {}", e.getMessage()); + LOG.warn("Trino server is not started: %s", e.getMessage()); return false; } } @@ -104,6 +108,198 @@ public class CatalogRegister { } } +<<<<<<< HEAD +======= + /** + * Builds the JDBC properties used by the internal connection to the Trino coordinator. + * + * <p>The properties derived from the dedicated {@code trino.jdbc.*} configurations are applied + * first, then the raw driver properties configured with the {@code trino.jdbc.properties.} prefix + * are applied on top of them, so that any driver property can be overridden. + * + * @param config the Gravitino configuration + * @return the JDBC properties + */ + @VisibleForTesting + static Properties buildJdbcProperties(GravitinoConfig config) { + boolean sslEnabled = config.isTrinoJdbcSslEnabled(); + String verification = config.getTrinoJdbcSslVerification(); + String truststorePath = config.getTrinoJdbcSslTruststorePath(); + String truststorePassword = config.getTrinoJdbcSslTruststorePassword(); + String truststoreType = config.getTrinoJdbcSslTruststoreType(); + String keystorePath = config.getTrinoJdbcSslKeystorePath(); + String keystorePassword = config.getTrinoJdbcSslKeystorePassword(); + String keystoreType = config.getTrinoJdbcSslKeystoreType(); + String roles = config.getTrinoJdbcRoles(); + + validateSslConfig( + sslEnabled, + verification, + truststorePath, + truststorePassword, + truststoreType, + keystorePath, + keystorePassword, + keystoreType); + + Properties properties = new Properties(); + properties.put("user", config.getTrinoUser()); + String password = config.getTrinoPassword(); + if (StringUtils.isNotEmpty(password)) { + properties.put("password", password); + } + + if (sslEnabled) { + properties.put("SSL", "true"); + properties.put("SSLVerification", verification); + if (StringUtils.isNotBlank(truststorePath)) { + properties.put("SSLTrustStorePath", truststorePath); + } + if (StringUtils.isNotEmpty(truststorePassword)) { + properties.put("SSLTrustStorePassword", truststorePassword); + } + if (StringUtils.isNotBlank(truststoreType)) { + properties.put("SSLTrustStoreType", truststoreType); + } + if (StringUtils.isNotBlank(keystorePath)) { + properties.put("SSLKeyStorePath", keystorePath); + } + if (StringUtils.isNotEmpty(keystorePassword)) { + properties.put("SSLKeyStorePassword", keystorePassword); + } + if (StringUtils.isNotBlank(keystoreType)) { + properties.put("SSLKeyStoreType", keystoreType); + } + } + + if (StringUtils.isNotBlank(roles)) { + properties.put("roles", roles); + } + + Map<String, String> extraProperties = config.getTrinoJdbcExtraProperties(); + if (!extraProperties.isEmpty()) { + // Log the names only, the values may contain credentials. + LOG.debug("Applying extra Trino JDBC properties: %s", extraProperties.keySet()); + extraProperties.keySet().stream() + .filter(key -> key.startsWith("SSL") && properties.containsKey(key)) + .forEach( + key -> + LOG.warn( + "Extra Trino JDBC property '%s' overrides the TLS setting derived from the " + + "dedicated configuration and is applied without validation", + key)); + properties.putAll(extraProperties); + } + return properties; + } + + private static void validateSslConfig( + boolean sslEnabled, + String verification, + String truststorePath, + String truststorePassword, + String truststoreType, + String keystorePath, + String keystorePassword, + String keystoreType) { + if (!SSL_VERIFICATION_MODES.contains(verification)) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, + String.format( + "Invalid value for config 'trino.jdbc.ssl.verification': expected one of %s, got: %s", + SSL_VERIFICATION_MODES, verification)); + } + + if (!sslEnabled) { + if (!SSL_VERIFICATION_FULL.equals(verification)) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, + "Config 'trino.jdbc.ssl.verification' requires TLS to be enabled either by an HTTPS " + + "'discovery.uri' or by 'trino.jdbc.ssl.enabled=true'"); + } + checkRequiresSslEnabled("trino.jdbc.ssl.truststore.path", truststorePath); + checkRequiresSslEnabled("trino.jdbc.ssl.truststore.password", truststorePassword); + checkRequiresSslEnabled("trino.jdbc.ssl.truststore.type", truststoreType); + checkRequiresSslEnabled("trino.jdbc.ssl.keystore.path", keystorePath); + checkRequiresSslEnabled("trino.jdbc.ssl.keystore.password", keystorePassword); + checkRequiresSslEnabled("trino.jdbc.ssl.keystore.type", keystoreType); + return; + } + + validateKeystoreConfig(verification, keystorePath, keystorePassword, keystoreType); + + if (StringUtils.isBlank(truststorePath)) { + // The driver falls back to the default JVM truststore, which the password and the type of a + // truststore that was never configured have nothing to apply to. + checkRequires( + "trino.jdbc.ssl.truststore.password", + truststorePassword, + "trino.jdbc.ssl.truststore.path"); + checkRequires( + "trino.jdbc.ssl.truststore.type", truststoreType, "trino.jdbc.ssl.truststore.path"); + return; + } + + if (SSL_VERIFICATION_NONE.equals(verification)) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, + "Config 'trino.jdbc.ssl.truststore.path' cannot be used with " + + "'trino.jdbc.ssl.verification' = NONE"); + } + if (!Files.exists(Path.of(truststorePath))) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_MISSING_CONFIG, + String.format( + "The truststore file configured by 'trino.jdbc.ssl.truststore.path' does not exist: %s", + truststorePath)); + } + } + + private static void validateKeystoreConfig( + String verification, String keystorePath, String keystorePassword, String keystoreType) { + if (StringUtils.isBlank(keystorePath)) { + checkRequires( + "trino.jdbc.ssl.keystore.password", keystorePassword, "trino.jdbc.ssl.keystore.path"); + checkRequires("trino.jdbc.ssl.keystore.type", keystoreType, "trino.jdbc.ssl.keystore.path"); + return; + } + if (SSL_VERIFICATION_NONE.equals(verification)) { + // The driver rejects the keystore properties in this combination, so fail with a config + // error here rather than letting it surface as a connection failure. + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, + "Config 'trino.jdbc.ssl.keystore.path' cannot be used with " + + "'trino.jdbc.ssl.verification' = NONE"); + } + if (!Files.exists(Path.of(keystorePath))) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_MISSING_CONFIG, + String.format( + "The keystore file configured by 'trino.jdbc.ssl.keystore.path' does not exist: %s", + keystorePath)); + } + } + + private static void checkRequires(String key, String value, String requiredKey) { + if (StringUtils.isNotEmpty(value)) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, + String.format("Config '%s' requires '%s' to be set", key, requiredKey)); + } + } + + private static void checkRequiresSslEnabled(String key, String value) { + if (StringUtils.isNotEmpty(value)) { + throw new TrinoException( + GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, + String.format( + "Config '%s' requires TLS to be enabled either by an HTTPS 'discovery.uri' or by " + + "'trino.jdbc.ssl.enabled=true'", + key)); + } + } + +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) private String generateCreateCatalogCommand(String name, GravitinoCatalog gravitinoCatalog) throws Exception { return String.format( @@ -151,12 +347,12 @@ public class CatalogRegister { } String createCatalogCommand = generateCreateCatalogCommand(name, catalog); executeSql(createCatalogCommand); - LOG.info("Register catalog {} successfully: {}", name, createCatalogCommand); + LOG.info("Register catalog %s successfully: %s", name, createCatalogCommand); } catch (SQLException e) { throw new TrinoException(GravitinoErrorCode.GRAVITINO_RUNTIME_ERROR, e.getMessage(), e); } catch (Exception e) { String message = String.format("Failed to register catalog %s", name); - LOG.error(message); + LOG.error(e, message); throw new TrinoException(GravitinoErrorCode.GRAVITINO_RUNTIME_ERROR, message, e); } } @@ -186,7 +382,7 @@ public class CatalogRegister { throw e; } catch (Exception e) { failedException = e; - LOG.warn("Failed to execute command: {}", showCatalogCommand, e); + LOG.warn(e, "Failed to execute command: %s", showCatalogCommand); Thread.sleep(EXECUTE_QUERY_BACKOFF_TIME_SECOND * 1000); } } @@ -212,7 +408,7 @@ public class CatalogRegister { throw e; } catch (Exception e) { failedException = e; - LOG.warn("Failed to execute command: {}", sql, e); + LOG.warn(e, "Failed to execute command: %s", sql); Thread.sleep(EXECUTE_QUERY_BACKOFF_TIME_SECOND * 1000); } } @@ -231,15 +427,15 @@ public class CatalogRegister { public void unregisterCatalog(String name) { try { if (!checkCatalogExist(name)) { - LOG.warn("Catalog {} does not exist", name); + LOG.warn("Catalog %s does not exist", name); return; } String dropCatalogCommand = generateDropCatalogCommand(name); executeSql(dropCatalogCommand); - LOG.info("Unregister catalog {} successfully: {}", name, dropCatalogCommand); + LOG.info("Unregister catalog %s successfully: %s", name, dropCatalogCommand); } catch (Exception e) { String message = String.format("Failed to unregister catalog %s", name); - LOG.error(message); + LOG.error(e, message); throw new TrinoException(GravitinoErrorCode.GRAVITINO_RUNTIME_ERROR, message, e); } } @@ -251,7 +447,7 @@ public class CatalogRegister { connection.close(); } } catch (SQLException e) { - LOG.error("Failed to close connection", e); + LOG.error(e, "Failed to close connection"); } } } 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..890293384a 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 @@ -18,6 +18,7 @@ */ package org.apache.gravitino.trino.connector.catalog; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import java.util.HashMap; import java.util.Set; @@ -31,12 +32,10 @@ import org.apache.gravitino.trino.connector.catalog.jdbc.postgresql.PostgreSQLCo import org.apache.gravitino.trino.connector.catalog.jdbc.trino.TrinoClusterConnectorAdapter; import org.apache.gravitino.trino.connector.catalog.memory.MemoryConnectorAdapter; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** This class use to create CatalogConnectorContext instance by given catalog. */ public class DefaultCatalogConnectorFactory implements CatalogConnectorFactory { - private static final Logger LOG = LoggerFactory.getLogger(DefaultCatalogConnectorFactory.class); + private static final Logger LOG = Logger.get(DefaultCatalogConnectorFactory.class); private static final String GLUE_CONNECTOR_PROVIDER_NAME = "glue"; private static final String HIVE_CONNECTOR_PROVIDER_NAME = "hive"; diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/jdbc/mysql/MySQLMetadataAdapter.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/jdbc/mysql/MySQLMetadataAdapter.java index 555170ef46..98f1606f65 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/jdbc/mysql/MySQLMetadataAdapter.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/jdbc/mysql/MySQLMetadataAdapter.java @@ -56,7 +56,6 @@ import org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadataAdap import org.apache.gravitino.trino.connector.catalog.jdbc.JdbcColumnDefaultValueConverter; import org.apache.gravitino.trino.connector.metadata.GravitinoColumn; import org.apache.gravitino.trino.connector.metadata.GravitinoTable; -import org.apache.logging.log4j.util.Strings; /** Transforming Apache Gravitino MySQL metadata to Trino. */ public class MySQLMetadataAdapter extends CatalogConnectorMetadataAdapter { @@ -281,7 +280,7 @@ public class MySQLMetadataAdapter extends CatalogConnectorMetadataAdapter { Arrays.stream(index.fieldNames()) .flatMap(Arrays::stream) .collect(Collectors.toUnmodifiableList()); - uniqueKeys.add(String.format("%s:%s", index.name(), Strings.join(columns, ','))); + uniqueKeys.add(String.format("%s:%s", index.name(), StringUtils.join(columns, ','))); break; default: throw new UnsupportedOperationException("Unsupported index type: " + index.type()); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/security/GravitinoAuthProvider.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/security/GravitinoAuthProvider.java index 21e2e25604..d1f17e8cff 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/security/GravitinoAuthProvider.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/security/GravitinoAuthProvider.java @@ -19,6 +19,11 @@ package org.apache.gravitino.trino.connector.security; import com.google.common.base.Preconditions; +<<<<<<< HEAD +======= +import io.airlift.log.Logger; +import io.trino.spi.TrinoException; +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) import io.trino.spi.connector.ConnectorSession; import java.io.File; import java.util.Locale; @@ -29,8 +34,12 @@ import org.apache.gravitino.client.GravitinoAdminClient; import org.apache.gravitino.client.GravitinoClientConfiguration; import org.apache.gravitino.client.KerberosTokenProvider; import org.apache.gravitino.trino.connector.GravitinoConfig; +<<<<<<< HEAD import org.slf4j.Logger; import org.slf4j.LoggerFactory; +======= +import org.apache.gravitino.trino.connector.GravitinoErrorCode; +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) /** * Builds a {@link GravitinoAdminClient} with the appropriate authentication provider based on the @@ -39,7 +48,7 @@ import org.slf4j.LoggerFactory; */ public class GravitinoAuthProvider { - private static final Logger LOG = LoggerFactory.getLogger(GravitinoAuthProvider.class); + private static final Logger LOG = Logger.get(GravitinoAuthProvider.class); /** Authentication type configuration key. */ public static final String AUTH_TYPE_KEY = @@ -287,7 +296,7 @@ public class GravitinoAuthProvider { // Remove leading slash from path if present String normalizedPath = path.startsWith("/") ? path.substring(1) : path; - LOG.info("Initializing OAuth2 token provider with server URI: {}", serverUri); + LOG.info("Initializing OAuth2 token provider with server URI: %s", serverUri); return DefaultOAuth2TokenProvider.builder() .withUri(serverUri) .withCredential(credential) @@ -331,7 +340,7 @@ public class GravitinoAuthProvider { kerberosBuilder.withKeyTabFile(keytabFile); } else { LOG.warn( - "No keytab file configured for Kerberos authentication ({}). " + "No keytab file configured for Kerberos authentication (%s). " + "Authentication will fail at runtime unless Kerberos credentials are already " + "present in the current security context.", KERBEROS_KEYTAB_FILE_PATH_KEY); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/AlterCatalogStoredProcedure.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/AlterCatalogStoredProcedure.java index 7126568911..55f479ca56 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/AlterCatalogStoredProcedure.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/AlterCatalogStoredProcedure.java @@ -20,6 +20,7 @@ package org.apache.gravitino.trino.connector.system.storedprocedure; import static io.trino.spi.type.VarcharType.VARCHAR; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import io.trino.spi.block.ArrayBlock; import io.trino.spi.procedure.Procedure; @@ -42,8 +43,6 @@ import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * Stored procedure implementation for altering an existing catalog in Gravitino. @@ -51,7 +50,7 @@ import org.slf4j.LoggerFactory; * <p>This procedure allows setting and removing catalog properties dynamically. */ public class AlterCatalogStoredProcedure extends GravitinoStoredProcedure { - private static final Logger LOG = LoggerFactory.getLogger(AlterCatalogStoredProcedure.class); + private static final Logger LOG = Logger.get(AlterCatalogStoredProcedure.class); private final CatalogConnectorManager catalogConnectorManager; private final String metalake; @@ -147,7 +146,7 @@ public class AlterCatalogStoredProcedure extends GravitinoStoredProcedure { GravitinoErrorCode.GRAVITINO_OPERATION_FAILED, "Update catalog failed due to the reloading process fails"); } - LOG.info("Alter catalog {} in metalake {} successfully.", catalogName, metalake); + LOG.info("Alter catalog %s in metalake %s successfully.", catalogName, metalake); } catch (NoSuchMetalakeException e) { throw new TrinoException( diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/CreateCatalogStoredProcedure.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/CreateCatalogStoredProcedure.java index 0de2383d95..48c67c729c 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/CreateCatalogStoredProcedure.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/CreateCatalogStoredProcedure.java @@ -21,6 +21,7 @@ package org.apache.gravitino.trino.connector.system.storedprocedure; import static io.trino.spi.type.BooleanType.BOOLEAN; import static io.trino.spi.type.VarcharType.VARCHAR; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import io.trino.spi.procedure.Procedure; import io.trino.spi.type.MapType; @@ -37,8 +38,6 @@ import org.apache.gravitino.exceptions.NoSuchMetalakeException; import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * Stored procedure implementation for creating a new catalog in Gravitino. @@ -46,7 +45,7 @@ import org.slf4j.LoggerFactory; * <p>This procedure allows creating a new catalog with specified properties and provider. */ public class CreateCatalogStoredProcedure extends GravitinoStoredProcedure { - private static final Logger LOG = LoggerFactory.getLogger(CreateCatalogStoredProcedure.class); + private static final Logger LOG = Logger.get(CreateCatalogStoredProcedure.class); private final CatalogConnectorManager catalogConnectorManager; private final String metalake; @@ -122,7 +121,7 @@ public class CreateCatalogStoredProcedure extends GravitinoStoredProcedure { "Create catalog failed due to the loading process fails"); } - LOG.info("Create catalog {} in metalake {} successfully.", catalogName, metalake); + LOG.info("Create catalog %s in metalake %s successfully.", catalogName, metalake); } catch (NoSuchMetalakeException e) { throw new TrinoException( diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/DropCatalogStoredProcedure.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/DropCatalogStoredProcedure.java index 85feea0072..c6646de1f0 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/DropCatalogStoredProcedure.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/storedprocedure/DropCatalogStoredProcedure.java @@ -21,6 +21,7 @@ package org.apache.gravitino.trino.connector.system.storedprocedure; import static io.trino.spi.type.BooleanType.BOOLEAN; import static io.trino.spi.type.VarcharType.VARCHAR; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import io.trino.spi.procedure.Procedure; import java.lang.invoke.MethodHandle; @@ -33,8 +34,6 @@ import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * Stored procedure implementation for dropping an existing catalog in Gravitino. @@ -42,7 +41,7 @@ import org.slf4j.LoggerFactory; * <p>This procedure allows dropping a catalog with optional ignoreNotExist flag. */ public class DropCatalogStoredProcedure extends GravitinoStoredProcedure { - private static final Logger LOG = LoggerFactory.getLogger(DropCatalogStoredProcedure.class); + private static final Logger LOG = Logger.get(DropCatalogStoredProcedure.class); private final CatalogConnectorManager catalogConnectorManager; private final String metalake; @@ -103,9 +102,8 @@ public class DropCatalogStoredProcedure extends GravitinoStoredProcedure { "Catalog " + NameIdentifier.of(metalake, catalogName) + " not exists."); } LOG.info( - "Drop catalog {} in metalake {} from server (no local connector) successfully.", - catalogName, - metalake); + "Drop catalog %s in metalake %s from server (no local connector) successfully.", + catalogName, metalake); return; } @@ -122,7 +120,7 @@ public class DropCatalogStoredProcedure extends GravitinoStoredProcedure { "Drop catalog failed due to the reloading process fails"); } - LOG.info("Drop catalog {} in metalake {} successfully.", catalogName, metalake); + LOG.info("Drop catalog %s in metalake %s successfully.", catalogName, metalake); } catch (NoSuchMetalakeException e) { throw new TrinoException( diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/util/json/JsonCodec.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/util/json/JsonCodec.java index 0566858223..c1eb134538 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/util/json/JsonCodec.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/util/json/JsonCodec.java @@ -31,6 +31,7 @@ import com.fasterxml.jackson.datatype.jdk8.Jdk8Module; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import com.fasterxml.jackson.module.paramnames.ParameterNamesModule; import io.airlift.json.RecordAutoDetectModule; +import io.airlift.log.Logger; import io.trino.spi.TrinoException; import io.trino.spi.block.Block; import io.trino.spi.block.BlockEncodingSerde; @@ -56,6 +57,11 @@ import org.apache.gravitino.trino.connector.GravitinoConnectorPluginManager; * Provides functionality for handling various Trino-specific types and plugin classes. */ public class JsonCodec { +<<<<<<< HEAD +======= + + private static final Logger LOG = Logger.get(JsonCodec.class); +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) private static volatile ObjectMapper mapper; private static volatile Type jsonType; @@ -110,7 +116,118 @@ public class JsonCodec { return (BlockEncodingSerde) internalBlockEncodingSerdeClass .getConstructor(blockEncodingManagerClass, TypeManager.class) +<<<<<<< HEAD .newInstance(blockEncodingManagerClass.getConstructor().newInstance(), typeManager); +======= + .newInstance(blockEncodingManager, typeManager); + } + + /** + * Instantiate BlockEncodingManager across Trino/Starburst variants. OSS Trino exposes a public + * no-arg constructor; Starburst replaces it with BlockEncodingManager(FeaturesConfig) (used to + * gate type-specific encodings via feature flags). Newer Trino branches additionally publish a + * {@code Set<BlockEncoding>} variant for Guice multibindings. We probe each known signature in + * order, then fall back to a generic constructor scan that fills unknown reference parameters + * with default values as a last-resort compatibility mechanism. + */ + @VisibleForTesting + static Object instantiateBlockEncodingManager( + Class<?> blockEncodingManagerClass, ClassLoader classLoader) throws Exception { + try { + Object instance = blockEncodingManagerClass.getConstructor().newInstance(); + LOG.debug("Instantiated BlockEncodingManager with its public no-argument constructor"); + return instance; + } catch (NoSuchMethodException ignored) { + // fall through to parameterized variants + } + + try { + Class<?> featuresConfigClass = classLoader.loadClass("io.trino.FeaturesConfig"); + Constructor<?> ctor = blockEncodingManagerClass.getConstructor(featuresConfigClass); + Object featuresConfig = featuresConfigClass.getConstructor().newInstance(); + Object instance = ctor.newInstance(featuresConfig); + LOG.debug("Instantiated BlockEncodingManager with FeaturesConfig"); + return instance; + } catch (NoSuchMethodException | ClassNotFoundException ignored) { + // fall through + } + + try { + Constructor<?> setCtor = blockEncodingManagerClass.getConstructor(Set.class); + Object instance = setCtor.newInstance(Collections.emptySet()); + LOG.debug("Instantiated BlockEncodingManager with an empty BlockEncoding set"); + return instance; + } catch (NoSuchMethodException ignored) { + // fall through to last-resort scan + } + + NoSuchMethodException lastError = null; + for (Constructor<?> ctor : blockEncodingManagerClass.getDeclaredConstructors()) { + if (!ctor.trySetAccessible()) { + continue; + } + + Class<?>[] paramTypes = ctor.getParameterTypes(); + Object[] args = new Object[paramTypes.length]; + boolean resolvable = true; + for (int i = 0; i < paramTypes.length && resolvable; i++) { + if (Set.class.isAssignableFrom(paramTypes[i])) { + args[i] = Collections.emptySet(); + } else if (paramTypes[i].isPrimitive()) { + args[i] = defaultPrimitive(paramTypes[i]); + } else { + Object defaultInstance = tryDefaultInstance(paramTypes[i]); + if (defaultInstance == null) { + resolvable = false; + } else { + args[i] = defaultInstance; + } + } + } + if (!resolvable) { + continue; + } + + try { + Object instance = ctor.newInstance(args); + LOG.debug("Instantiated BlockEncodingManager with fallback constructor %s", ctor); + return instance; + } catch (ReflectiveOperationException | RuntimeException e) { + lastError = + new NoSuchMethodException( + "Failed invoking BlockEncodingManager constructor " + ctor + ": " + e.getMessage()); + } + } + + throw new NoSuchMethodException( + "No usable BlockEncodingManager constructor found in " + + blockEncodingManagerClass.getName() + + (lastError == null ? "" : " (last error: " + lastError.getMessage() + ")")); + } + + private static Object tryDefaultInstance(Class<?> type) { + try { + Constructor<?> ctor = type.getDeclaredConstructor(); + if (!ctor.trySetAccessible()) { + return null; + } + return ctor.newInstance(); + } catch (ReflectiveOperationException | RuntimeException e) { + return null; + } + } + + private static Object defaultPrimitive(Class<?> type) { + if (type == boolean.class) return false; + if (type == char.class) return '\0'; + if (type == byte.class) return (byte) 0; + if (type == short.class) return (short) 0; + if (type == int.class) return 0; + if (type == long.class) return 0L; + if (type == float.class) return 0f; + if (type == double.class) return 0d; + throw new IllegalArgumentException("Unsupported primitive type: " + type); +>>>>>>> 52b8f5341 ([#12634] improvement(trino-connector): Log via io.airlift.log.Logger (#12635)) } static void registerHandleSerializationModule(ObjectMapper objectMapper) {
