This is an automated email from the ASF dual-hosted git repository.
voonhous pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new c798ccabbbd0 fix(trino): add the reflective HudiTrinoStorage
constructor (#19882)
c798ccabbbd0 is described below
commit c798ccabbbd0307b7de7bbc40290b6629382e516
Author: Ojash Kushwaha <[email protected]>
AuthorDate: Tue Sep 15 22:02:09 2026 +0530
fix(trino): add the reflective HudiTrinoStorage constructor (#19882)
TrinoStorageConfiguration registers HudiTrinoStorage as
hoodie.storage.class, and HoodieStorageUtils.getStorage builds that
class by reflection through a (StoragePath, StorageConfiguration)
constructor that HudiTrinoStorage never had. The read path injects
storage directly and never hit this, but initTable and timeline instant
writes resolve storage from configuration and failed with "Unable to
create HudiTrinoStorage".
Add the constructor. A configuration holds only strings while
HudiTrinoStorage needs the session's TrinoFileSystem, so
TrinoStorageConfiguration gains overloads that carry the file system as
a transient field, and newInstance() and getInline() keep it. The
existing constructors pass null, leaving read-path callers unchanged. A
non-Trino configuration or one without a file system fails with an
IllegalArgumentException.
Tests in TestHudiTrinoStorage cover reflective resolution, derived
configurations, both failure messages, and initTable writing
hoodie.properties and a requested instant through the passed file
system.
---
.../plugin/hudi/storage/HudiTrinoStorage.java | 11 +++
.../hudi/storage/TrinoStorageConfiguration.java | 24 +++++-
.../plugin/hudi/storage/TestHudiTrinoStorage.java | 89 ++++++++++++++++++++++
3 files changed, 122 insertions(+), 2 deletions(-)
diff --git
a/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/HudiTrinoStorage.java
b/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/HudiTrinoStorage.java
index 1edff63e8a7b..cc2906cb5ad6 100644
---
a/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/HudiTrinoStorage.java
+++
b/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/HudiTrinoStorage.java
@@ -54,6 +54,17 @@ public class HudiTrinoStorage
this.fileSystem = fileSystem;
}
+ public HudiTrinoStorage(StoragePath path, StorageConfiguration<?>
storageConf)
+ {
+ if (!(storageConf instanceof TrinoStorageConfiguration trinoConf)) {
+ throw new IllegalArgumentException("Storage configuration for " +
path + " is not a TrinoStorageConfiguration");
+ }
+ TrinoFileSystem configuredFileSystem = trinoConf.getFileSystem()
+ .orElseThrow(() -> new IllegalArgumentException("Storage
configuration for " + path + " carries no file system"));
+ super(trinoConf);
+ this.fileSystem = configuredFileSystem;
+ }
+
public static Location convertToLocation(StoragePath path)
{
return Location.of(path.toString());
diff --git
a/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/TrinoStorageConfiguration.java
b/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/TrinoStorageConfiguration.java
index 48939ae84c5a..e69acecdc64b 100644
---
a/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/TrinoStorageConfiguration.java
+++
b/hudi-trino/src/main/java/io/trino/plugin/hudi/storage/TrinoStorageConfiguration.java
@@ -13,12 +13,14 @@
*/
package io.trino.plugin.hudi.storage;
+import io.trino.filesystem.TrinoFileSystem;
import io.trino.plugin.hudi.io.HudiTrinoIOFactory;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.storage.StorageConfiguration;
import java.util.HashMap;
import java.util.Map;
+import java.util.Optional;
import static
org.apache.hudi.common.config.HoodieStorageConfig.HOODIE_IO_FACTORY_CLASS;
import static
org.apache.hudi.common.config.HoodieStorageConfig.HOODIE_STORAGE_CLASS;
@@ -28,14 +30,32 @@ public class TrinoStorageConfiguration
{
private final Map<String, String> configMap;
+ private final transient TrinoFileSystem fileSystem;
+
public TrinoStorageConfiguration()
{
- this(getDefaultConfigs());
+ this(getDefaultConfigs(), null);
+ }
+
+ public TrinoStorageConfiguration(TrinoFileSystem fileSystem)
+ {
+ this(getDefaultConfigs(), fileSystem);
}
public TrinoStorageConfiguration(Map<String, String> configMap)
+ {
+ this(configMap, null);
+ }
+
+ public TrinoStorageConfiguration(Map<String, String> configMap,
TrinoFileSystem fileSystem)
{
this.configMap = configMap;
+ this.fileSystem = fileSystem;
+ }
+
+ public Optional<TrinoFileSystem> getFileSystem()
+ {
+ return Optional.ofNullable(fileSystem);
}
public static Map<String, String> getDefaultConfigs()
@@ -49,7 +69,7 @@ public class TrinoStorageConfiguration
@Override
public StorageConfiguration newInstance()
{
- return new TrinoStorageConfiguration(new HashMap<>(configMap));
+ return new TrinoStorageConfiguration(new HashMap<>(configMap),
fileSystem);
}
@Override
diff --git
a/hudi-trino/src/test/java/io/trino/plugin/hudi/storage/TestHudiTrinoStorage.java
b/hudi-trino/src/test/java/io/trino/plugin/hudi/storage/TestHudiTrinoStorage.java
index a085a73af551..6f8bfdf1dd3b 100644
---
a/hudi-trino/src/test/java/io/trino/plugin/hudi/storage/TestHudiTrinoStorage.java
+++
b/hudi-trino/src/test/java/io/trino/plugin/hudi/storage/TestHudiTrinoStorage.java
@@ -17,8 +17,16 @@ import io.trino.filesystem.FileEntry;
import io.trino.filesystem.Location;
import io.trino.filesystem.TrinoFileSystem;
import io.trino.filesystem.memory.MemoryFileSystem;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.timeline.HoodieInstant;
+import org.apache.hudi.common.util.HoodieStorageUtils;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StorageConfiguration;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathInfo;
+import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration;
import org.junit.jupiter.api.Test;
import java.io.IOException;
@@ -26,10 +34,14 @@ import java.time.Instant;
import java.util.List;
import java.util.Optional;
+import static
org.apache.hudi.common.config.HoodieStorageConfig.HOODIE_STORAGE_CLASS;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
class TestHudiTrinoStorage
{
+ private static final StoragePath EXTENSION_POINT_PATH = new
StoragePath("memory:///warehouse/table");
+
@Test
void testConvertToPathInfo()
{
@@ -132,6 +144,83 @@ class TestHudiTrinoStorage
assertThat(entries.get(0).getBlockSize()).isEqualTo(20);
}
+ @Test
+ void testStorageResolvesFromConfiguration()
+ {
+ TrinoFileSystem fileSystem = new MemoryFileSystem();
+ StorageConfiguration<?> conf = new
TrinoStorageConfiguration(fileSystem);
+
+ HoodieStorage storage =
HoodieStorageUtils.getStorage(EXTENSION_POINT_PATH, conf);
+
+ assertThat(storage).isInstanceOf(HudiTrinoStorage.class);
+ assertThat(storage.getConf()).isSameAs(conf);
+ }
+
+ @Test
+ void testResolvedStorageCarriesPassedFileSystem()
+ {
+ TrinoFileSystem fileSystem = new MemoryFileSystem();
+
+ HoodieStorage storage = HoodieStorageUtils.getStorage(
+ EXTENSION_POINT_PATH, new
TrinoStorageConfiguration(fileSystem));
+
+ assertThat(((HudiTrinoStorage)
storage).getFileSystem()).isSameAs(fileSystem);
+ }
+
+ @Test
+ void testConfigurationWithoutFileSystemFailsClearly()
+ {
+ assertThatThrownBy(() ->
HoodieStorageUtils.getStorage(EXTENSION_POINT_PATH, new
TrinoStorageConfiguration()))
+ .rootCause()
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("carries no file system");
+ }
+
+ @Test
+ void testForeignConfigurationFailsClearly()
+ {
+ HadoopStorageConfiguration conf = new HadoopStorageConfiguration(new
Configuration());
+ conf.set(HOODIE_STORAGE_CLASS.key(), HudiTrinoStorage.class.getName());
+
+ assertThatThrownBy(() ->
HoodieStorageUtils.getStorage(EXTENSION_POINT_PATH, conf))
+ .rootCause()
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("is not a TrinoStorageConfiguration");
+ }
+
+ @Test
+ void testDerivedConfigurationsKeepTheFileSystem()
+ {
+ TrinoFileSystem fileSystem = new MemoryFileSystem();
+ TrinoStorageConfiguration conf = new
TrinoStorageConfiguration(fileSystem);
+
+ assertThat(HoodieStorageUtils.getStorage(EXTENSION_POINT_PATH,
conf.newInstance()))
+ .isInstanceOf(HudiTrinoStorage.class);
+ assertThat(HoodieStorageUtils.getStorage(EXTENSION_POINT_PATH,
conf.getInline()))
+ .isInstanceOf(HudiTrinoStorage.class);
+ assertThat(((TrinoStorageConfiguration)
conf.newInstance()).getFileSystem()).containsSame(fileSystem);
+ assertThat(((TrinoStorageConfiguration)
conf.getInline()).getFileSystem()).containsSame(fileSystem);
+ }
+
+ @Test
+ void testInitTableWritesThroughExtensionPoint()
+ throws IOException
+ {
+ TrinoFileSystem fileSystem = new MemoryFileSystem();
+
+ HoodieTableMetaClient metaClient =
HoodieTableMetaClient.newTableBuilder()
+ .setTableName("t")
+ .setTableType(HoodieTableType.COPY_ON_WRITE)
+ .initTable(new TrinoStorageConfiguration(fileSystem),
EXTENSION_POINT_PATH);
+ metaClient.getActiveTimeline().createNewInstant(
+ metaClient.createNewInstant(HoodieInstant.State.REQUESTED,
"commit", "001"));
+
+ assertThat(fileSystem.newInputFile(
+ Location.of(EXTENSION_POINT_PATH +
"/.hoodie/hoodie.properties")).exists()).isTrue();
+ assertThat(fileSystem.newInputFile(
+ Location.of(EXTENSION_POINT_PATH +
"/.hoodie/timeline/001.commit.requested")).exists()).isTrue();
+ }
+
private static HudiTrinoStorage createStorageWithFiles()
throws IOException
{