Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Closed
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
80 changes: 70 additions & 10 deletions cpp/src/arrow/filesystem/gcsfs.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,13 +25,23 @@
#include "arrow/result.h"
#include "arrow/util/checked_cast.h"

#define ARROW_GCS_RETURN_NOT_OK(expr) \
if (!expr.ok()) return internal::ToArrowStatus(expr)

namespace arrow {
namespace fs {
namespace {

namespace gcs = google::cloud::storage;

auto constexpr kSep = '/';
// Change the default upload buffer size. In general, sending larger buffers is more
// efficient with GCS, as each buffer requires a roundtrip to the service. With formatted
// output (when using `operator<<`), keeping a larger buffer in memory before uploading
// makes sense. With unformatted output (the only choice given gcs::io::OutputStream's
// API) it is better to let the caller provide as large a buffer as they want. The GCS C++
// client library will upload this buffer with zero copies if possible.
auto constexpr kUploadBufferSize = 256 * 1024;

struct GcsPath {
std::string full_path;
Expand DownExpand Up@@ -83,18 +93,14 @@ class GcsInputStream : public arrow::io::InputStream {

Result<int64_t> Read(int64_t nbytes, void* out) override {
stream_.read(static_cast<char*>(out), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
return stream_.gcount();
}

Result<std::shared_ptr<Buffer>> Read(int64_t nbytes) override {
ARROW_ASSIGN_OR_RAISE(auto buffer, arrow::AllocateResizableBuffer(nbytes));
stream_.read(reinterpret_cast<char*>(buffer->mutable_data()), nbytes);
if (!stream_.status().ok()) {
return internal::ToArrowStatus(stream_.status());
}
ARROW_GCS_RETURN_NOT_OK(stream_.status());
RETURN_NOT_OK(buffer->Resize(stream_.gcount(), true));
return buffer;
}
Expand All@@ -103,6 +109,43 @@ class GcsInputStream : public arrow::io::InputStream {
mutable gcs::ObjectReadStream stream_;
};

class GcsOutputStream : public arrow::io::OutputStream {
public:
explicit GcsOutputStream(gcs::ObjectWriteStream stream) : stream_(std::move(stream)) {}
~GcsOutputStream() override = default;

Status Close() override {
stream_.Close();
return internal::ToArrowStatus(stream_.last_status());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does last_status also clear the error status or is it sticky? If it's sticky, then a failed Write would also return an error when calling Close?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is sticky. And yes, a failed Write() will make subsequent Close() fail. It is unadvisable to finalize a stream that failed, you don't know what is its state.

}

Result<int64_t> Tell() const override {
if (!stream_) {
return Status::IOError("invalid stream");
}
return tell_;
}

bool closed() const override { return !stream_.IsOpen(); }

Status Write(const void* data, int64_t nbytes) override {
if (stream_.write(reinterpret_cast<const char*>(data), nbytes)) {
tell_ += nbytes;
return Status::OK();
}
return internal::ToArrowStatus(stream_.last_status());
}

Status Flush() override {
stream_.flush();
return Status::OK();
}

private:
gcs::ObjectWriteStream stream_;
int64_t tell_ = 0;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It appears tell_ is never updated anywhere. Should you do it in Write perhaps?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Thanks.

};

} // namespace

google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
Expand All@@ -116,6 +159,7 @@ google::cloud::Options AsGoogleCloudOptions(const GcsOptions& o) {
options.set<google::cloud::UnifiedCredentialsOption>(
google::cloud::MakeInsecureCredentials());
}
options.set<gcs::UploadBufferSizeOption>(kUploadBufferSize);
if (!o.endpoint_override.empty()) {
options.set<gcs::RestEndpointOption>(scheme + "://" + o.endpoint_override);
}
Expand All@@ -140,12 +184,27 @@ class GcsFileSystem::Impl {

Result<std::shared_ptr<io::InputStream>> OpenInputStream(const GcsPath& path) {
auto stream = client_.ReadObject(path.bucket, path.object);
if (!stream.status().ok()) {
return internal::ToArrowStatus(stream.status());
}
ARROW_GCS_RETURN_NOT_OK(stream.status());
return std::make_shared<GcsInputStream>(std::move(stream));
}

Result<std::shared_ptr<io::OutputStream>> OpenOutputStream(
const GcsPath& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
gcs::EncryptionKey encryption_key;
ARROW_ASSIGN_OR_RAISE(encryption_key, internal::ToEncryptionKey(metadata));
gcs::PredefinedAcl predefined_acl;
ARROW_ASSIGN_OR_RAISE(predefined_acl, internal::ToPredefinedAcl(metadata));
gcs::KmsKeyName kms_key_name;
ARROW_ASSIGN_OR_RAISE(kms_key_name, internal::ToKmsKeyName(metadata));
gcs::WithObjectMetadata with_object_metadata;
ARROW_ASSIGN_OR_RAISE(with_object_metadata, internal::ToObjectMetadata(metadata));

auto stream = client_.WriteObject(path.bucket, path.object, encryption_key,
predefined_acl, kms_key_name, with_object_metadata);
ARROW_GCS_RETURN_NOT_OK(stream.last_status());
return std::make_shared<GcsOutputStream>(std::move(stream));
}

private:
static Result<FileInfo> GetFileInfoImpl(const GcsPath& path,
const google::cloud::Status& status,
Expand DownExpand Up@@ -245,7 +304,8 @@ Result<std::shared_ptr<io::RandomAccessFile>> GcsFileSystem::OpenInputFile(

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenOutputStream(
const std::string& path, const std::shared_ptr<const KeyValueMetadata>& metadata) {
return Status::NotImplemented("The GCS FileSystem is not fully implemented");
ARROW_ASSIGN_OR_RAISE(auto p, GcsPath::FromString(path));
return impl_->OpenOutputStream(p, metadata);
}

Result<std::shared_ptr<io::OutputStream>> GcsFileSystem::OpenAppendStream(
Expand Down
131 changes: 131 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -17,9 +17,13 @@

#include "arrow/filesystem/gcsfs_internal.h"

#include <absl/time/time.h> // NOLINT
#include <google/cloud/storage/client.h>

#include <sstream>
#include <unordered_map>

#include "arrow/util/key_value_metadata.h"

namespace arrow {
namespace fs {
Expand DownExpand Up@@ -62,6 +66,133 @@ Status ToArrowStatus(const google::cloud::Status& s) {
return Status::OK();
}

namespace gcs = ::google::cloud::storage;

Result<gcs::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::EncryptionKey{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "encryptionKeyBase64") {
return gcs::EncryptionKey::FromBase64Key(values[i]);
}
}
return gcs::EncryptionKey{};
}

Result<gcs::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::KmsKeyName{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "kmsKeyName") {
return gcs::KmsKeyName(values[i]);
}
}
return gcs::KmsKeyName{};
}

Result<gcs::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::PredefinedAcl{};
}

const auto& keys = metadata->keys();
const auto& values = metadata->values();

for (std::size_t i = 0; i < keys.size(); ++i) {
if (keys[i] == "predefinedAcl") {
return gcs::PredefinedAcl(values[i]);
}
}
return gcs::PredefinedAcl{};
}

Result<gcs::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata) {
if (!metadata) {
return gcs::WithObjectMetadata{};
}

static auto const setters = [] {
using setter = std::function<Status(gcs::ObjectMetadata&, const std::string&)>;
return std::unordered_map<std::string, setter>{
{"Cache-Control",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_cache_control(v);
return Status::OK();
}},
{"Content-Disposition",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_disposition(v);
return Status::OK();
}},
{"Content-Encoding",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_encoding(v);
return Status::OK();
}},
{"Content-Language",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_language(v);
return Status::OK();
}},
{"Content-Type",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_content_type(v);
return Status::OK();
}},
{"customTime",
[](gcs::ObjectMetadata& m, const std::string& v) {
std::string err;
absl::Time t;
if (!absl::ParseTime(absl::RFC3339_full, v, &t, &err)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is absl already an include-time dependency of GCS?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. Some Abseil types are exposed in the public API for google-cloud-cpp.

return Status::Invalid("Error parsing RFC-3339 timestamp: '", v, "': ", err);
}
m.set_custom_time(absl::ToChronoTime(t));
return Status::OK();
}},
{"storageClass",
[](gcs::ObjectMetadata& m, const std::string& v) {
m.set_storage_class(v);
return Status::OK();
}},
{"predefinedAcl",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"encryptionKeyBase64",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
{"kmsKeyName",
[](gcs::ObjectMetadata&, const std::string&) { return Status::OK(); }},
};
}();

const auto& keys = metadata->keys();
const auto& values = metadata->values();

gcs::ObjectMetadata object_metadata;
for (std::size_t i = 0; i < keys.size(); ++i) {
auto it = setters.find(keys[i]);
if (it != setters.end()) {
auto status = it->second(object_metadata, values[i]);
if (!status.ok()) return status;
} else {
object_metadata.upsert_metadata(keys[i], values[i]);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is meant to allow inserting arbitrary metadata strings? Will GCS balk if the user throws some unrecognized metadata keys here?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

GCS accepts (mostly) arbitrary metadata keys:

https://cloud.google.com/storage/docs/metadata#custom-metadata

}
}
return gcs::WithObjectMetadata(std::move(object_metadata));
}

} // namespace internal
} // namespace fs
} // namespace arrow
15 changes: 15 additions & 0 deletions cpp/src/arrow/filesystem/gcsfs_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -18,6 +18,9 @@
#pragma once

#include <google/cloud/status.h>
#include <google/cloud/storage/object_metadata.h>
#include <google/cloud/storage/well_known_headers.h>
#include <google/cloud/storage/well_known_parameters.h>

#include <memory>
#include <string>
Expand All@@ -31,6 +34,18 @@ namespace internal {

Status ToArrowStatus(const google::cloud::Status& s);

Result<google::cloud::storage::EncryptionKey> ToEncryptionKey(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::PredefinedAcl> ToPredefinedAcl(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::KmsKeyName> ToKmsKeyName(
const std::shared_ptr<const KeyValueMetadata>& metadata);

Result<google::cloud::storage::WithObjectMetadata> ToObjectMetadata(
const std::shared_ptr<const KeyValueMetadata>& metadata);

} // namespace internal
} // namespace fs
} // namespace arrow
Loading