Hi Samrat, Following up on my previous email, I have updated the FLIP-597[1] based on the previous discussion. I have also submitted a draft PR[2] containing only relevant interfaces as a quick overview. Please have a look when you have time.
If there is no further blocking feedback, I suggest we start voting on this FLIP to get further work for Azure and GCS unblocked. Looking forward to hearing from you soon. Kind regards, Alex [1] https://cwiki.apache.org/confluence/x/9gDuGQ [2] https://github.com/apache/flink/pull/28976 On Wed, 12 Aug 2026 at 17:40, Aleksandr Iushmanov <[email protected]> wrote: > Hi Samrat, > > Thank you for coming back. > > I agree with feedback and here are my thoughts below, please let me know > if this would satisfy the requirements and if so, I will incorporate the > changes in the FLIP document. > > 1. Scope and interface layering > > > Glad we are aligned. I will keep an action item to add cleaner separation > of the interfaces on the FLIP. > > 2. Thread safety and cancellation > > 3. Stream ownership, abort, buffering, and CSE > > > I understand the concern about potentially hanging abort. We discussed > with @rkhachatryan the option of adding `abort` as a separate method > exposed through top level interfaces and this requires changes to @Public > interfaces which is rather invasive. As a middle ground to preserve > compatibility on cancellation paths I suggest we add the thread ownership > check and do close gracefully as long as it is called from the owner thread > and close `raw` sdk stream first without lock if called from any other > thread leading to fast IOExceptions on any read paths. > > This solution is not ideal, because it introduces a very shady contract > for stream termination, but I consider it a lesser of two evils. Please let > me know if this would address the concern. > > 4. Existing open questions > > Thank you, I will update the FLIP accordingly. > > 5. Exception classification > > > I agree, exceptions should be preserved, I will update Azure code to > preserve original exceptions and update the FLIP to reflect this intention. > > 6. Migration and GCS recoverable state > > > I will add a note to main FLIP to require child FLIPs to evaluate > feasibility of keeping recoverable writer state compatible between Hadoop > and Native FS when possible, and disclose limitations when not. > > 7. One additional issue found in the reference implementation > > > This was done intentionally, I haven't found examples in flink code that > would write to stream after `sync` call. All cases of this API use are > within try-with-resources and essentially we either call `close` explicitly > or `sync + close` without doing anything in between. In discussions with > Roman and Piotr we concluded that javadoc on `sync` API is likely > inaccurate. In streaming mode, we persist data remotely frequently (every > blob) and until we publish a file by closing the stream, it shouldn't be > consumed by any reader. To keep it clear and simple, we propose to treat > `sync` and `close` the same way regarding data consumption readiness. > > Also this implementation is not about recoverable writer setup, which has > different semantics regarding persisting "partial" data. I believe that if > we need to persist partial data with the expectation that it can be > consumed after that, we should be using a recoverable writer instead. > Please let me know if this addresses your concern. > > There is one more change that I would like to introduce regarding common > interfaces that proved useful for `re-encyrpt` semantics and have been > missed before on FLIPs, I will publish updated versions soon. But in the > meantime, you can find proposed interfaces changes below: > ``` > public interface ObjectStorageFileSystem { > > /** Lazily enumerates all file paths under {@code prefix}, > recursively. Must be closed. */ > CloseableIterator<Path> pathsList(Path prefix) throws IOException; > > /** Fetches file status with user metadata and ETag populated (one > HEAD per call). */ > RichFileStatus getRichFileStatus(Path path) throws IOException; > > /** > * ETag-guarded atomic rename: moves {@code src} → {@code dst} only if > ETags match. > * Empty {@code dstETag} means dst must not exist. Returns {@code > false} if any guard fired; > * throws {@link IOException} if the storage call itself failed (src > state undefined). > */ > boolean moveVerified(Path src, Path dst, String srcETag, String > dstETag) throws IOException; > } > > /** {@link FileStatus} enriched with cloud-object user metadata and an > ETag. */ > @Internal > @Experimental > public interface RichFileStatus extends FileStatus { > > /** User-defined blob metadata (e.g. Azure {@code > PathProperties.getMetadata()}). */ > Map<String, String> getMetadata(); > > /** ETag of the object at fetch time; used for {@link > ObjectStorageFileSystem#moveVerified}. */ > String getETag(); > } > ``` > > I will work on PR tomorrow to include all interface changes proposed as > part of this FLIP to allow GCS / AWS testing. > > Thank you for being active in testing the changes and driving productive > discussion. > > Kind regards, > Aleksandr Iushmanov > > > On Wed, 29 Jul 2026 at 13:58, Samrat Deb <[email protected]> wrote: > >> Hi Aleksandr, >> >> >> Thank you for the detailed response. I checked out PR #28547 [1], along >> with the existing S3 and GCS. The implementation resolves some of my >> earlier questions, while exposing a few lifecycle points that I think we >> should address before freezing the shared contract. >> >> >> >> 1. Scope and interface layering >> >> >> >> After reading the implementation, I agree that InputStreamOpener, >> >> >> InputStreamExtension, and RawAndWrappedInputStreams have distinct >> >> responsibilities: >> >> >> >> - InputStreamOpener captures the provider-specific SDK operation. >> >> >> - InputStreamExtension composes buffering or decryption. >> >> - RawAndWrappedInputStreams distinguishes the SDK resource from the >> readable >> >> stream. >> - The context objects provide an evolvable contract without changing >> >> >> functional-interface signatures. >> >> >> >> >> I no longer think the context wrappers are a performance concern. >> StreamContextImpl >> >> is allocated once per ObjectStorageInputStream; ReadContext and >> >> RawAndWrappedInputStreams are created only when opening or reopening the >> stream, >> >> not on every read. Plain writes also reuse >> WriteContext.EMPTY_WRITE_CONTEXT. >> >> >> >> I therefore withdraw the suggestion to collapse these interfaces solely >> to reduce >> >> their number. I would keep the separation, but strengthen the lifecycle >> contract >> >> as described below. >> >> >> >> >> I also agree with placing these types in flink-core. This gives all >> filesystem >> >> plugins one canonical type identity and avoids requiring every plugin to >> package >> another base module. I see this as a packaging and maintainability >> decision rather >> >> than a hard classloading requirement. >> >> >> >> >> 2. Thread safety and cancellation >> >> >> >> >> I do not yet have a reproducible benchmark artefact strong enough to >> support my >> >> earlier throughput statement, so I suggest that we remove lock cost from >> the >> decision criteria until such data is available. >> >> >> >> >> The cancellation concern remains independent of lock overhead. >> >> >> >> >> ObjectStorageInputStream.read(byte[], ...) holds its ReentrantLock while >> executing >> >> the potentially blocking wrappedStream.read(...). close() then waits for >> the same >> lock using non-interruptible lock() [2]. >> >> >> >> >> Consequently, if the SDK read does not return, close() cannot reach the >> underlying >> >> resource to terminate it. Using a non-interruptible lock ensures that >> the >> closer >> does not abandon lock acquisition; it does not ensure that cleanup >> completes. >> >> >> >> >> The current input concurrency test only exercises lock contention during >> stream >> >> opening: the first thread blocks inside the opener while holding the >> lock, and a >> >> second reader is then interrupted while waiting to acquire it [3]. No >> test drives >> close() from a cancelling thread while another thread is already blocked >> inside a >> >> read of an open SDK stream. The output side, by contrast, already >> specifies the >> >> analogous ordering: close() blocks until a concurrent write releases the >> lock [4]. >> >> On the read path, that same ordering turns a stuck SDK read into a stuck >> >> >> cancellation. >> >> >> >> >> Provider timeouts reduce the probability of an indefinite wait, but I do >> not think >> >> they should be the correctness mechanism for task cancellation. Their >> values are >> >> provider and configuration-dependent, and forceful slot release does not >> >> >> guarantee prompt cleanup of the underlying HTTP connection. >> >> >> >> >> My preferred contract is therefore: >> >> >> - Normal read, write, seek, and position operations are thread-confined. >> >> >> - Graceful close is ordered with normal operations and may finish >> protocol work. >> - Cancellation has a separate, idempotent abort path that may run >> concurrently and >> >> can terminate the native SDK resource without waiting for the blocked >> >> operation. >> >> >> - After abort, no data from that stream is considered valid or >> consumable. >> >> This is also consistent with the existing FSDataOutputStream contract, >> which says >> >> streams are generally not thread-safe but requires close() to be safe >> when invoked >> concurrently during cancellation [5]. >> >> >> Before accepting the shared implementation, I suggest adding a contract >> test with a >> >> delegate blocked in read() and asserting that cancellation can invoke >> the >> provider >> abort operation without first releasing the read. >> >> >> >> 3. Stream ownership, abort, buffering, and CSE >> >> >> >> The abort requirements I had in mind are: >> >> >> >> - Abort is idempotent and does not drain unread data. >> >> >> - Abort can be invoked during a blocked read. >> >> - Providers retain control over their native termination operation. >> >> >> - Graceful close and abort have explicitly different semantics. >> >> >> - The contract defines ownership and close ordering for the decorated >> stream and >> >> raw SDK resource. >> >> >> - Cleanup attempts preserve the primary exception and attach later >> failures as >> >> suppressed exceptions. >> >> >> >> >> >> The current RawAndWrappedInputStreams contains only two InputStream >> references. >> >> ObjectStorageInputStream closes the wrapper first and closes the raw SDK >> stream >> >> directly only if wrapper close fails [6]. This relies on the wrapper >> owning and >> >> closing the raw stream, but that ownership is described as "typically" >> true rather >> than enforced by the type. >> >> >> >> >> There is also no provider-specific abort operation. Existing native S3 >> >> >> demonstrates why that matters: NativeS3InputStream explicitly calls >> >> ResponseInputStream.abort() before closing the buffered and response >> streams, and >> >> its tests assert abort-before-close behaviour [7]. >> >> >> >> >> >> I suggest evolving the returned pair into an owned handle, or augmenting >> it with >> >> lifecycle callbacks, along these lines: >> >> >> - readable stream >> >> - graceful close >> >> >> - cancellation abort >> >> Buffering remains an independent decoration of the readable stream. The >> ownership >> >> proposal is not intended to remove BufferedInputStream or its >> forward-seek >> optimization. >> >> >> >> >> For CSE, I agree that the encryption-format details belong in FLIP-601. >> The >> >> reference implementation also reinforces the need for distinct >> graceful-close and >> abort semantics: AesGcmDecryptingInputStream.close() drains all unread >> ciphertext >> >> to authenticate the final GCM tag [8]. For a partially consumed large >> object, >> >> cancellation latency can therefore include reading the remainder of the >> object, >> >> not merely waiting for one SDK timeout. >> >> >> >> >> >> A graceful close may reasonably finish authentication. Cancellation >> should abort >> >> without draining because the partially read data will be discarded. >> >> >> >> >> I would also avoid assuming that Flink only reads files sequentially. >> The >> Parquet >> adapter exposes FSDataInputStream.seek() directly [9]; the Avro adapter >> does the >> same and additionally seeks to the checkpointed sync offset during state >> >> >> restoration [10]. This does not require chunked CSE to be solved in >> FLIP-597, but >> >> the shared stream contract should leave room for logical-to-physical >> offset mapping >> >> in FLIP-601. >> >> >> >> >> >> 4. Existing open questions >> >> >> RecoverableWriter >> I agree that FLIP-597 should not introduce another common >> RecoverableWriter >> >> implementation. RecoverableWriter is already the common Flink contract; >> persistence >> protocols and state should remain provider-specific, with shared >> contract >> tests >> >> where useful. >> >> >> GCS is a concrete example. Although GSFileSystem delegates ordinary >> filesystem >> >> operations to Hadoop, its recoverable writer already operates directly >> on >> the GCS >> SDK: >> >> >> >> - GSFileSystem constructs GSBlobStorageImpl from the SDK Storage client. >> >> >> - GSRecoverableWriter supports resume and exposes versioned serializers >> [11]. >> - persist() closes the current immutable component object and records >> the >> component >> >> UUIDs and position. >> >> >> - Commit hierarchically composes component objects because GCS compose >> accepts at >> most 32 sources. >> >> >> - GSRecoverableWriterITCase exercises write, persist, recover, and >> commit >> against a >> real bucket. >> >> >> FLIP-603 can therefore carry this provider-specific implementation >> forward rather >> than introducing a new shared recovery protocol. >> >> >> >> >> PathsCopyingFileSystem >> >> >> I agree that this should be a follow-up. The existing interface already >> represents >> >> the capability. >> >> >> >> The current S3 implementation is tightly coupled to S3TransferManager, >> performs >> >> provider-specific batching, leaves completed and partial destination >> cleanup to its >> caller, and currently does not use the supplied close registry. This is >> useful >> >> implementation experience, but not yet evidence of a reusable >> cross-cloud >> >> >> algorithm. >> >> >> >> >> >> ObjectStorageOperations >> >> >> I agree with the provider-local-first approach. >> >> >> >> PR #28547 already follows it for Azure through the package-private >> >> >> DataLakeStorageOperations interface and its in-memory test >> implementation >> [13]. >> Existing GCS code similarly has GSBlobStorage and a test double for >> writer >> >> operations. >> >> >> >> >> >> For a native GCS filesystem, the required operations would be >> approximately: >> >> - Metadata lookup. >> >> >> - Prefix-and-delimiter listing. >> >> - Ranged reads through ReadChannel. >> >> >> - Conditional object creation using generation preconditions. >> - Server-side copy. >> >> >> - Compose with the 32-source limit. >> >> - Batch deletion. >> >> >> The semantics differ materially from Azure: ordinary GCS buckets have >> inferred >> >> directories, marker objects, copy-and-delete rename, and compose-based >> recovery, >> whereas ADLS Gen2 has real directories and atomic rename. >> >> >> I suggest keeping these adapters provider-local, sharing filesystem >> behaviour tests, >> >> and extracting a common operations interface only after S3, Azure, and >> GCS >> demonstrate a stable semantic intersection. >> >> >> >> >> 5. Exception classification >> >> >> >> I was not proposing an ExceptionClassifyingFileSystem wrapper as part of >> FLIP-597. >> >> I suggest a narrower rule for the initial work: >> >> >> - Filesystem methods expose the established Flink exceptions, such as >> >> >> FileNotFoundException and IOException. >> - When translating an SDK exception, the original throwable is retained >> as the >> >> cause. >> >> - A later child FLIP can define provider-specific classifiers for >> authentication, >> throttling, timeout, cancellation, and precondition failures. >> >> >> >> >> >> The Azure reference implementation currently illustrates why even that >> narrow rule >> >> should be explicit. The 404 branches drop the SDK cause: getFileStatus() >> and >> >> open(), for example, throw FileNotFoundException without retaining the >> >> >> DataLakeStorageException, and the same pattern recurs in listStatus() >> and >> the >> >> internal path-existence checks. The create() path, by contrast, does >> retain the SDK >> >> cause for 409 and 412 responses [14]. >> >> >> >> >> >> I suggest fixing this consistency issue. Always retaining the original >> throwable as >> >> the cause without introducing a cross-cloud exception taxonomy in >> FLIP-597. >> >> >> >> 6. Migration and GCS recoverable state >> >> >> >> I agree that deprecation must be evaluated per filesystem plugin. >> FLIP-597 should >> >> not imply deprecation of flink-fs-hadoop-shaded as a whole because HDFS >> and other >> Hadoop integrations remain outside this effort. >> >> >> >> >> >> For recoverable-state migration, the general warning remains valid: >> state >> is not >> >> automatically portable when two plugins use different upload protocols >> and >> >> serializers. >> >> >> >> >> GCS is more favourable because the existing Hadoop plugin already uses >> the native >> >> compose-based writer, and both its recoverable serializers are at >> version >> 1. The >> state splits across two records [12]: >> >> >> >> >> - The commit recoverable carries the final bucket and object name and >> the >> >> >> component-object UUIDs. >> >> >> - The resume recoverable extends it with the current logical write >> position and >> >> whether the stream was closed. >> >> >> >> >> A native GCS plugin can be bidirectionally compatible if it preserves >> the >> >> >> serializer version and bytes, component naming rules, and commit >> protocol. >> >> >> >> There is one important configuration constraint: the serialized state >> does not >> contain the temporary bucket or entropy mode. >> GSCommitRecoverable.getComponentBlobIds() >> >> reconstructs component paths using the current >> gs.writer.temporary.bucket.name and >> >> gs.filesink.entropy.enabled values [12]. Restoring with different values >> will >> >> resolve different object paths unless the component objects have been >> moved >> >> serializer version and bytes, component naming rules, and commit >> protocol. >> >> There is one important configuration constraint: the serialized state >> does not >> contain the temporary bucket or entropy mode. >> GSCommitRecoverable.getComponentBlobIds() >> reconstructs component paths using the current >> gs.writer.temporary.bucket.name and >> gs.filesink.entropy.enabled values [12]. Restoring with different values >> will >> resolve different object paths unless the component objects have been >> moved >> accordingly. >> >> FLIP-603 should therefore require: >> >> - Identical recoverable serializers and component naming. >> - Matching temporary-bucket and entropy configuration during migration. >> - Roll-forward and rollback integration tests between flink-gs-fs-hadoop >> and >> flink-gcs-fs-native. >> - Documentation that changing those options requires relocating >> in-progress >> components. >> >> I have captured this design in the FLIP-603 draft. >> >> 7. One additional issue found in the reference implementation >> >> ObjectStorageOutputStream.sync() currently calls commitAndClose(). Tests >> >> >> explicitly require subsequent write(), flush(), and getPos() calls to >> fail [15]. >> >> >> >> That is different from existing local and Hadoop streams, where sync() >> persists >> >> data without terminating the stream, and from the FSDataOutputStream >> description of >> sync() as fsync-like. >> >> >> >> >> Object stores may be unable to make an object visible without completing >> the >> >> upload, so a terminal operation may be unavoidable for some SDKs. >> However, I think >> we should explicitly resolve whether: >> >> >> >> >> >> - sync() remains non-terminal, >> >> - the operation is unsupported for this common stream, or >> >> >> - terminal sync becomes a documented object-storage-specific semantic. >> >> >> >> >> >> I would prefer resolving this before treating ObjectStorageOutputStream >> as a >> >> cloud-generic implementation. >> >> >> >> >> >> With these adjustments, my proposed convergence is: >> >> >> >> >> - Keep the context and extension separation. >> >> >> - Keep the shared stream types in flink-core. >> >> - Add explicit graceful-close versus cancellation-abort semantics. >> >> >> - Validate lifecycle behaviour with blocked-read cancellation tests and >> S3/GCS >> adapters. >> >> >> - Keep recovery, bulk copy, storage operations, and exception >> classification >> >> provider-specific or in follow-up FLIPs. >> >> >> - Use FLIP-603 as the second provider validation before extracting more >> common >> >> surface area. >> >> >> >> >> >> Bests >> Samrat >> >> >> >> >> [1] https://github.com/apache/flink/pull/28547 >> [2] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-core/src/main/java/org/apache/flink/core/fs/ObjectStorageInputStream.java#L303-L347 >> >> [3] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-core/src/test/java/org/apache/flink/core/fs/ObjectStorageInputStreamTest.java#L499-L541 >> >> [4] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-core/src/test/java/org/apache/flink/core/fs/ObjectStorageOutputStreamTest.java#L441-L519 >> [5] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-core/src/main/java/org/apache/flink/core/fs/FSDataOutputStream.java#L39-L48 >> >> [6] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-core/src/main/java/org/apache/flink/core/fs/ObjectStorageInputStream.java#L122-L161 >> >> [7] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-s3-fs-native/src/main/java/org/apache/flink/fs/s3native/NativeS3InputStream.j >> ava#L197-L245 >> >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-s3-fs-native/src/test/java/org/apache/flink/fs/s3native/NativeS3InputStreamTe >> st.java#L231-L239 >> [8] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-fs-cse-aes-gcm/src/main/java/org/apache/flink/fs/cse/aes/gcm/AesGcmDecrypting >> InputStream.java#L216-L240 >> [9] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/ParquetInputFile.java#L53- >> L74 >> [10] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-formats/flink-avro/src/main/java/org/apache/flink/formats/avro/utils/FSDataInputStreamWrapper.j >> ava#L53-L56 >> >> >> >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-formats/flink-avro/src/main/java/org/apache/flink/formats/avro/AvroInputFormat.java#L204-L210 >> >> >> >> [11] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-gs-fs-hadoop/src/main/java/org/apache/flink/fs/gs/writer/GSRecoverableWriter >> . >> java#L59-L114 >> [12] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-gs-fs-hadoop/src/main/java/org/apache/flink/fs/gs/writer/GSCommitRecoverable >> . >> java#L61-L90 >> >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-gs-fs-hadoop/src/main/java/org/apache/flink/fs/gs/writer/GSResumeRecoverable >> . >> java#L37-L58 >> [13] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-azure-fs-native/src/main/java/org/apache/flink/fs/azurefs/DataLakeStorageOper >> ations.java#L35-L119 >> [14] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-azure-fs-native/src/main/java/org/apache/flink/fs/azurefs/AzureDataLakeFileSy >> stem.java#L262-L268 >> >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-azure-fs-native/src/main/java/org/apache/flink/fs/azurefs/AzureDataLakeFileSy >> stem.java#L326-L330 >> >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-filesystems/flink-azure-fs-native/src/main/java/org/apache/flink/fs/azurefs/AzureDataLakeFileSy >> stem.java#L531-L537 >> [15] >> >> https://github.com/apache/flink/blob/ed35c2a57a5152503ec914ccfa1298a4f04ad110/flink-core/src/main/java/org/apache/flink/core/fs/ObjectStorageOutputStream.java#L159-L205 >> >> >> >> [16] >> >> https://cwiki.apache.org/confluence/spaces/FLINK/pages/444334267/FLIP-603+FileSystems+GCS+SDK-native+filesystem >> >> >> >> >> >> >> On Mon, Jul 13, 2026 at 9:29 PM Aleksandr Iushmanov <[email protected]> >> wrote: >> >> > Hi Samrat, >> > >> > Thank you for looking into FLIP-597[1] in great detail. And especially >> > for trying it out for GCS/AWS implementations as this is the best >> > way to find early issues with the interface design. >> > >> > Please find my thoughts below: >> > >> > 1. Scope and interface layering >> > >> > *re: layers*: I can probably add a diagram to explain the layers and the >> > flow. Explicitly calling out 3 layers makes sense to me. >> > >> > *re: handle*: Speaking of the interfaces, could you please elaborate on >> > your proposal? As per my current understanding we can discuss the >> > following tradeoffs: >> > a. ReadContext, StreamContext, WriteContext >> > Currently, they wrap only small subset of arguments to pass: >> > >> > ReadContext -> long (position to re-open) >> > StreamContext -> long, long (position, contentLength) >> > WriteContext -> map<string, string> (file metadata) >> > >> > Yes, these wraps could be inlined but the values I >> > see in them are: more explicit contract; better symmetry in the >> > interfaces. From a maintainability point of view, it seems to me that >> > making changes to these wrappers is easier than tracing free passed >> > arguments. Do you have performance concerns around them? >> > >> > As for InputStreamOpener, InputStreamExtension, and >> > RawAndWrappedInputStreams, >> > they all have distinct roles: the opener is the provider-specific SDK >> call, >> > the extension for cross-cutting concerns like CSE[2], and the stream >> pair >> > to separate >> > the readable stream from SDK connection. >> > We could in theory collapse some of them, but imho it would hurt >> > testability >> > and contract separation. Could you please provide more details on >> > how to better outline contracts? >> > >> > I would also like to understand the reason for placing the >> implementations >> > > in >> > > flink-core >> > >> > >> > Main motivation as mentioned in rejected alternatives section [3] is to >> > avoid introduction of yet another dependency that every plugin would >> have >> > to >> > declare alongside core. It made more sense to me to put new interfaces >> > alongside `FileSystem / FSDataInputStream`. >> > >> > 2. Thread safety and cancellation. >> > >> > > locking every read affected the throughput >> > >> > >> > That is a slightly surprising finding. Given how we use >> > streams in flink, access to the reentrant lock would be 99% uncontended >> > (except the cancellation case). And most reads would be 8KB reads for >> the >> > download path. >> > I would imagine that unless you issue single byte reads, latency for >> > acquiring >> > uncontended locks should be negligible relative to the HTTP latency on >> the >> > SDK side (multiple orders of magnitude smaller). Do you have more >> details >> > on a benchmark scenario in which it mattered? >> > >> > during a blocking SDK read >> > >> > IIUC, SDK reads internally are non-interruptible. If you are concerned >> > about >> > locking in `close` non-interruptibly vs interruptibly elsewhere, I don't >> > think >> > It makes a big difference in terms of a hang risk. When stream is closed >> > gracefully, it is closed from the same thread, hence it can't be stuck >> > on the SDK read. So this hang could only ever have an impact during task >> > cancellation. Which is asynchronous and doesn't prevent job >> re-submission. >> > It would also be limited to SDK timeout, which is typically well below >> > default timeout for task cancellation before the forceful slot release. >> > >> > If full thread safety is required >> > >> > >> > Thread safety is required for CSE. >> > The issue stems from the fact that CSE requires verification of >> > read stream, which requires reading until the GCM verification tag. >> > Unless/until >> > we implement chunked encryption, we would have to "read input stream >> fully" >> > to ensure that we can trust the data. Given that we need to read the >> stream >> > on >> > close, it requires full thread safety on reads between main reader and >> > closer. There are only 2 ways I see that could avoid locking on reads: >> > a. We implement chunked encryption CSE (worse performance on sequential >> > reads, >> > but much better random access, which is normally not used in flink). >> With >> > chunked encryption, we could "read ahead" and only serve already >> verified >> > data and eliminate the need to "read on close". >> > b. We clearly separate graceful close (read data is used) from abortion >> > (read data is never used). Abortion case would not require thread >> safety. >> > And >> > graceful close has to be called from the same thread. >> > >> > 3. Stream ownership, CRT, and encryption >> > >> > > S3 may require >> > > aborting a response rather than closing and draining it >> > >> > >> > Draining is only implemented in the CSE extension itself, it doesn't >> happen >> > for standard buffering implementation. The only difference is the >> potential >> > cancellation delay mentioned above if we had a hanging/stuck read. >> > >> > RawAndWrappedInputStreams may therefore need more >> > > explicit rules for ownership, close ordering, and abort behaviour. >> > >> > >> > Could you please provide additional requirements regarding the abort >> > behaviour you have mentioned? >> > >> > This might let each provider adapt its native stream >> > > without introducing an intermediate buffer >> > >> > >> > The default buffering was introduced to provide forward seek >> optimisation. >> > If you do forward skips within the block that is already pulled by SDK, >> > it might entirely ignore the fact that you have data locally and instead >> > open a new range read with the cloud provider. Having a buffering >> extension >> > allows faster in-memory buffer drain instead of extra HTTP requests. >> > >> > If your concern is that you would like to skip buffering, we could add >> > another "no-op" implementation, that essentially gives you a RawAndRaw >> > input stream. In combination with idempotent `close`, I don't see it >> > causing any issues. >> > >> > Random access will require more than metadata >> > > injection >> > >> > >> > For the client side encryption, the largest question is whether we >> > have to go for chunked implementation already or not. Flink FS hot paths >> > don't involve "random access". Files are read from beginning to the end, >> > hence the benefits of encrypting small "chunks" vs the whole file is >> > marginal. >> > Standard format for algorithms like AES-GCM would be >> > [Initialization Vector(IV)][Ciphertext][Verification Tag(GCM)]. In case >> of >> > chunked encryption, it would be [IV1][Chunk1][GCM1][IV2][Chunk2][GCM2] >> etc. >> > Initial implementation I have suggested in FLIP-601 [2] goes for the >> whole >> > file encryption approach as a simpler and slightly faster first >> > implementation. >> > However, if you want to go for chunked encryption, we don't need much >> > extra. >> > Metadata injection gives you "chunk size/IV length/GCM tag length", then >> > remaining work is to implement a couple of position mappers that would >> > track reopen stream not from 0 position but "beginning of the >> > closest chunk". In addition, for chunked encryption, we would swap >> "drain" >> > for "read ahead". >> > >> > I will share the header-injection >> > > and range-read requirements once that discussion concludes. >> > >> > >> > Sounds great, I suggest we move this requirements discussion in >> > FLIP-601 thread. >> > >> > 4. Existing open questions >> > >> > only shared utilities and contract tests? >> > >> > >> > This part, but as we have discussed, cloud specific parts might be >> > too diverse around recoverable writers, so juice probably not worth the >> > squeeze and I would be happy to keep this one out of the scope. >> > >> > For PathsCopyingFileSystem, I believe the existing interface already >> > > provides the common capability >> > >> > >> > In general, I agree. The reason I mention it is that we have provided >> > this optimisation for AWS in many different ways, but we don't really >> > have it supported for other clouds, and I was wondering about where >> > the value of the implementation sits. If it is heavily re-using SDK >> > capabilities with thin logic, then yes, doing it per-cloud makes sense. >> > However, if the added value is in handling concurrency over simple SDK >> > calls, that could be largely cloud agnostic and be easily re-used for >> > other clouds too. >> > >> > Having said that, I would suggest considering this as a follow up >> > improvement rather than the main FLIP scope. >> > >> > demonstrate a stable intersection >> > >> > >> > I am heavy +1 to that, let's see how it looks for all 3 and if there is >> > a lot of code duplication, let's extract it. >> > >> > If the intended interface is narrower >> > >> > >> > Intended interface is certainly narrower than all ObjectStorage >> operations >> > and is naturally limited to operations that FS implementation demands. >> > I have example of such implementation for Azure, and we use >> S3AccessHelper >> > [4] >> > For Native S3, most likely, these 2 will already be very representative, >> > but if we can analyse what would be needed for GCS too, it would be >> > helpful. >> > >> > >> > 5. Exception classification >> > >> > Is your proposal to add a special exception wrapper interface that would >> > intercept all FS exceptions and re-wrap them with pre-classification? >> > It makes sense to me to have per-cloud classifier, but I am not sure >> > where would this interface live and when would we re-wrap the >> exceptions. >> > Would it be some FS wrapper like ExceptionClassifingFileSystem? If so, >> > I would suggest having it as a followup/child FLIP. >> > Could you please suggest how this classification would integrate with >> FS? >> > >> > 6. Migration >> > >> > I assume deprecation will be >> > > evaluated separately for each cloud plugin >> > >> > >> > It makes sense to consider deprecation timelines "per plugin". This >> > FLIP focuses on hadoop-less flink filesystems, and we are on track to >> > provide native versions for all main clouds. Hence I called it out as >> our >> > general intention. There are more Hadoop uses in the Flink codebase but >> > we are not aiming to address all of them in one go. >> > >> > We should also document that in-progress RecoverableWriter >> > > State may not be portable between Hadoop and native implementations >> > >> > >> > This is an important callout, it might block migration of jobs with >> > object storage sinks. Do you have a detailed writeup on the challenges >> > there that I could use as a reference? >> > >> > Thank you again for taking time to review and analyse the proposal. >> > I look forward to converging this discussion to unblock further work for >> > GCS and Azure support. >> > >> > Kind regards, >> > Aleksandr Iushmanov >> > >> > References: >> > [1] FLIP-597: Hadoop-less Flink Filesystems — >> > https://cwiki.apache.org/confluence/x/9gDuGQ >> > [2] FLIP-601: Client-Side Encryption Extension — >> > https://cwiki.apache.org/confluence/x/P4E_Gg >> > [3] >> > >> > >> https://cwiki.apache.org/confluence/spaces/FLINK/pages/435028214/FLIP-597+Filesystems+Hadoop-less+Flink+filesystems#FLIP597%3A%5BFilesystems%5DHadooplessFlinkfilesystems-Separateflink-fs-commonmodule >> > [4] S3AccessHelper: >> > >> > >> https://github.com/apache/flink/blob/master/flink-filesystems/flink-s3-fs-base/src/main/java/org/apache/flink/fs/s3/common/writer/S3AccessHelper.java >> > >> > >> > On Sun, 12 Jul 2026 at 06:44, Samrat Deb <[email protected]> wrote: >> > >> > > Hi Aleksandr, >> > > >> > > I have been applying the proposed interfaces to the native GCS >> > prototype. >> > > I >> > > encountered a few questions that I would like to understand better. >> > Some >> > > may >> > > come from gaps in my prototype or my understanding of the intended >> > > layering, >> > > so I would appreciate corrections and guidance. >> > > >> > > 1. Scope and interface layering >> > > >> > > Would it be helpful for the FLIP to show three explicit layers? >> > > >> > > - Existing Flink contracts, such as FileSystem and RecoverableWriter >> > > - Shared object-storage mechanics such as position tracking and >> ranged >> > > reopening >> > > - Provider-specific SDK adapters and capabilities >> > > >> > > My current impression is that ReadContext, StreamContext, >> > > InputStreamOpener, >> > > InputStreamExtension, and RawAndWrappedInputStreams overlap >> somewhat. >> > > Could we >> > > initially reduce this to a position-aware stream factory or handle, >> and >> > > add >> > > further context only when a second use case requires it? >> > > >> > > I would also like to understand the reason for placing the >> > > implementations in >> > > flink-core rather than a small flink-object-storage-fs-base module. >> > There >> > > may >> > > be plugin classloading constraints that I have missed. >> > > >> > > 2. Thread safety and cancellation >> > > >> > > In the GCS prototype, locking every read affected the throughput. >> There >> > > is also a >> > > Cancellation concern: when a lock is held during a blocking SDK >> read, a >> > > Concurrent close cannot acquire the lock to abort that operation. >> > > >> > > Would a narrower contract be sufficient: reads and writes remain >> > thread- >> > > Is confined, while close or abort is safe to invoke concurrently >> during >> > > task >> > > cancellation? If full thread safety is required, it would help to >> > > document the >> > > Flink caller that requires it. >> > > >> > > I will share the GCS benchmark results and a blocked-read >> cancellation >> > > test so >> > > that we can evaluate this with data rather than assumptions. >> > > >> > > 3. Stream ownership, CRT, and encryption >> > > >> > > I understand the need to support normal SDK streams, AWS CRT, >> > buffering, >> > > and >> > > future client-side encryption through one model. >> > > >> > > The remaining question for me is lifecycle ownership. S3 may require >> > > aborting >> > > a response rather than closing and draining it, while GCS and Azure >> > > expose >> > > different channel semantics. RawAndWrappedInputStreams may therefore >> > need >> > > more >> > > explicit rules for ownership, close ordering, and abort behaviour. >> > > >> > > Could the opener return an owned handle containing the readable >> stream >> > > and its >> > > abort/close operation? This might let each provider adapt its native >> > > stream >> > > without introducing an intermediate buffer. I will test this >> approach >> > > with >> > > both the GCS prototype and the native S3 implementation. >> > > >> > > For client-side encryption, I am still gathering requirements from >> the >> > > security discussion. Random access will require more than metadata >> > > injection >> > > because logical offsets, encrypted chunk boundaries, integrity tags, >> > and >> > > physical object lengths must be considered. I will share the >> > > header-injection >> > > and range-read requirements once that discussion concludes. >> > > >> > > 4. Existing open questions >> > > >> > > For RecoverableWriter, my current understanding is that >> > RecoverableWriter >> > > itself is already the cloud-neutral abstraction, while the recovery >> > > protocol >> > > should remain provider-specific. S3 uses multipart state, whereas >> the >> > > existing >> > > GCS writer uses immutable component objects and composes so that >> older >> > > checkpoint handles remain valid after writing continues. Did you >> > envisage >> > > a >> > > new common implementation here, or only shared utilities and >> contract >> > > tests? >> > > >> > > For PathsCopyingFileSystem, I believe the existing interface already >> > > provides >> > > the common capability. Each provider could implement it >> independently >> > > when an >> > > optimised transfer path is available. >> > > >> > > For ObjectStorageOperations, I understand the motivation, especially >> > > given the >> > > declining availability of reliable emulators. One possible >> incremental >> > > approach would be to use provider-local SDK adapters initially, and >> > share >> > > FileSystem behaviour tests, and extract a common operations >> interface >> > > once S3, >> > > GCS and Azure demonstrate a stable intersection. If the intended >> > > interface is >> > > narrower than a complete object-store abstraction, documenting the >> > exact >> > > operations and semantics would help me evaluate it against GCS. >> > > >> > > 5. Exception classification >> > > >> > > I am interested in connecting this work with the >> > exception-classification >> > > mechanism. It would be useful to align common categories such as >> > > authentication, not-found, precondition failure, throttling, >> timeout, >> > and >> > > cancellation while retaining the original SDK exception as the >> cause. >> > > >> > > Should this classification be part of FLIP-597, or should the >> > filesystem >> > > preserve provider exceptions and let the classification plugin >> handle >> > > them in >> > > a follow-up? I would be happy to prepare a comparison of the AWS, >> GCS, >> > > and >> > > Azure exception models. >> > > >> > > One final clarification concerns migration. I assume deprecation >> will >> > be >> > > evaluated separately for each cloud plugin rather than applying to >> > > flink-fs- >> > > hadoop-shaded generally, because HDFS and the remaining Hadoop-based >> > > integrations >> > > still require it. We should also document that in-progress >> > > RecoverableWriter >> > > State may not be portable between Hadoop and native implementations. >> > > >> > > I will follow up with the GCS prototype, performance notes, the >> S3/GCS >> > > behaviour comparison, and the encryption requirements. I am also >> happy >> > to >> > > help >> > > refine the interface documentation or add shared contract tests. >> > > >> > > These questions are intended to help us arrive at the smallest >> > > abstraction >> > > that preserves correctness and performance across providers. I >> support >> > > the >> > > Hadoop-less direction, and I am keen to learn from the Azure and S3 >> > > perspectives >> > > so that we can help finalise FLIP-597. >> > > >> > > Bests, >> > > Samrat >> > > >> > > On Thu, Jun 25, 2026 at 4:39 PM Aleksandr Iushmanov < >> [email protected] >> > > >> > > wrote: >> > > >> > > > Hi everyone, >> > > > >> > > > I'd like to start a discussion on FLIP-597: Hadoop-less Flink >> > Filesystems >> > > > [1]. >> > > > >> > > > As FLIP-555 [2] established, Flink's Hadoop-based filesystem plugins >> > > > carry significant maintenance burdens: transitive dependency >> conflicts, >> > > > classpath issues, and CVE exposure. FLIP-555 addressed this for S3 >> by >> > > > introducing a native SDK-based filesystem. This FLIP extends that >> > effort >> > > > as an umbrella for all cloud providers (Azure, GCS, and potentially >> > > > others) with the end goal to deprecate Hadoop-based variants. >> > > > >> > > > The FLIP proposes common, cloud-agnostic I/O abstractions placed in >> > > > flink-core that all SDK-native filesystem implementations share as >> > > > a step towards Hadoop-less Flink deployments: >> > > > >> > > > - ObjectStorageInputStream / ObjectStorageOutputStream: thread-safe >> > > > stream implementations. >> > > > - InputStreamExtension / InputStreamOpener: composable read pipeline >> > > > supporting buffering, decryption via decoration >> > > > - WriteContext / ReadContext: metadata descriptors enabling features >> > > > like client-side encryption without changing the FileSystem API >> > > > >> > > > There are a number of open questions for discussion: >> > > > 1. Should a cloud-agnostic RecoverableWriter abstraction be part of >> > > > this FLIP given that implementation across clouds may be too >> > specific >> > > > for common abstractions to be useful. >> > > > 2. Should PathsCopyingFileSystem (bulk copy) be in scope? Transfer >> > > managers >> > > > for clouds vary. And currently we only support this >> implementation >> > in >> > > > AWS. >> > > > 3. Should this FLIP include an ObjectStorageOperations interface >> > > > (analogous to FLIP-555's S3AccessHelper) for a generalized >> > > > ObjectStorageFileSystem base, or should we defer premature >> > > > generalisation >> > > > of testing setup? Main motivation for this would be the ability >> to >> > > > verify >> > > > common FS parts with testing doubles with reduced reliance on >> other >> > > > libraries >> > > > like MinIO [3] (archived), LocalStack [4] (archived), Azurite [5] >> > (no >> > > > SDK V2 support) >> > > > that demonstrates a decline in support. >> > > > >> > > > Looking forward to your feedback. >> > > > >> > > > Best regards, >> > > > Aleksandr Iushmanov >> > > > >> > > > [1] https://cwiki.apache.org/confluence/x/9gDuGQ >> > > > [2] https://cwiki.apache.org/confluence/x/uYqmFw >> > > > [3] https://github.com/minio/minio >> > > > [4] https://github.com/localstack/localstack >> > > > [5] https://github.com/Azure/Azurite >> > > > >> > > >> > >> >
