anoopj commented on code in PR #18172:
URL: https://github.com/apache/iceberg/pull/18172#discussion_r4146494967


##########
docs/docs/fileio.md:
##########
@@ -39,3 +39,21 @@ Different FileIO implementations are used depending on the 
type of storage. Iceb
 - Object Service Storage (including https)
 - Dell Enterprise Cloud Storage
 - Hadoop (adapts any Hadoop FileSystem implementation)
+
+## Google Cloud Storage FileIO

Review Comment:
   Do we really need to update this docs? It has the risk that it will go 
stale, probably the code should be the authoritiative source. Also none of the 
other `FileIO` implementations seem to be documenting their usage here. 



##########
gcp/src/main/java/org/apache/iceberg/gcp/gcs/GCSOutputStream.java:
##########
@@ -81,26 +96,81 @@ public long getPos() {
 
   @Override
   public void flush() throws IOException {
-    stream.flush();
+    Preconditions.checkState(!closed, "Already closed.");
+    if (stream != null) {
+      stream.flush();
+    }
   }
 
   @Override
   public void write(int b) throws IOException {
+    Preconditions.checkState(!closed, "Already closed.");
     stream.write(b);
     pos += 1;
     writeBytes.increment();
     writeOperations.increment();
+
+    if (!useWriteChannel && pos >= writeThreshold) {
+      switchToWriteChannel();
+    }
   }
 
   @Override
   public void write(byte[] b, int off, int len) throws IOException {
-    stream.write(b, off, len);
-    pos += len;
+    Preconditions.checkState(!closed, "Already closed.");
+    int remaining = len;
+    int offset = off;
+
+    if (!useWriteChannel && pos + remaining >= writeThreshold) {
+      int toThreshold = writeThreshold - (int) pos;
+      if (toThreshold > 0) {
+        stream.write(b, offset, toThreshold);
+        pos += toThreshold;
+        offset += toThreshold;
+        remaining -= toThreshold;
+      }
+      switchToWriteChannel();
+    }
+
+    if (remaining > 0) {
+      stream.write(b, offset, remaining);
+      pos += remaining;
+    }
+
     writeBytes.increment(len);
     writeOperations.increment();
   }
 
-  private void openStream() {
+  /**
+   * Open the existing WriteChannel streaming path. Once {@link 
Storage#writer} succeeds, close()
+   * must not fall through to {@link Storage#create}.
+   */
+  private void switchToWriteChannel() throws IOException {
+    OutputStream channelStream = openWriteChannel();
+    try {
+      buffer.writeTo(channelStream);
+      buffer = null;
+    } catch (IOException e) {
+      try {
+        // Best-effort abort. A failed close can leave a truncated GCS object; 
useWriteChannel
+        // stays true so close() does not also Storage.create the same path.
+        channelStream.close();
+      } catch (IOException closeException) {
+        e.addSuppressed(closeException);
+      }
+      stream = null;
+      throw e;
+    }
+  }
+
+  private OutputStream openWriteChannel() {

Review Comment:
   We are setting `useWriteChannel`  only after storage.writer(...), which can 
throw a `StorageException` (an RTE). If it throws at the switch point, 
`useWriteChannel` stays `false` and the buffer still holds the threshold-byte 
prefix, so the ensuing `close()` falls through to `storage.create(...)` and 
uploads a truncated object ?
   
   I think the fix might be trivial: commit to the channel path before the 
fallible call.
   
   ```
   private void openWriteChannel() {
      useWriteChannel = true;
      stream = newWriteChannelStream();[...]
   }



##########
gcp/src/main/java/org/apache/iceberg/gcp/GCPProperties.java:
##########
@@ -47,6 +47,18 @@ public class GCPProperties implements Serializable {
   public static final String GCS_CHANNEL_READ_CHUNK_SIZE = 
"gcs.channel.read.chunk-size-bytes";
   public static final String GCS_CHANNEL_WRITE_CHUNK_SIZE = 
"gcs.channel.write.chunk-size-bytes";
 
+  /**
+   * Max size for a single-shot GCS upload. Larger objects stream through a 
WriteChannel.
+   *
+   * <p>Default: 8 MiB so typical Iceberg metadata, manifests, and delete 
files use one insert. Data

Review Comment:
   Field docs should say what the value is, not the design rationale (they go 
stale).



##########
gcp/src/main/java/org/apache/iceberg/gcp/GCPProperties.java:
##########
@@ -47,6 +47,18 @@ public class GCPProperties implements Serializable {
   public static final String GCS_CHANNEL_READ_CHUNK_SIZE = 
"gcs.channel.read.chunk-size-bytes";
   public static final String GCS_CHANNEL_WRITE_CHUNK_SIZE = 
"gcs.channel.write.chunk-size-bytes";
 
+  /**
+   * Max size for a single-shot GCS upload. Larger objects stream through a 
WriteChannel.
+   *
+   * <p>Default: 8 MiB so typical Iceberg metadata, manifests, and delete 
files use one insert. Data
+   * files are larger and stream after this prefix. GCS recommends simple 
upload under about 5 MiB;
+   * 8 MiB is a coverage tradeoff and should not be raised without evidence. 
Set to 0 to always use
+   * WriteChannel.
+   */
+  public static final String GCS_WRITE_THRESHOLD_BYTES = 
"gcs.write.threshold-bytes";
+
+  public static final long GCS_WRITE_THRESHOLD_BYTES_DEFAULT = 8L * 1024 * 
1024;

Review Comment:
   This is new public API, so the type is permanent once it ships. Why `long`? 
We are already validating that  it is <= `Integer.MAX_VALUE`. Should probably 
be an int?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to