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]