Skip to content

Commit 7c38bd4

Browse files
feat(io): add ResolvingFileIO to resolve FileIO by location scheme and forward vended credentials (#828)
1 parent 9d9ef8c commit 7c38bd4

16 files changed

Lines changed: 676 additions & 162 deletions

‎src/iceberg/CMakeLists.txt‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,7 @@ set(ICEBERG_SOURCES
8383
partition_field.cc
8484
partition_spec.cc
8585
partition_summary.cc
86+
resolving_file_io.cc
8687
row/arrow_array_wrapper.cc
8788
row/manifest_wrapper.cc
8889
row/partition_values.cc

‎src/iceberg/arrow/s3/arrow_s3_file_io.cc‎

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@
3535
#include "iceberg/arrow/arrow_io_util.h"
3636
#include "iceberg/arrow/arrow_status_internal.h"
3737
#include "iceberg/arrow/s3/s3_properties.h"
38+
#include "iceberg/logging/log_macros.h"
3839
#include "iceberg/util/macros.h"
3940
#include "iceberg/util/string_util.h"
4041

@@ -90,9 +91,11 @@ std::string SplitEndpointScheme(std::string_view endpoint,
9091
return std::string(endpoint);
9192
}
9293

94+
// Location prefixes this FileIO can serve: must cover every scheme
95+
// ResolveFileIOName routes here, or such a credential would be dropped.
9396
bool IsS3FileIOCredentialPrefix(std::string_view prefix) {
9497
return prefix == "s3" || prefix.starts_with("s3://") || prefix.starts_with("s3a://") ||
95-
prefix.starts_with("s3n://");
98+
prefix.starts_with("s3n://") || prefix.starts_with("oss://");
9699
}
97100

98101
} // namespace
@@ -181,6 +184,7 @@ Result<std::shared_ptr<::arrow::fs::FileSystem>> BuildArrowS3FileSystem(
181184
return std::shared_ptr<::arrow::fs::FileSystem>(std::move(fs));
182185
}
183186

187+
// Keep in sync with ResolveFileIOName (resolving_file_io.cc).
184188
std::string CanonicalizeS3Scheme(std::string_view location) {
185189
for (std::string_view scheme : {"s3a://", "s3n://", "oss://"}) {
186190
if (location.starts_with(scheme)) {
@@ -235,10 +239,11 @@ Status ArrowS3FileIO::SetStorageCredentials(
235239
// TODO(gangwu): Refresh vended credentials via credentials.uri before tokens expire.
236240
for (const auto& credential : storage_credentials) {
237241
ICEBERG_RETURN_UNEXPECTED(credential.Validate());
242+
// A server may vend credentials for several storage systems at once;
243+
// non-S3 prefixes are skipped, not rejected (Java S3FileIO filters
244+
// credentials by the "s3" prefix).
238245
if (!IsS3FileIOCredentialPrefix(credential.prefix)) {
239-
return NotSupported(
240-
"Storage credential prefix '{}' is unsupported by Arrow S3 FileIO",
241-
credential.prefix);
246+
continue;
242247
}
243248
auto properties = default_properties_;
244249
for (const auto& [key, value] : credential.config) {
@@ -249,6 +254,14 @@ Status ArrowS3FileIO::SetStorageCredentials(
249254
CanonicalizeS3Scheme(credential.prefix),
250255
std::make_unique<ArrowFileSystemFileIO>(std::move(fs)));
251256
}
257+
if (file_io_by_prefix.empty() && !storage_credentials.empty()) {
258+
// Silent skipping of every vended credential is hard to diagnose: S3 access
259+
// would proceed with the default credentials and fail only at IO time.
260+
ICEBERG_LOG_WARN(
261+
"None of the {} vended storage credential(s) has an S3-compatible prefix; "
262+
"S3 access will use the default credentials",
263+
storage_credentials.size());
264+
}
252265
file_io_by_prefix_ = std::move(file_io_by_prefix);
253266
storage_credentials_ = storage_credentials;
254267
return {};

‎src/iceberg/catalog/rest/rest_file_io.cc‎

Lines changed: 6 additions & 66 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,6 @@
2121

2222
#include <string>
2323
#include <unordered_map>
24-
#include <utility>
2524
#include <vector>
2625

2726
#include "iceberg/catalog/rest/types.h"
@@ -33,11 +32,6 @@ namespace iceberg::rest {
3332

3433
namespace {
3534

36-
bool IsBuiltinImpl(std::string_view io_impl) {
37-
return io_impl == FileIORegistry::kArrowLocalFileIO ||
38-
io_impl == FileIORegistry::kArrowS3FileIO;
39-
}
40-
4135
std::unordered_map<std::string, std::string> MergeFileIOProperties(
4236
const std::unordered_map<std::string, std::string>& catalog_config,
4337
const std::unordered_map<std::string, std::string>& table_config) {
@@ -50,56 +44,13 @@ std::unordered_map<std::string, std::string> MergeFileIOProperties(
5044

5145
} // namespace
5246

53-
Result<BuiltinFileIOKind> DetectBuiltinFileIO(std::string_view location) {
54-
const auto pos = location.find("://");
55-
if (pos == std::string_view::npos) {
56-
return BuiltinFileIOKind::kArrowLocal;
57-
}
58-
59-
const auto scheme = location.substr(0, pos);
60-
if (scheme == "file") {
61-
return BuiltinFileIOKind::kArrowLocal;
62-
}
63-
if (scheme == "s3" || scheme == "s3a" || scheme == "s3n") {
64-
return BuiltinFileIOKind::kArrowS3;
65-
}
66-
67-
return NotSupported("URI scheme '{}' is not supported for automatic FileIO resolution",
68-
scheme);
69-
}
70-
71-
std::string_view BuiltinFileIOName(BuiltinFileIOKind kind) {
72-
switch (kind) {
73-
case BuiltinFileIOKind::kArrowLocal:
74-
return FileIORegistry::kArrowLocalFileIO;
75-
case BuiltinFileIOKind::kArrowS3:
76-
return FileIORegistry::kArrowS3FileIO;
77-
}
78-
std::unreachable();
79-
}
80-
8147
Result<std::unique_ptr<FileIO>> MakeCatalogFileIO(const RestCatalogProperties& config) {
8248
std::string io_impl = config.Get(RestCatalogProperties::kIOImpl);
83-
std::string warehouse = config.Get(RestCatalogProperties::kWarehouse);
84-
8549
if (io_impl.empty()) {
86-
if (warehouse.empty()) {
87-
return InvalidArgument(R"("{}" or "{}" property is required to create FileIO)",
88-
RestCatalogProperties::kIOImpl.key(),
89-
RestCatalogProperties::kWarehouse.key());
90-
}
91-
ICEBERG_ASSIGN_OR_RAISE(const auto detected_kind, DetectBuiltinFileIO(warehouse));
92-
io_impl = std::string(BuiltinFileIOName(detected_kind));
93-
}
94-
95-
if (!warehouse.empty() && IsBuiltinImpl(io_impl)) {
96-
ICEBERG_ASSIGN_OR_RAISE(const auto detected_kind, DetectBuiltinFileIO(warehouse));
97-
const auto detected_name = BuiltinFileIOName(detected_kind);
98-
if (io_impl != detected_name) {
99-
return InvalidArgument(
100-
R"("io-impl" value '{}' is incompatible with warehouse '{}')", io_impl,
101-
warehouse);
102-
}
50+
// Resolve the FileIO per file-path scheme instead of guessing from
51+
// `warehouse`, which is often a logical identifier rather than a storage
52+
// URI (Java defaults to ResolvingFileIO likewise).
53+
io_impl = std::string(FileIORegistry::kResolvingFileIO);
10354
}
10455

10556
// TODO(gangwu): Support Java-style customized FileIO creation flows instead of
@@ -112,19 +63,8 @@ Result<std::unique_ptr<FileIO>> MakeTableFileIO(
11263
const std::unordered_map<std::string, std::string>& table_config,
11364
const std::vector<StorageCredential>& storage_credentials) {
11465
const auto default_properties = MergeFileIOProperties(catalog_config, table_config);
115-
const auto properties = RestCatalogProperties::FromMap(default_properties);
116-
auto io_impl = properties.Get(RestCatalogProperties::kIOImpl);
117-
if (io_impl.empty()) {
118-
const auto warehouse = properties.Get(RestCatalogProperties::kWarehouse);
119-
if (warehouse.empty()) {
120-
return InvalidArgument(R"("{}" or "{}" property is required to create FileIO)",
121-
RestCatalogProperties::kIOImpl.key(),
122-
RestCatalogProperties::kWarehouse.key());
123-
}
124-
ICEBERG_ASSIGN_OR_RAISE(const auto detected_kind, DetectBuiltinFileIO(warehouse));
125-
io_impl = std::string(BuiltinFileIOName(detected_kind));
126-
}
127-
ICEBERG_ASSIGN_OR_RAISE(auto io, FileIORegistry::Load(io_impl, default_properties));
66+
ICEBERG_ASSIGN_OR_RAISE(
67+
auto io, MakeCatalogFileIO(RestCatalogProperties::FromMap(default_properties)));
12868

12969
if (storage_credentials.empty()) {
13070
return io;

‎src/iceberg/catalog/rest/rest_file_io.h‎

Lines changed: 2 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -22,9 +22,7 @@
2222
/// \file iceberg/catalog/rest/rest_file_io.h
2323
/// \brief Provide helpers to create FileIO instances for REST catalog responses.
2424

25-
#include <cstdint>
2625
#include <memory>
27-
#include <string_view>
2826
#include <unordered_map>
2927
#include <vector>
3028

@@ -37,16 +35,8 @@
3735

3836
namespace iceberg::rest {
3937

40-
enum class BuiltinFileIOKind : uint8_t {
41-
kArrowLocal,
42-
kArrowS3,
43-
};
44-
45-
ICEBERG_REST_EXPORT Result<BuiltinFileIOKind> DetectBuiltinFileIO(
46-
std::string_view location);
47-
48-
ICEBERG_REST_EXPORT std::string_view BuiltinFileIOName(BuiltinFileIOKind kind);
49-
38+
/// \brief Build the catalog FileIO: the configured `io-impl`, or the
39+
/// scheme-resolving FileIO by default.
5040
ICEBERG_REST_EXPORT Result<std::unique_ptr<FileIO>> MakeCatalogFileIO(
5141
const RestCatalogProperties& config);
5242

‎src/iceberg/file_io_registry.cc‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,13 +22,24 @@
2222
#include <mutex>
2323
#include <utility>
2424

25+
#include "iceberg/resolving_file_io.h"
26+
2527
namespace iceberg {
2628

2729
namespace {
2830

2931
struct RegistryState {
3032
std::mutex mutex;
3133
std::unordered_map<std::string, FileIORegistry::Factory> registry;
34+
35+
RegistryState() {
36+
// Always available: the scheme-resolving FileIO lives in the core library.
37+
registry[std::string(FileIORegistry::kResolvingFileIO)] =
38+
[](const std::unordered_map<std::string, std::string>& properties)
39+
-> Result<std::unique_ptr<FileIO>> {
40+
return std::make_unique<ResolvingFileIO>(properties);
41+
};
42+
}
3243
};
3344

3445
RegistryState& State() {

‎src/iceberg/file_io_registry.h‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,8 @@ class ICEBERG_EXPORT FileIORegistry {
4343
public:
4444
static constexpr std::string_view kArrowLocalFileIO = "arrow-fs-local";
4545
static constexpr std::string_view kArrowS3FileIO = "arrow-fs-s3";
46+
/// Always registered; resolves the concrete FileIO per file-path scheme.
47+
static constexpr std::string_view kResolvingFileIO = "resolving-file-io";
4648

4749
/// Factory function type for creating FileIO instances.
4850
using Factory = std::function<Result<std::unique_ptr<FileIO>>(

‎src/iceberg/meson.build‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,7 @@ iceberg_sources = files(
136136
'partition_spec.cc',
137137
'partition_summary.cc',
138138
'puffin_dv_io.cc',
139+
'resolving_file_io.cc',
139140
'row/arrow_array_wrapper.cc',
140141
'row/manifest_wrapper.cc',
141142
'row/partition_values.cc',
@@ -329,6 +330,7 @@ install_headers(
329330
'name_mapping.h',
330331
'partition_field.h',
331332
'partition_spec.h',
333+
'resolving_file_io.h',
332334
'result.h',
333335
'schema_field.h',
334336
'schema.h',

‎src/iceberg/resolving_file_io.cc‎

Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
#include "iceberg/resolving_file_io.h"
21+
22+
#include <utility>
23+
24+
#include "iceberg/file_io_registry.h"
25+
#include "iceberg/resolving_file_io_internal.h"
26+
#include "iceberg/util/macros.h"
27+
28+
namespace iceberg {
29+
30+
ResolvingFileIO::ResolvingFileIO(std::unordered_map<std::string, std::string> properties)
31+
: properties_(std::move(properties)) {}
32+
33+
ResolvingFileIO::~ResolvingFileIO() = default;
34+
35+
Result<std::string_view> ResolveFileIOName(std::string_view location) {
36+
const auto pos = location.find("://");
37+
if (pos == std::string_view::npos) {
38+
return FileIORegistry::kArrowLocalFileIO;
39+
}
40+
41+
const auto scheme = location.substr(0, pos);
42+
if (scheme == "file") {
43+
return FileIORegistry::kArrowLocalFileIO;
44+
}
45+
// S3-compatible schemes served by the S3 FileIO (Java: SCHEME_TO_FILE_IO).
46+
// Keep in sync with CanonicalizeS3Scheme in arrow_s3_file_io.cc.
47+
if (scheme == "s3" || scheme == "s3a" || scheme == "s3n" || scheme == "oss") {
48+
return FileIORegistry::kArrowS3FileIO;
49+
}
50+
51+
return NotSupported("URI scheme '{}' is not supported for FileIO resolution", scheme);
52+
}
53+
54+
Result<FileIO*> ResolvingFileIO::FileIOForPath(std::string_view location) {
55+
ICEBERG_ASSIGN_OR_RAISE(const auto name, ResolveFileIOName(location));
56+
57+
std::lock_guard lock(mutex_);
58+
auto it = io_by_name_.find(name);
59+
if (it == io_by_name_.end()) {
60+
ICEBERG_ASSIGN_OR_RAISE(auto io,
61+
FileIORegistry::Load(std::string(name), properties_));
62+
// Forward all credentials; each implementation applies the prefixes it
63+
// understands.
64+
if (!storage_credentials_.empty()) {
65+
if (auto* credentialed = io->AsSupportsStorageCredentials()) {
66+
ICEBERG_RETURN_UNEXPECTED(
67+
credentialed->SetStorageCredentials(storage_credentials_));
68+
}
69+
}
70+
it = io_by_name_.emplace(std::string(name), std::move(io)).first;
71+
}
72+
return it->second.get();
73+
}
74+
75+
Result<std::unique_ptr<InputFile>> ResolvingFileIO::NewInputFile(
76+
std::string file_location) {
77+
ICEBERG_ASSIGN_OR_RAISE(auto* io, FileIOForPath(file_location));
78+
return io->NewInputFile(std::move(file_location));
79+
}
80+
81+
Result<std::unique_ptr<InputFile>> ResolvingFileIO::NewInputFile(
82+
std::string file_location, size_t length) {
83+
ICEBERG_ASSIGN_OR_RAISE(auto* io, FileIOForPath(file_location));
84+
return io->NewInputFile(std::move(file_location), length);
85+
}
86+
87+
Result<std::unique_ptr<OutputFile>> ResolvingFileIO::NewOutputFile(
88+
std::string file_location) {
89+
ICEBERG_ASSIGN_OR_RAISE(auto* io, FileIOForPath(file_location));
90+
return io->NewOutputFile(std::move(file_location));
91+
}
92+
93+
Status ResolvingFileIO::DeleteFile(const std::string& file_location) {
94+
ICEBERG_ASSIGN_OR_RAISE(auto* io, FileIOForPath(file_location));
95+
return io->DeleteFile(file_location);
96+
}
97+
98+
Status ResolvingFileIO::DeleteFiles(const std::vector<std::string>& file_locations) {
99+
std::unordered_map<FileIO*, std::vector<std::string>> locations_by_io;
100+
for (const auto& file_location : file_locations) {
101+
ICEBERG_ASSIGN_OR_RAISE(auto* io, FileIOForPath(file_location));
102+
locations_by_io[io].push_back(file_location);
103+
}
104+
for (auto& [io, locations] : locations_by_io) {
105+
ICEBERG_RETURN_UNEXPECTED(io->DeleteFiles(locations));
106+
}
107+
return {};
108+
}
109+
110+
Status ResolvingFileIO::SetStorageCredentials(
111+
const std::vector<StorageCredential>& storage_credentials) {
112+
// Rebuild delegates lazily with the new credentials. Updating live delegates
113+
// instead would leave the resolver inconsistent if one of them rejected them.
114+
std::lock_guard lock(mutex_);
115+
storage_credentials_ = storage_credentials;
116+
io_by_name_.clear();
117+
return {};
118+
}
119+
120+
const std::vector<StorageCredential>& ResolvingFileIO::credentials() const {
121+
return storage_credentials_;
122+
}
123+
124+
} // namespace iceberg

0 commit comments

Comments
 (0)