diff --git a/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergConfig.java b/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergConfig.java index 676d7c1e6d33..550ca20153d0 100644 --- a/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergConfig.java +++ b/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergConfig.java @@ -14,6 +14,7 @@ package io.trino.plugin.iceberg; import io.airlift.configuration.Config; +import io.airlift.configuration.ConfigDescription; import io.trino.plugin.hive.HiveCompressionCodec; import org.apache.iceberg.FileFormat; @@ -27,6 +28,7 @@ public class IcebergConfig private IcebergFileFormat fileFormat = ORC; private HiveCompressionCodec compressionCodec = GZIP; private boolean useFileSizeFromMetadata = true; + private boolean purgeDataOnTableDrop; @NotNull public FileFormat getFileFormat() @@ -74,4 +76,17 @@ public IcebergConfig setUseFileSizeFromMetadata(boolean useFileSizeFromMetadata) this.useFileSizeFromMetadata = useFileSizeFromMetadata; return this; } + + public boolean isPurgeDataOnTableDrop() + { + return purgeDataOnTableDrop; + } + + @Config("iceberg.delete-files-on-table-drop") + @ConfigDescription("Recursively delete table data and metadata on drop") + public IcebergConfig setPurgeDataOnTableDrop(boolean purgeDataOnTableDrop) + { + this.purgeDataOnTableDrop = purgeDataOnTableDrop; + return this; + } } diff --git a/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadata.java b/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadata.java index 340a1f5923cc..2a94020337d3 100644 --- a/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadata.java +++ b/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadata.java @@ -165,6 +165,7 @@ public class IcebergMetadata private final HdfsEnvironment hdfsEnvironment; private final TypeManager typeManager; private final JsonCodec commitTaskCodec; + private final boolean purgeDataOnDrop; private final Map> snapshotIds = new ConcurrentHashMap<>(); @@ -174,12 +175,14 @@ public IcebergMetadata( HiveMetastore metastore, HdfsEnvironment hdfsEnvironment, TypeManager typeManager, - JsonCodec commitTaskCodec) + JsonCodec commitTaskCodec, + boolean purgeDataOnDrop) { this.metastore = requireNonNull(metastore, "metastore is null"); this.hdfsEnvironment = requireNonNull(hdfsEnvironment, "hdfsEnvironment is null"); this.typeManager = requireNonNull(typeManager, "typeManager is null"); this.commitTaskCodec = requireNonNull(commitTaskCodec, "commitTaskCodec is null"); + this.purgeDataOnDrop = purgeDataOnDrop; } @Override @@ -606,7 +609,14 @@ public Optional getInfo(ConnectorTableHandle tableHandle) public void dropTable(ConnectorSession session, ConnectorTableHandle tableHandle) { IcebergTableHandle handle = (IcebergTableHandle) tableHandle; - metastore.dropTable(new HiveIdentity(session), handle.getSchemaName(), handle.getTableName(), true); + Optional table = metastore.getTable(new HiveIdentity(session), handle.getSchemaName(), handle.getTableName()); + + metastore.dropTable(new HiveIdentity(session), handle.getSchemaName(), handle.getTableName(), false); + if (purgeDataOnDrop && table.isPresent()) { + String tableLocation = table.get().getStorage().getLocation(); + HdfsContext context = new HdfsContext(session, handle.getSchemaName(), handle.getTableName()); + deleteDirRecursive(context, hdfsEnvironment, new Path(tableLocation)); + } } @Override @@ -1074,6 +1084,16 @@ private Map> getMaterializedViewToken(ConnectorSess return viewToken; } + private void deleteDirRecursive(HdfsContext context, HdfsEnvironment hdfsEnvironment, Path path) + { + try { + hdfsEnvironment.getFileSystem(context, path).delete(path, true); + } + catch (IOException | RuntimeException e) { + log.error(e, "Failed to delete table data files in path: " + path.toString()); + } + } + private static class TableToken { // Current Snapshot ID of the table diff --git a/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadataFactory.java b/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadataFactory.java index 64fdfdb911ed..7678f4ce8c0b 100644 --- a/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadataFactory.java +++ b/plugin/trino-iceberg/src/main/java/io/trino/plugin/iceberg/IcebergMetadataFactory.java @@ -28,6 +28,7 @@ public class IcebergMetadataFactory private final HdfsEnvironment hdfsEnvironment; private final TypeManager typeManager; private final JsonCodec commitTaskCodec; + private final boolean purgeTableDataOnDrop; @Inject public IcebergMetadataFactory( @@ -37,23 +38,25 @@ public IcebergMetadataFactory( TypeManager typeManager, JsonCodec commitTaskDataJsonCodec) { - this(metastore, hdfsEnvironment, typeManager, commitTaskDataJsonCodec); + this(metastore, hdfsEnvironment, typeManager, commitTaskDataJsonCodec, config.isPurgeDataOnTableDrop()); } public IcebergMetadataFactory( HiveMetastore metastore, HdfsEnvironment hdfsEnvironment, TypeManager typeManager, - JsonCodec commitTaskCodec) + JsonCodec commitTaskCodec, + boolean purgeTableData) { this.metastore = requireNonNull(metastore, "metastore is null"); this.hdfsEnvironment = requireNonNull(hdfsEnvironment, "hdfsEnvironment is null"); this.typeManager = requireNonNull(typeManager, "typeManager is null"); this.commitTaskCodec = requireNonNull(commitTaskCodec, "commitTaskCodec is null"); + this.purgeTableDataOnDrop = purgeTableData; } public IcebergMetadata create() { - return new IcebergMetadata(metastore, hdfsEnvironment, typeManager, commitTaskCodec); + return new IcebergMetadata(metastore, hdfsEnvironment, typeManager, commitTaskCodec, purgeTableDataOnDrop); } } diff --git a/plugin/trino-iceberg/src/test/java/io/trino/plugin/iceberg/TestIcebergConfig.java b/plugin/trino-iceberg/src/test/java/io/trino/plugin/iceberg/TestIcebergConfig.java index ec8cd479f897..634c2263bcd6 100644 --- a/plugin/trino-iceberg/src/test/java/io/trino/plugin/iceberg/TestIcebergConfig.java +++ b/plugin/trino-iceberg/src/test/java/io/trino/plugin/iceberg/TestIcebergConfig.java @@ -34,7 +34,8 @@ public void testDefaults() assertRecordedDefaults(recordDefaults(IcebergConfig.class) .setFileFormat(ORC) .setCompressionCodec(GZIP) - .setUseFileSizeFromMetadata(true)); + .setUseFileSizeFromMetadata(true) + .setPurgeDataOnTableDrop(false)); } @Test @@ -44,12 +45,14 @@ public void testExplicitPropertyMappings() .put("iceberg.file-format", "Parquet") .put("iceberg.compression-codec", "NONE") .put("iceberg.use-file-size-from-metadata", "false") + .put("iceberg.delete-files-on-table-drop", "true") .build(); IcebergConfig expected = new IcebergConfig() .setFileFormat(PARQUET) .setCompressionCodec(HiveCompressionCodec.NONE) - .setUseFileSizeFromMetadata(false); + .setUseFileSizeFromMetadata(false) + .setPurgeDataOnTableDrop(true); assertFullMapping(properties, expected); }