plusplusjiajia commented on code in PR #898:
URL: https://github.com/apache/iceberg-cpp/pull/898#discussion_r4110389201
##########
src/iceberg/arrow/s3/arrow_s3_file_io.cc:
##########
@@ -244,62 +287,91 @@ Status ArrowS3FileIO::SetStorageCredentials(
properties[key] = value;
}
ICEBERG_ASSIGN_OR_RAISE(auto fs, BuildArrowS3FileSystem(properties));
- file_io_by_prefix.emplace_back(
- CanonicalizeS3Scheme(credential.prefix),
- std::make_unique<ArrowFileSystemFileIO>(std::move(fs)));
+ delegates.emplace_back(CanonicalizeS3Scheme(credential.prefix),
+
std::make_shared<ArrowFileSystemFileIO>(std::move(fs)));
}
- if (file_io_by_prefix.empty() && !storage_credentials.empty()) {
+ if (delegates.empty() && !storage_credentials.empty()) {
// Silent skipping of every vended credential is hard to diagnose: S3
access
// would proceed with the default credentials and fail only at IO time.
ICEBERG_LOG_WARN(
"None of the {} vended storage credential(s) has an S3-compatible
prefix; "
"S3 access will use the default credentials",
storage_credentials.size());
}
- file_io_by_prefix_ = std::move(file_io_by_prefix);
- storage_credentials_ = storage_credentials;
- return {};
+ return delegates;
}
-ArrowFileSystemFileIO& ArrowS3FileIO::FileIOForPath(std::string_view location)
{
- if (file_io_by_prefix_.empty()) {
- return default_file_io_;
+void ArrowS3FileIO::InstallCredentials(
+ std::vector<StorageCredential>& storage_credentials, DelegatesByPrefix&
delegates) {
+ file_io_by_prefix_.swap(delegates);
+ storage_credentials_.swap(storage_credentials);
+}
+
+std::shared_ptr<ArrowFileSystemFileIO> ArrowS3FileIO::MatchDelegate(
+ const std::shared_ptr<ArrowFileSystemFileIO>& fallback,
+ const DelegatesByPrefix& by_prefix, std::string_view location) {
+ if (by_prefix.empty()) {
+ return fallback;
}
const std::string canonical = CanonicalizeS3Scheme(location);
- ArrowFileSystemFileIO* best = &default_file_io_;
+ auto best = fallback;
size_t best_len = 0;
- for (const auto& [prefix, file_io] : file_io_by_prefix_) {
+ for (const auto& [prefix, file_io] : by_prefix) {
if (prefix.size() > best_len && canonical.starts_with(prefix)) {
- best = file_io.get();
+ best = file_io;
best_len = prefix.size();
}
}
- return *best;
+ return best;
+}
+
+std::shared_ptr<ArrowFileSystemFileIO> ArrowS3FileIO::FileIOForPath(
+ std::string_view location) {
+ std::shared_lock lock(mutex_);
+ return MatchDelegate(default_file_io_, file_io_by_prefix_, location);
}
Result<std::unique_ptr<InputFile>> ArrowS3FileIO::NewInputFile(
std::string file_location) {
- return FileIOForPath(file_location).NewInputFile(std::move(file_location));
+ return FileIOForPath(file_location)->NewInputFile(std::move(file_location));
}
Result<std::unique_ptr<InputFile>> ArrowS3FileIO::NewInputFile(std::string
file_location,
size_t length) {
- return FileIOForPath(file_location).NewInputFile(std::move(file_location),
length);
+ return FileIOForPath(file_location)->NewInputFile(std::move(file_location),
length);
}
Result<std::unique_ptr<OutputFile>> ArrowS3FileIO::NewOutputFile(
std::string file_location) {
- return FileIOForPath(file_location).NewOutputFile(std::move(file_location));
+ return FileIOForPath(file_location)->NewOutputFile(std::move(file_location));
}
Status ArrowS3FileIO::DeleteFile(const std::string& file_location) {
- return FileIOForPath(file_location).DeleteFile(file_location);
+ return FileIOForPath(file_location)->DeleteFile(file_location);
}
Status ArrowS3FileIO::DeleteFiles(const std::vector<std::string>&
file_locations) {
- std::unordered_map<ArrowFileSystemFileIO*, std::vector<std::string>>
locations_by_io;
+ // One snapshot so the whole batch matches the same delegate generation; only
+ // ever a handful of delegates, so a linear scan beats hashing.
+ std::shared_ptr<ArrowFileSystemFileIO> fallback;
+ DelegatesByPrefix by_prefix;
+ {
+ std::shared_lock lock(mutex_);
+ fallback = default_file_io_;
+ by_prefix = file_io_by_prefix_;
+ }
+ std::vector<std::pair<std::shared_ptr<ArrowFileSystemFileIO>,
std::vector<std::string>>>
Review Comment:
@wgtmac Thanks! I'd keep that out of scope here and follow up separately.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]