Skip to content

Commit 22a2c87

Browse files
committed
docs(storage): add zonal bucket pre-warmed writer pool sample
Adds storage_optimize_write_latency_pool sample (region tag: storage_optimize_write_latency_pool) demonstrating a pre-warmed pool of AsyncAppendableObjectWriter instances with finalize_on_close=False to avoid object creation and finalization metadata overhead on the critical write path. Verified with both mock unit tests and live integration testing against a Rapid (zonal) bucket in us-central1-a: Running live Python test against bucket=<zonal-bucket>, prefix=live_py_pool_1790087453494 Python 1. Init pool (3 writers): 862.92 ms Python 2. Write+flush: 82.28 ms Python 4. Read back: b'0123456789', pool size after refill: 3 Ran 1 test in 1.612s OK
1 parent 430a87a commit 22a2c87

3 files changed

Lines changed: 148 additions & 0 deletions

File tree

‎storage/samples/snippets/zonal_buckets/README.md‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,4 +75,12 @@ This snippet downloads a range of bytes from multiple objects concurrently.
7575

7676
```bash
7777
python samples/snippets/zonal_buckets/storage_open_multiple_objects_ranged_read.py --bucket_name <bucket_name> --object_names <object_name_1> <object_name_2>
78+
```
79+
80+
### Optimize write latency with a pre-warmed writer pool
81+
82+
This snippet uses a pre-warmed pool of writers for a zonal bucket.
83+
84+
```bash
85+
python samples/snippets/zonal_buckets/storage_optimize_write_latency_pool.py --bucket_name <bucket_name> --key_prefix <key_prefix>
7886
```
Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
#!/usr/bin/env python
2+
3+
# Copyright 2026 Google Inc. All Rights Reserved.
4+
#
5+
# Licensed under the Apache License, Version 2.0 (the 'License');
6+
# you may not use this file except in compliance with the License.
7+
# You may obtain a copy of the License at
8+
#
9+
# http://www.apache.org/licenses/LICENSE-2.0
10+
#
11+
# Unless required by applicable law or agreed to in writing, software
12+
# distributed under the License is distributed on an "AS IS" BASIS,
13+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
# See the License for the specific language governing permissions and
15+
# limitations under the License.
16+
17+
import argparse
18+
import asyncio
19+
from io import BytesIO
20+
21+
from google.cloud.storage.asyncio.async_appendable_object_writer import (
22+
AsyncAppendableObjectWriter,
23+
)
24+
from google.cloud.storage.asyncio.async_grpc_client import AsyncGrpcClient
25+
from google.cloud.storage.asyncio.async_multi_range_downloader import (
26+
AsyncMultiRangeDownloader,
27+
)
28+
29+
30+
# [START storage_optimize_write_latency_pool]
31+
async def storage_optimize_write_latency_pool(
32+
bucket_name: str, key_prefix: str, pool_size: int = 3, grpc_client=None
33+
):
34+
"""Uses a pre-warmed pool of writers for a zonal bucket.
35+
36+
grpc_client: an existing grpc_client to use, this is only for testing.
37+
"""
38+
# The ID of your GCS zonal bucket
39+
# bucket_name = "your-unique-bucket-name"
40+
41+
# The prefix for your pooled GCS objects
42+
# key_prefix = "pooled-object"
43+
44+
if grpc_client is None:
45+
grpc_client = AsyncGrpcClient()
46+
47+
next_object_name = f"{key_prefix}_{pool_size}"
48+
49+
async def new_prewarmed_writer(name: str) -> AsyncAppendableObjectWriter:
50+
w = AsyncAppendableObjectWriter(
51+
client=grpc_client,
52+
bucket_name=bucket_name,
53+
object_name=name,
54+
generation=0,
55+
)
56+
await w.open()
57+
await w.flush() # Forces 0-byte object creation in the background.
58+
return w
59+
60+
# 1. Init pool: Flushing incurs operation charges, so size the pool
61+
# carefully.
62+
pool = [
63+
await new_prewarmed_writer(f"{key_prefix}_{i}")
64+
for i in range(pool_size)
65+
]
66+
67+
# 2. Write: Pop a pre-warmed writer; append() writes and flushes data
68+
# (~1-2 ms).
69+
writer = pool.pop(0)
70+
await writer.append(b"0123456789")
71+
72+
# 3. Pool maintenance (run asynchronously off the critical write path):
73+
# Close the used writer without finalizing, refill the pool, and discard
74+
# stale writers.
75+
async def maintain_pool(used: AsyncAppendableObjectWriter, next_name: str):
76+
await used.close(finalize_on_close=False)
77+
pool.append(await new_prewarmed_writer(next_name))
78+
79+
maintenance_task = asyncio.create_task(
80+
maintain_pool(writer, next_object_name)
81+
)
82+
83+
# 4. Read: Unfinalized objects are readable after flush().
84+
mrd = AsyncMultiRangeDownloader(
85+
grpc_client, bucket_name, f"{key_prefix}_0"
86+
)
87+
await mrd.open()
88+
buf = BytesIO()
89+
await mrd.download_ranges([(0, 0, buf)])
90+
await mrd.close()
91+
92+
await maintenance_task
93+
for rem in pool:
94+
await rem.close(finalize_on_close=False)
95+
96+
print(
97+
f"Read unfinalized object {key_prefix}_0: "
98+
f"{buf.getvalue().decode('utf-8')}"
99+
)
100+
101+
102+
# [END storage_optimize_write_latency_pool]
103+
104+
105+
if __name__ == "__main__":
106+
parser = argparse.ArgumentParser(
107+
description=__doc__,
108+
formatter_class=argparse.RawDescriptionHelpFormatter,
109+
)
110+
parser.add_argument(
111+
"--bucket_name", help="Your Cloud Storage zonal bucket name."
112+
)
113+
parser.add_argument(
114+
"--key_prefix", help="Prefix for pooled object names."
115+
)
116+
args = parser.parse_args()
117+
118+
asyncio.run(
119+
storage_optimize_write_latency_pool(
120+
bucket_name=args.bucket_name,
121+
key_prefix=args.key_prefix,
122+
)
123+
)

‎storage/samples/snippets/zonal_buckets/zonal_snippets_test.py‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
import storage_open_object_multiple_ranged_read
3333
import storage_open_object_read_full_object
3434
import storage_open_object_single_ranged_read
35+
import storage_optimize_write_latency_pool
3536
import storage_pause_and_resume_appendable_upload
3637
import storage_read_appendable_object_tail
3738

@@ -258,3 +259,19 @@ def test_storage_open_multiple_objects_ranged_read(
258259
blob2 = json_client.bucket(_ZONAL_BUCKET).blob(blob2_name)
259260
blob1.delete()
260261
blob2.delete()
262+
263+
264+
def test_storage_optimize_write_latency_pool(
265+
async_grpc_client, json_client, event_loop, capsys
266+
):
267+
key_prefix = f"test-writer-pool-{uuid.uuid4()}"
268+
event_loop.run_until_complete(
269+
storage_optimize_write_latency_pool.storage_optimize_write_latency_pool(
270+
_ZONAL_BUCKET, key_prefix, pool_size=3, grpc_client=async_grpc_client
271+
)
272+
)
273+
out, _ = capsys.readouterr()
274+
assert f"Read unfinalized object {key_prefix}_0: 0123456789" in out
275+
bucket = json_client.bucket(_ZONAL_BUCKET)
276+
for i in range(4):
277+
bucket.blob(f"{key_prefix}_{i}").delete()

0 commit comments

Comments
 (0)