Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import org.apache.iceberg.io.FileInfo;
import org.apache.iceberg.io.InputFile;
import org.apache.iceberg.io.OutputFile;
import org.apache.iceberg.io.PrefixListingPage;
import org.apache.iceberg.io.SupportsPrefixOperations;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.relocated.com.google.common.collect.Iterables;
Expand Down Expand Up @@ -237,6 +238,16 @@ public Iterable<FileInfo> listPrefix(String prefix) {
return prefixIo.listPrefix(prefix);
}

@Override
public boolean supportsPrefixListingWithDelimiter(String prefix, String delimiter) {
return prefixIo.supportsPrefixListingWithDelimiter(prefix, delimiter);
}

@Override
public Iterable<PrefixListingPage> listPrefixWithDelimiter(String prefix, String delimiter) {
return prefixIo.listPrefixWithDelimiter(prefix, delimiter);
}

@Override
public void deletePrefix(String prefix) {
prefixIo.deletePrefix(prefix);
Expand All @@ -261,6 +272,16 @@ public Iterable<FileInfo> listPrefix(String prefix) {
return delegateFileIO.listPrefix(prefix);
}

@Override
public boolean supportsPrefixListingWithDelimiter(String prefix, String delimiter) {
return delegateFileIO.supportsPrefixListingWithDelimiter(prefix, delimiter);
}

@Override
public Iterable<PrefixListingPage> listPrefixWithDelimiter(String prefix, String delimiter) {
return delegateFileIO.listPrefixWithDelimiter(prefix, delimiter);
}

@Override
public void deletePrefix(String prefix) {
delegateFileIO.deletePrefix(prefix);
Expand Down
49 changes: 49 additions & 0 deletions api/src/main/java/org/apache/iceberg/io/PrefixListingPage.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.iceberg.io;

/** A page in a delimited prefix listing. */
public interface PrefixListingPage {

/** Files that do not contain the delimiter after the listed prefix. */
Iterable<FileInfo> files();

/**
* Common prefixes through the first delimiter after the listed prefix.
*
* <p>Each common prefix includes the delimiter and is a location suitable for a subsequent
* listing operation.
*/
Iterable<String> subPrefixes();

/** Create a page from the given files and common prefixes. */
static PrefixListingPage of(Iterable<FileInfo> files, Iterable<String> subPrefixes) {
return new PrefixListingPage() {
@Override
public Iterable<FileInfo> files() {
return files;
}

@Override
public Iterable<String> subPrefixes() {
return subPrefixes;
}
};
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,43 @@ public interface SupportsPrefixOperations extends FileIO {
*/
Iterable<FileInfo> listPrefix(String prefix);

/**
* Lists files and common prefixes under a prefix, grouped by a delimiter.
*
* <p>A file is returned in {@link PrefixListingPage#files()} when the part of its location after
* {@code prefix} does not contain {@code delimiter}. When the remaining part contains the
* delimiter, the file is not returned directly. Instead, {@link PrefixListingPage#subPrefixes()}
* contains the common prefix through the first occurrence of the delimiter. Common prefixes are
* unique, include the delimiter, and are suitable for use in a subsequent listing operation.
*
* <p>Implementations can restrict the supported delimiters. Callers must use {@link
* #supportsPrefixListingWithDelimiter(String, String)} before calling this method.
*
* @param prefix prefix to list
* @param delimiter non-empty delimiter used to group matching locations
* @return iterable of pages containing files and common prefixes directly below the prefix
* @throws UnsupportedOperationException if prefix listing with the delimiter is not supported
*/
default Iterable<PrefixListingPage> listPrefixWithDelimiter(String prefix, String delimiter) {
throw new UnsupportedOperationException(
String.format("Prefix listing with delimiter '%s' is not supported", delimiter));
}

/**
* Returns whether this implementation supports prefix listing with the given delimiter.
*
* <p>Support can vary by prefix when a FileIO selects a storage implementation from the location
* or when the target has additional restrictions, such as an S3 directory bucket. Callers must
* check support for each prefix and delimiter pair they intend to list.
*
* @param prefix prefix to list
* @param delimiter non-empty delimiter used to group matching locations
* @return {@code true} if prefix listing with the delimiter is supported
*/
default boolean supportsPrefixListingWithDelimiter(String prefix, String delimiter) {
return false;
}

/**
* Delete all files under a prefix.
*
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.iceberg.io;

import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

import java.util.Collections;
import org.junit.jupiter.api.Test;

class TestSupportsPrefixOperations {

@Test
void prefixListingPageContainsFilesAndSubPrefixes() {
FileInfo file = new FileInfo("file:/table/file.parquet", 10L, 20L);
PrefixListingPage page =
PrefixListingPage.of(
Collections.singletonList(file), Collections.singletonList("file:/table/partition/"));

assertThat(page.files()).containsExactly(file);
assertThat(page.subPrefixes()).containsExactly("file:/table/partition/");
}

@Test
void delimitedListingIsUnsupportedByDefault() {
SupportsPrefixOperations io = new TestFileIO();

assertThat(io.supportsPrefixListingWithDelimiter("file:/table/", "/")).isFalse();
assertThatThrownBy(() -> io.listPrefixWithDelimiter("file:/table/", "/"))
.isInstanceOf(UnsupportedOperationException.class)
.hasMessage("Prefix listing with delimiter '/' is not supported");
}

private static class TestFileIO implements SupportsPrefixOperations {
@Override
public InputFile newInputFile(String path) {
throw new UnsupportedOperationException();
}

@Override
public OutputFile newOutputFile(String path) {
throw new UnsupportedOperationException();
}

@Override
public void deleteFile(String path) {}

@Override
public Iterable<FileInfo> listPrefix(String prefix) {
return Collections.emptyList();
}

@Override
public void deletePrefix(String prefix) {}
}
}
83 changes: 77 additions & 6 deletions aws/src/main/java/org/apache/iceberg/aws/s3/S3FileIO.java
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@
import org.apache.iceberg.io.FileInfo;
import org.apache.iceberg.io.InputFile;
import org.apache.iceberg.io.OutputFile;
import org.apache.iceberg.io.PrefixListingPage;
import org.apache.iceberg.io.StorageCredential;
import org.apache.iceberg.io.SupportsRecoveryOperations;
import org.apache.iceberg.io.SupportsStorageCredentials;
Expand Down Expand Up @@ -76,10 +77,12 @@
import software.amazon.awssdk.services.s3.model.GetObjectTaggingRequest;
import software.amazon.awssdk.services.s3.model.GetObjectTaggingResponse;
import software.amazon.awssdk.services.s3.model.ListObjectsV2Request;
import software.amazon.awssdk.services.s3.model.ListObjectsV2Response;
import software.amazon.awssdk.services.s3.model.ObjectIdentifier;
import software.amazon.awssdk.services.s3.model.ObjectVersion;
import software.amazon.awssdk.services.s3.model.PutObjectTaggingRequest;
import software.amazon.awssdk.services.s3.model.S3Exception;
import software.amazon.awssdk.services.s3.model.S3Object;
import software.amazon.awssdk.services.s3.model.Tag;
import software.amazon.awssdk.services.s3.model.Tagging;
import software.amazon.awssdk.services.s3.paginators.ListObjectVersionsIterable;
Expand Down Expand Up @@ -343,15 +346,83 @@ public Iterable<FileInfo> listPrefix(String prefix) {
return () ->
client.s3().listObjectsV2Paginator(request).stream()
.flatMap(r -> r.contents().stream())
.map(
o ->
new FileInfo(
String.format("%s://%s/%s", s3uri.scheme(), s3uri.bucket(), o.key()),
o.size(),
o.lastModified().toEpochMilli()))
.map(o -> createFileInfo(s3uri, o))
.iterator();
}

@Override
public Iterable<PrefixListingPage> listPrefixWithDelimiter(String prefix, String delimiter) {
PrefixedS3Client client = clientForStoragePath(prefix);

S3URI uri = new S3URI(prefix, client.s3FileIOProperties().bucketToAccessPointMapping());
if (!supportsPrefixListingWithDelimiter(uri, client.s3FileIOProperties(), delimiter)) {
throw new UnsupportedOperationException(
String.format("Prefix listing with delimiter '%s' is not supported", delimiter));
}

if (uri.useS3DirectoryBucket()
&& client.s3FileIOProperties().isS3DirectoryBucketListPrefixAsDirectory()) {
uri = uri.toDirectoryPath();
}

S3URI s3uri = uri;
ListObjectsV2Request request =
ListObjectsV2Request.builder()
.bucket(s3uri.bucket())
.prefix(s3uri.key())
.delimiter(delimiter)
.build();

return () ->
client.s3().listObjectsV2Paginator(request).stream()
.map(response -> createPrefixListingPage(s3uri, response))
.iterator();
}

@Override
public boolean supportsPrefixListingWithDelimiter(String prefix, String delimiter) {
if (delimiter == null || delimiter.isEmpty()) {
return false;
}

PrefixedS3Client client = clientForStoragePath(prefix);
S3URI uri = new S3URI(prefix, client.s3FileIOProperties().bucketToAccessPointMapping());
return supportsPrefixListingWithDelimiter(uri, client.s3FileIOProperties(), delimiter);
}

private static boolean supportsPrefixListingWithDelimiter(
S3URI uri, S3FileIOProperties properties, String delimiter) {
if (delimiter == null || delimiter.isEmpty()) {
return false;
} else if (!uri.useS3DirectoryBucket()) {
return true;
}

return "/".equals(delimiter)
&& (uri.key().isEmpty()
|| uri.key().endsWith("/")
|| properties.isS3DirectoryBucketListPrefixAsDirectory());
}

private PrefixListingPage createPrefixListingPage(S3URI s3uri, ListObjectsV2Response response) {
List<FileInfo> files = Lists.newArrayList();
List<String> subPrefixes = Lists.newArrayList();
response.contents().forEach(object -> files.add(createFileInfo(s3uri, object)));
response
.commonPrefixes()
.forEach(commonPrefix -> subPrefixes.add(toUri(s3uri, commonPrefix.prefix())));
return PrefixListingPage.of(files, subPrefixes);
}

private FileInfo createFileInfo(S3URI s3uri, S3Object object) {
return new FileInfo(
toUri(s3uri, object.key()), object.size(), object.lastModified().toEpochMilli());
}

private String toUri(S3URI s3uri, String key) {
return String.format("%s://%s/%s", s3uri.scheme(), s3uri.bucket(), key);
}

/**
* This method provides a "best-effort" to delete all objects under the given prefix.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,9 @@
import static org.mockito.Mockito.when;

import com.azure.core.http.rest.PagedIterable;
import com.azure.core.http.rest.PagedResponse;
import com.azure.core.http.rest.Response;
import com.azure.core.util.IterableStream;
import com.azure.storage.blob.models.BlobStorageException;
import com.azure.storage.file.datalake.DataLakeFileClient;
import com.azure.storage.file.datalake.DataLakeFileSystemClient;
Expand All @@ -49,6 +51,7 @@
import org.apache.iceberg.io.FileInfo;
import org.apache.iceberg.io.InputFile;
import org.apache.iceberg.io.OutputFile;
import org.apache.iceberg.io.PrefixListingPage;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.junit.jupiter.api.Test;
Expand Down Expand Up @@ -169,13 +172,48 @@ public void testListPrefixOperations() {

// assert that only files were returned and not directories
FileInfo fileInfo = result.next();
assertThat(fileInfo.location()).isEqualTo("dir/file");
assertThat(fileInfo.location())
.isEqualTo("abfs://container@account.dfs.core.windows.net/dir/file");
assertThat(fileInfo.size()).isEqualTo(123L);
assertThat(fileInfo.createdAtMillis()).isEqualTo(now.toInstant().toEpochMilli());

assertThat(result.hasNext()).isFalse();
}

/** Azurite does not support ADLSv2 directory operations yet so use mocks here. */
@SuppressWarnings("unchecked")
@Test
void listPrefixWithDelimiter() {
String prefix = "abfs://container@account.dfs.core.windows.net/dir";
OffsetDateTime now = OffsetDateTime.now();
PathItem dir =
new PathItem("tag", now, 0L, "group", true, "dir/sub", "owner", "permissions", now, null);
PathItem file =
new PathItem(
"tag", now, 123L, "group", false, "dir/file", "owner", "permissions", now, null);

PagedIterable<PathItem> response = mock(PagedIterable.class);
PagedResponse<PathItem> pageResponse = mock(PagedResponse.class);
when(pageResponse.getElements()).thenReturn(new IterableStream<>(ImmutableList.of(dir, file)));
when(response.iterableByPage())
.thenReturn(new IterableStream<>(ImmutableList.of(pageResponse)));

DataLakeFileSystemClient client = mock(DataLakeFileSystemClient.class);
when(client.listPaths(any(), any())).thenReturn(response);

ADLSFileIO io = spy(new ADLSFileIO());
io.initialize(ImmutableMap.of());
doReturn(client).when(io).client(any(ADLSLocation.class));

Iterable<PrefixListingPage> listing = io.listPrefixWithDelimiter(prefix, "/");
PrefixListingPage page = listing.iterator().next();
assertThat(page.files())
.extracting(FileInfo::location)
.containsExactly("abfs://container@account.dfs.core.windows.net/dir/file");
assertThat(page.subPrefixes())
.containsExactly("abfs://container@account.dfs.core.windows.net/dir/sub/");
}

/** Azurite does not support ADLSv2 directory operations yet so use mocks here. */
@SuppressWarnings("unchecked")
@Test
Expand Down
Loading
Loading