Skip to content

Commit df06290

Browse files
authored
docs(samples): pre-warm writer pool channels with flush() after open() (#14537)
`open()` is lazy for new appendable objects: the stream is opened and the 0-byte object is created on the first `flush()` or write. Add an empty `flush()` after `open()` so pooled channels are actually pre-warmed.
1 parent af6a871 commit df06290

1 file changed

Lines changed: 23 additions & 11 deletions

File tree

‎java-storage/samples/snippets/src/main/java/com/example/storage/object/OptimizeWriteLatencyPool.java‎

Lines changed: 23 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import com.google.cloud.storage.ReadProjectionConfigs;
2929
import com.google.cloud.storage.Storage;
3030
import com.google.cloud.storage.StorageOptions;
31+
import java.io.IOException;
3132
import java.nio.ByteBuffer;
3233
import java.nio.charset.StandardCharsets;
3334
import java.util.Queue;
@@ -38,6 +39,25 @@
3839
import java.util.concurrent.TimeUnit;
3940

4041
public class OptimizeWriteLatencyPool {
42+
private static AppendableUploadWriteableByteChannel newPrewarmedChannel(
43+
Storage storage, BlobInfo info, BlobAppendableUploadConfig config) throws IOException {
44+
AppendableUploadWriteableByteChannel channel =
45+
storage.blobAppendableUpload(info, config, Storage.BlobWriteOption.doesNotExist()).open();
46+
// open() is lazy. flush() creates the 0-byte object.
47+
try {
48+
channel.flush();
49+
} catch (IOException e) {
50+
// Close the channel; attach any close error to the flush error.
51+
try {
52+
channel.closeWithoutFinalizing();
53+
} catch (IOException closeException) {
54+
e.addSuppressed(closeException);
55+
}
56+
throw e;
57+
}
58+
return channel;
59+
}
60+
4161
public static void optimizeWriteLatencyPool(String bucketName, String keyPrefix)
4262
throws Exception {
4363
// The ID of your GCS zonal bucket
@@ -56,14 +76,10 @@ public static void optimizeWriteLatencyPool(String bucketName, String keyPrefix)
5676
Queue<AppendableUploadWriteableByteChannel> pool = new ConcurrentLinkedQueue<>();
5777
ExecutorService executor = Executors.newSingleThreadExecutor();
5878
try {
59-
// 1. Init pool: Sized to ensure pre-warmed channels are always available.
79+
// 1. Init pool: Flushing incurs operation charges, so size the pool carefully.
6080
for (int i = 0; i < poolSize; i++) {
6181
BlobInfo info = BlobInfo.newBuilder(bucketName, keyPrefix + "_" + i).build();
62-
// open() establishes the stream and creates the 0-byte object in the background.
63-
pool.add(
64-
storage
65-
.blobAppendableUpload(info, config, Storage.BlobWriteOption.doesNotExist())
66-
.open());
82+
pool.add(newPrewarmedChannel(storage, info, config));
6783
}
6884

6985
// 2. Write: Pop a pre-warmed writer and commit with flush() (~1-2 ms)
@@ -91,11 +107,7 @@ public static void optimizeWriteLatencyPool(String bucketName, String keyPrefix)
91107
() -> {
92108
channel.closeWithoutFinalizing();
93109
BlobInfo nextInfo = BlobInfo.newBuilder(bucketName, nextObjectName).build();
94-
pool.add(
95-
storage
96-
.blobAppendableUpload(
97-
nextInfo, config, Storage.BlobWriteOption.doesNotExist())
98-
.open());
110+
pool.add(newPrewarmedChannel(storage, nextInfo, config));
99111
return null;
100112
});
101113

0 commit comments

Comments
 (0)