This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch branch-1.2.x in repository https://gitbox.apache.org/repos/asf/hudi.git
commit 4c822168c434df41c46179cdb9b7b68803bed8e1 Author: voonhous <[email protected]> AuthorDate: Mon Sep 21 12:07:49 2026 +0800 feat(trino): extract HudiExecutorModule for lakehouse reuse (#20008) Trino's lakehouse connector cannot install HudiModule, so it has to provide the bindings HudiSplitManager needs. The metastore getter calls the package-private HudiMetadata.getMetastore(), so lakehouse falls back to creating a fresh metastore per getSplits call, which bypasses the per-transaction cache. Move the three executors and the transaction-scoped metastore getter into a public HudiExecutorModule that HudiModule installs. Trino's module had this name before RFC-105, so LakehouseHudiModule can go back to binder.install(new HudiExecutorModule()). Closes #20002 (cherry picked from commit 2b802c79a9460fb29da01b59c6ae8ae47bcc7d20) --- .../{HudiModule.java => HudiExecutorModule.java} | 49 ++++--------------- .../main/java/io/trino/plugin/hudi/HudiModule.java | 55 +--------------------- 2 files changed, 9 insertions(+), 95 deletions(-) diff --git a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiModule.java b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiExecutorModule.java similarity index 56% copy from hudi-trino/src/main/java/io/trino/plugin/hudi/HudiModule.java copy to hudi-trino/src/main/java/io/trino/plugin/hudi/HudiExecutorModule.java index 7dec679db63d..039aed2bba1e 100644 --- a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiModule.java +++ b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiExecutorModule.java @@ -17,66 +17,33 @@ import com.google.inject.Binder; import com.google.inject.Key; import com.google.inject.Module; import com.google.inject.Provides; -import com.google.inject.Scopes; import com.google.inject.Singleton; import io.trino.metastore.HiveMetastore; -import io.trino.plugin.base.metrics.FileFormatDataSourceStats; -import io.trino.plugin.base.session.SessionPropertiesProvider; -import io.trino.plugin.hive.HideDeltaLakeTables; -import io.trino.plugin.hive.HiveNodePartitioningProvider; import io.trino.plugin.hive.HiveTransactionHandle; -import io.trino.plugin.hive.parquet.ParquetReaderConfig; -import io.trino.plugin.hive.parquet.ParquetWriterConfig; import io.trino.plugin.hudi.stats.ForHudiTableStatistics; -import io.trino.spi.connector.ConnectorNodePartitioningProvider; -import io.trino.spi.connector.ConnectorPageSourceProvider; -import io.trino.spi.connector.ConnectorSplitManager; import io.trino.spi.security.ConnectorIdentity; import java.util.concurrent.ExecutorService; import java.util.concurrent.ScheduledExecutorService; import java.util.function.BiFunction; -import static com.google.inject.multibindings.Multibinder.newSetBinder; -import static io.airlift.concurrent.Threads.daemonThreadsNamed; -import static io.airlift.configuration.ConfigBinder.configBinder; import static io.airlift.bootstrap.ClosingBinder.closingBinder; +import static io.airlift.concurrent.Threads.daemonThreadsNamed; import static java.util.concurrent.Executors.newCachedThreadPool; import static java.util.concurrent.Executors.newScheduledThreadPool; -import static org.weakref.jmx.guice.ExportBinder.newExporter; -public class HudiModule +/** + * Executors and the transaction-scoped metastore getter that {@link HudiSplitManager} and + * {@link HudiMetadataFactory} need. Installed by {@link HudiModule} and by Trino's lakehouse + * connector, which cannot install {@link HudiModule} itself. Requires {@link HudiConfig} and + * {@link HudiTransactionManager} to be bound. + */ +public class HudiExecutorModule implements Module { @Override public void configure(Binder binder) { - binder.bind(HudiTransactionManager.class).in(Scopes.SINGLETON); - - configBinder(binder).bindConfig(HudiConfig.class); - - binder.bind(boolean.class).annotatedWith(HideDeltaLakeTables.class).toInstance(false); - - newSetBinder(binder, SessionPropertiesProvider.class).addBinding().to(HudiSessionProperties.class).in(Scopes.SINGLETON); - binder.bind(HudiTableProperties.class).in(Scopes.SINGLETON); - - binder.bind(ConnectorSplitManager.class).to(HudiSplitManager.class).in(Scopes.SINGLETON); - binder.bind(ConnectorPageSourceProvider.class).to(HudiPageSourceProvider.class).in(Scopes.SINGLETON); - binder.bind(ConnectorNodePartitioningProvider.class).to(HiveNodePartitioningProvider.class).in(Scopes.SINGLETON); - - configBinder(binder).bindConfig(ParquetReaderConfig.class); - configBinder(binder).bindConfig(ParquetWriterConfig.class); - - binder.bind(HudiMetadataFactory.class).in(Scopes.SINGLETON); - - binder.bind(FileFormatDataSourceStats.class).in(Scopes.SINGLETON); - newExporter(binder).export(FileFormatDataSourceStats.class).withGeneratedName(); - - // HudiCacheKeyProvider is deliberately not bound on release-1.2.1, so Trino's default - // provider is used. It implements the Trino 483 CacheKeyProvider contract, which changed - // after 483; binding it breaks a plugin assembled against a newer Trino at the first cached - // read. Restore the binding when the connector targets a Trino with the new contract. - closingBinder(binder).registerExecutor(Key.get(ExecutorService.class, ForHudiTableStatistics.class)); closingBinder(binder).registerExecutor(Key.get(ExecutorService.class, ForHudiSplitManager.class)); closingBinder(binder).registerExecutor(Key.get(ScheduledExecutorService.class, ForHudiSplitSource.class)); diff --git a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiModule.java b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiModule.java index 7dec679db63d..b4c2071797b7 100644 --- a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiModule.java +++ b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiModule.java @@ -14,35 +14,20 @@ package io.trino.plugin.hudi; import com.google.inject.Binder; -import com.google.inject.Key; import com.google.inject.Module; -import com.google.inject.Provides; import com.google.inject.Scopes; -import com.google.inject.Singleton; -import io.trino.metastore.HiveMetastore; import io.trino.plugin.base.metrics.FileFormatDataSourceStats; import io.trino.plugin.base.session.SessionPropertiesProvider; import io.trino.plugin.hive.HideDeltaLakeTables; import io.trino.plugin.hive.HiveNodePartitioningProvider; -import io.trino.plugin.hive.HiveTransactionHandle; import io.trino.plugin.hive.parquet.ParquetReaderConfig; import io.trino.plugin.hive.parquet.ParquetWriterConfig; -import io.trino.plugin.hudi.stats.ForHudiTableStatistics; import io.trino.spi.connector.ConnectorNodePartitioningProvider; import io.trino.spi.connector.ConnectorPageSourceProvider; import io.trino.spi.connector.ConnectorSplitManager; -import io.trino.spi.security.ConnectorIdentity; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.ScheduledExecutorService; -import java.util.function.BiFunction; import static com.google.inject.multibindings.Multibinder.newSetBinder; -import static io.airlift.concurrent.Threads.daemonThreadsNamed; import static io.airlift.configuration.ConfigBinder.configBinder; -import static io.airlift.bootstrap.ClosingBinder.closingBinder; -import static java.util.concurrent.Executors.newCachedThreadPool; -import static java.util.concurrent.Executors.newScheduledThreadPool; import static org.weakref.jmx.guice.ExportBinder.newExporter; public class HudiModule @@ -77,44 +62,6 @@ public class HudiModule // after 483; binding it breaks a plugin assembled against a newer Trino at the first cached // read. Restore the binding when the connector targets a Trino with the new contract. - closingBinder(binder).registerExecutor(Key.get(ExecutorService.class, ForHudiTableStatistics.class)); - closingBinder(binder).registerExecutor(Key.get(ExecutorService.class, ForHudiSplitManager.class)); - closingBinder(binder).registerExecutor(Key.get(ScheduledExecutorService.class, ForHudiSplitSource.class)); - } - - @Provides - @Singleton - @ForHudiTableStatistics - public ExecutorService createTableStatisticsExecutor(HudiConfig hudiConfig) - { - return newScheduledThreadPool( - hudiConfig.getTableStatisticsExecutorParallelism(), - daemonThreadsNamed("hudi-table-statistics-executor-%s")); - } - - @Provides - @Singleton - @ForHudiSplitManager - public ExecutorService createExecutorService() - { - return newCachedThreadPool(daemonThreadsNamed("hudi-split-manager-%s")); - } - - @Provides - @Singleton - @ForHudiSplitSource - public ScheduledExecutorService createSplitLoaderExecutor(HudiConfig hudiConfig) - { - return newScheduledThreadPool( - hudiConfig.getSplitLoaderParallelism(), - daemonThreadsNamed("hudi-split-loader-%s")); - } - - @Provides - @Singleton - public BiFunction<ConnectorIdentity, HiveTransactionHandle, HiveMetastore> createHiveMetastoreGetter(HudiTransactionManager transactionManager) - { - return (identity, transactionHandle) -> - transactionManager.get(transactionHandle, identity).getMetastore(); + binder.install(new HudiExecutorModule()); } }
