diff --git a/google/cloud/storage/examples/storage_async_samples.cc b/google/cloud/storage/examples/storage_async_samples.cc index c37b9417ee8ea..cdc252b5049b6 100644 --- a/google/cloud/storage/examples/storage_async_samples.cc +++ b/google/cloud/storage/examples/storage_async_samples.cc @@ -26,6 +26,7 @@ #include #include #include +#include #include #include #include @@ -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 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 { + 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> 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> { + auto close_status = co_await writer.Close(); + if (!close_status.ok()) throw std::runtime_error(close_status.message()); + auto [new_writer, new_token] = + (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); + 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 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 const&) { + std::cerr << "AsyncClient::OptimizeWriteLatencyPool() example requires " + "coroutines\n"; +} + void ReadAppendableObjectTail(google::cloud::storage::AsyncClient&, std::vector const&) { std::cerr << "AsyncClient::ReadAppendableObjectTail() example requires " @@ -1492,6 +1576,13 @@ void AutoRun(std::vector 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", {""}, + OptimizeWriteLatencyPool), make_entry("rewrite-object", {""}, RewriteObject), make_entry("resume-rewrite-object", {""}, ResumeRewrite),