-
Notifications
You must be signed in to change notification settings - Fork 470
docs(storage): add zonal bucket pre-warmed writer pool sample #16487
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
NickGoog
wants to merge
5
commits into
googleapis:main
Choose a base branch
from
NickGoog:docs-zonal-bucket-writer-pool
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
Show all changes
5 commits
Select commit
Hold shift + click to select a range
afc7261
docs(storage): add zonal bucket pre-warmed writer pool sample
NickGoog 4e9da63
fix(storage): address review feedback in writer pool sample
NickGoog 4a7e02e
docs(storage): describe Flush() as faster instead of citing latency
NickGoog e9e5dfb
Merge branch 'main' into docs-zonal-bucket-writer-pool
kalragauri 7c925c7
docs(storage): use make_bucket_entry and move args in writer pool sample
NickGoog File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -26,6 +26,7 @@ | |
| #include <algorithm> | ||
| #include <cstdint> | ||
| #include <cstdlib> | ||
| #include <deque> | ||
| #include <fstream> | ||
| #include <iostream> | ||
| #include <map> | ||
|
|
@@ -903,6 +904,83 @@ void FinalizeAppendableObjectUpload(google::cloud::storage::AsyncClient& client, | |
| std::cout << "Finalized object: " << object.DebugString() << "\n"; | ||
| } | ||
|
|
||
| void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient& client, | ||
| std::vector<std::string> const& argv) { | ||
| //! [optimize-write-latency-pool] | ||
| // [START storage_optimize_write_latency_pool] | ||
| namespace gcs = google::cloud::storage; | ||
| auto coro = [](gcs::AsyncClient& client, std::string bucket_name, | ||
| std::string key_prefix, | ||
| int pool_size) -> google::cloud::future<void> { | ||
| std::string const next_object_name = | ||
| key_prefix + "_" + std::to_string(pool_size); | ||
|
|
||
| // 1. Init pool: Sized to ensure pre-warmed writers are always available. | ||
| std::deque<std::pair<gcs::AsyncWriter, gcs::AsyncToken>> pool; | ||
| for (int i = 0; i < pool_size; ++i) { | ||
| auto [writer, token] = (co_await client.StartAppendableObjectUpload( | ||
| gcs::BucketName(bucket_name), | ||
| key_prefix + "_" + std::to_string(i))) | ||
| .value(); | ||
| pool.emplace_back(std::move(writer), std::move(token)); | ||
| } | ||
|
|
||
| // 2. Write: Pop a pre-warmed writer and commit with the faster Flush() | ||
| // instead of Finalize(). | ||
| auto [writer, token] = std::move(pool.front()); | ||
| pool.pop_front(); | ||
| token = (co_await writer.Write(std::move(token), | ||
| gcs::WritePayload("0123456789"))) | ||
| .value(); | ||
| auto flush_status = co_await writer.Flush(); | ||
| if (!flush_status.ok()) throw std::runtime_error(flush_status.message()); | ||
|
|
||
| // 3. Pool maintenance (run asynchronously off the critical write path): | ||
| // Close the used writer without finalizing and refill the pool. | ||
| auto maintain_pool = [](gcs::AsyncClient client, std::string bucket_name, | ||
| std::string next_object_name, | ||
| gcs::AsyncWriter writer) | ||
| -> google::cloud::future<std::pair<gcs::AsyncWriter, gcs::AsyncToken>> { | ||
| auto close_status = co_await writer.Close(); | ||
| if (!close_status.ok()) throw std::runtime_error(close_status.message()); | ||
| auto [new_writer, new_token] = | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit: since bucket_name and next_object_name are passed by value into maintain_pool and not used again, moving them avoids an extra copy. Consider passing |
||
| (co_await client.StartAppendableObjectUpload( | ||
| gcs::BucketName(std::move(bucket_name)), | ||
| std::move(next_object_name))) | ||
| .value(); | ||
| co_return {std::move(new_writer), std::move(new_token)}; | ||
| }; | ||
| auto maintenance_future = | ||
| maintain_pool(client, bucket_name, next_object_name, std::move(writer)); | ||
|
|
||
| // 4. Read: Unfinalized objects are readable after Flush(). | ||
| gcs::ObjectDescriptor descriptor = | ||
| (co_await client.Open(gcs::BucketName(bucket_name), key_prefix + "_0")) | ||
| .value(); | ||
| auto [reader, read_token] = descriptor.Read(0, 10); | ||
|
NickGoog marked this conversation as resolved.
|
||
| std::string contents; | ||
| while (read_token.valid()) { | ||
| auto [payload, t] = (co_await reader.Read(std::move(read_token))).value(); | ||
| read_token = std::move(t); | ||
| for (auto const& buffer : payload.contents()) { | ||
| contents.append(buffer.begin(), buffer.end()); | ||
| } | ||
| } | ||
| std::cout << "Read unfinalized object " << key_prefix << "_0: " << contents | ||
| << "\n"; | ||
|
|
||
| auto [new_writer, new_token] = co_await std::move(maintenance_future); | ||
| pool.emplace_back(std::move(new_writer), std::move(new_token)); | ||
| for (auto& [rem_writer, rem_token] : pool) { | ||
| auto close_status = co_await rem_writer.Close(); | ||
| if (!close_status.ok()) throw std::runtime_error(close_status.message()); | ||
| } | ||
| }; | ||
| // [END storage_optimize_write_latency_pool] | ||
| //! [optimize-write-latency-pool] | ||
| coro(client, argv.at(0), argv.at(1), 3).get(); | ||
| } | ||
|
|
||
| void ReadAppendableObjectTail(google::cloud::storage::AsyncClient& client, | ||
| std::vector<std::string> const& argv) { | ||
| //! [read-appendable-object-tail] | ||
|
|
@@ -1145,6 +1223,12 @@ void FinalizeAppendableObjectUpload(google::cloud::storage::AsyncClient&, | |
| "coroutines\n"; | ||
| } | ||
|
|
||
| void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient&, | ||
| std::vector<std::string> const&) { | ||
| std::cerr << "AsyncClient::OptimizeWriteLatencyPool() example requires " | ||
| "coroutines\n"; | ||
| } | ||
|
|
||
| void ReadAppendableObjectTail(google::cloud::storage::AsyncClient&, | ||
| std::vector<std::string> const&) { | ||
| std::cerr << "AsyncClient::ReadAppendableObjectTail() example requires " | ||
|
|
@@ -1492,6 +1576,13 @@ void AutoRun(std::vector<std::string> const& argv) { | |
| scheduled_for_delete.push_back(std::move(object_name)); | ||
| object_name = examples::MakeRandomObjectName(generator, "object-"); | ||
|
|
||
| std::cout << "Running OptimizeWriteLatencyPool() example" << std::endl; | ||
| OptimizeWriteLatencyPool(client, {bucket_name, object_name}); | ||
| for (int i = 0; i != 4; ++i) { | ||
| scheduled_for_delete.push_back(object_name + "_" + std::to_string(i)); | ||
| } | ||
| object_name = examples::MakeRandomObjectName(generator, "object-"); | ||
|
|
||
| std::cout << "Running ReadAppendableObjectTail() example" << std::endl; | ||
| // Create a dummy object for the tail example to read. In a real | ||
| // application another process would be writing to this object. | ||
|
|
@@ -1657,6 +1748,8 @@ int main(int argc, char* argv[]) try { | |
| PauseAndResumeAppendableUpload), | ||
| make_entry("finalize-appendable-object-upload", {}, | ||
| FinalizeAppendableObjectUpload), | ||
| make_bucket_entry("optimize-write-latency-pool", {"<key-prefix>"}, | ||
| OptimizeWriteLatencyPool), | ||
|
|
||
| make_entry("rewrite-object", {"<destination>"}, RewriteObject), | ||
| make_entry("resume-rewrite-object", {"<destination>"}, ResumeRewrite), | ||
|
|
||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.