Skip to content

HDDS-16396. Support concurrent positioned reads for EC and encrypted keys - #11348

Open
peterxcli wants to merge 21 commits into
apache:masterfrom
peterxcli:HDDS-16396-positioned-read-api
Open

peterxcli wants to merge 21 commits into
apache:masterfrom
peterxcli:HDDS-16396-positioned-read-api

Conversation

@peterxcli

@peterxcli peterxcli commented Sep 28, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

Hadoop positioned reads now use independent readers across Ozone's client and filesystem streams, including EC and encrypted keys. Wrappers delegate native reads and report only capabilities supported by their delegates. This builds on the merged #11314.

OzoneInputStream now extends ExtendedInputStream, and consequently Hadoop's FSInputStream, to expose the positioned-read API.

some context:

What is the link to the Apache JIRA

https://issues.apache.org/jira/browse/HDDS-16396

How was this patch tested?

Added and updated unit and integration tests covering concurrent EC and encrypted positioned reads, multipart boundaries, cursor preservation, EOF handling, capability reporting, and client cleanup.

Generated-by: Codex (GPT-6)

Copilot AI lite review requested due to automatic review settings September 28, 2026 12:48

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Resolve the positioned-read client leak and correct capability delegation in OzoneInputStream.

Review effort: Lite
Findings: 2 Medium severity

Open (2)
What changed in this PR

This pull request standardizes positioned-read APIs across Ozone client, filesystem, EC, streaming, multipart, and encrypted streams.

Changes:

  • Adds positioned-read delegation and capability reporting.
  • Implements dedicated EC and streaming readers.
  • Expands concurrency, failover, encryption, and wrapper test coverage.
File Reviewed change
hadoop-ozone/​ozonefs-common/​src/​test/​java/​org/​apache/​hadoop/​fs/​ozone/​TestOzoneFSInputStream.java Tests wrapper delegation and unsupported streams.
hadoop-ozone/​ozonefs-common/​src/​main/​java/​org/​apache/​hadoop/​fs/​ozone/​OzoneFSInputStream.java Delegates positioned reads.
hadoop-ozone/​ozonefs-common/​src/​main/​java/​org/​apache/​hadoop/​fs/​ozone/​CapableOzoneFSInputStream.java Reports read capabilities.
hadoop-ozone/​integration-test/​src/​test/​java/​org/​apache/​hadoop/​ozone/​client/​rpc/​TestOzoneAtRestEncryption.java Tests encrypted multipart reads.
hadoop-ozone/​integration-test/​src/​test/​java/​org/​apache/​hadoop/​fs/​ozone/​TestOzoneFSInputStream.java Adds EC and streaming coverage.
hadoop-ozone/​integration-test/​pom.xml Adds the client test-JAR dependency.
hadoop-ozone/​client/​src/​test/​java/​org/​apache/​hadoop/​ozone/​client/​io/​TestOzoneCryptoInputStream.java Tests concurrent encrypted reads.
hadoop-ozone/​client/​src/​main/​java/​org/​apache/​hadoop/​ozone/​client/​io/​OzoneInputStream.java Adds positioned-read delegation.
hadoop-ozone/​client/​src/​main/​java/​org/​apache/​hadoop/​ozone/​client/​io/​OzoneCryptoInputStream.java Supports concurrent positioned decryption.
hadoop-hdds/​client/​src/​test/​java/​org/​apache/​hadoop/​ozone/​client/​io/​TestECBlockReconstructedInputStream.java Tests concurrent EC reconstruction.
hadoop-hdds/​client/​src/​test/​java/​org/​apache/​hadoop/​ozone/​client/​io/​TestECBlockInputStreamProxy.java Tests EC overlap, failover, and cleanup.
hadoop-hdds/​client/​src/​test/​java/​org/​apache/​hadoop/​ozone/​client/​io/​TestECBlockInputStream.java Adds positioned-read tests.
hadoop-hdds/​client/​src/​test/​java/​org/​apache/​hadoop/​ozone/​client/​io/​ECStreamTestUtil.java Adds positioned-read test utilities.
hadoop-hdds/​client/​src/​test/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​TestStreamBlockInputStream.java Tests streaming positioned reads and cleanup.
hadoop-hdds/​client/​src/​test/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​TestMultipartInputStream.java Tests multipart positioned reads.
hadoop-hdds/​client/​src/​test/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​PositionedReadTestHelper.java Improves concurrent offset distribution.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​ozone/​client/​io/​ECBlockReconstructedInputStream.java Rejects unsupported positioned reads.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​ozone/​client/​io/​ECBlockInputStreamProxy.java Creates independent EC readers.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​ozone/​client/​io/​ECBlockInputStream.java Marks direct positioned reads unsupported.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​StreamBlockInputStream.java Implements streaming positioned reads.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​PartInputStream.java Extends the positioned-read contract.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​MultipartInputStream.java Routes reads across parts.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​ExtendedInputStream.java Defines the positioned-read API.
hadoop-hdds/​client/​src/​main/​java/​org/​apache/​hadoop/​hdds/​scm/​storage/​BlockInputStream.java Exposes block positioned reads.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

@peterxcli
peterxcli marked this pull request as ready for review October 1, 2026 05:04
@peterxcli
peterxcli force-pushed the HDDS-16396-positioned-read-api branch from 8f3291b to 0a5b8af Compare October 1, 2026 05:30
Comment on lines +63 to +64
int offset = (int) (((long) threadId * source.length / THREAD_COUNT + i * 17L)
% (source.length - BUFFER_SIZE));

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The old formula read only the first 13 KiB of a 512 KiB source, and the new formula reads more chunks and parts.

OLD · 512 KiB source
       T0 T1 T2 T3 T4 T5 T6 T7
        \  \  \  \  /  /  /  /
                 ↓
       +---------+-------------------------------------------------------------------------------+
       | READS   |                                                                               |
       +---------+-------------------------------------------------------------------------------+
       0      12,779 B                                                                       512 KiB


NEW · 512 KiB source

          T0          T1          T2          T3          T4          T5          T6          T7
          ↓           ↓           ↓           ↓           ↓           ↓           ↓           ↓
       +-----------------------------------------------------------------------------------------------+
       |  [===]       [===]       [===]       [===]       [===]       [===]       [===]       [===]    |
       +-----------------------------------------------------------------------------------------------+
          0          64          128         192         256         320         384         448      512 KiB


ONE THREAD · 100 reads · 4,096 bytes/read

                 start offset                                         end offset
                      │                                                    │
Read 0           +0 B ├───────────────── 4,096 B ──────────────────────────┤ +4,096 B
Read 1          +17 B  ├───────────────── 4,096 B ──────────────────────────┤ +4,113 B
Read 2          +34 B   ├───────────────── 4,096 B ──────────────────────────┤ +4,130 B
                       ⋮
Read 99      +1,683 B           ├───────────────── 4,096 B ──────────────────────────┤ +5,779 B
                      │                                                              │
                      └──────────────── [===] = 5,779 B ─────────────────────────────┘

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks for improving this test helper

@peterxcli
peterxcli requested review from smengcl and szetszwo October 1, 2026 06:39
@peterxcli

Copy link
Copy Markdown
Member Author

cc @yandrey321 @taklwu ptal!

@peterxcli peterxcli changed the title HDDS-16396. Standardize positioned-read APIs across client streams HDDS-16396. Support concurrent positioned reads for EC and encrypted keys Oct 1, 2026
@peterxcli
peterxcli requested a review from ivandika3 October 2, 2026 03:58

@taklwu taklwu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment on lines +63 to +64
int offset = (int) (((long) threadId * source.length / THREAD_COUNT + i * 17L)
% (source.length - BUFFER_SIZE));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks for improving this test helper

*/
public interface PartInputStream
extends CanUnbuffer, Seekable {
extends CanUnbuffer, Seekable, ByteBufferPositionedReadable {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

with ByteBufferPositionedReadable, that looks better.

assertInstanceOf(MultipartInputStream.class, inputStream.getInputStream());
if (numParts == 2) {
inputStream.seek(123);
PositionedReadTestHelper.runConcurrentPositionedReads(inputData, inputStream::readFully);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we also use this helper in testPutKeyWithEncryption to cover concurrent positioned reads for a regular encrypted key? That path uses Hadoop's CryptoInputStream directly, so it would be useful to verify the data and unchanged sequential position through that wrapper as well.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ditto. Could we extend testPutKeyWithEncryption with concurrent byte-array and
ByteBuffer positioned reads, checking the returned plaintext and unchanged sequential
position? Regular encrypted keys use Hadoop's CryptoInputStream directly, which the
multipart coverage doesn't exercise.

Please use enough data to cross a crypto-buffer boundary and reuse the existing
cluster/helper.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sure, done in a4ebc7f

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: We could mix byte-array reads into this helper and add an explicit read across a crypto-buffer boundary. This currently calls only the ByteBuffer overload, and none of the helper's generated ranges crosses the default 8 KiB boundary.

--- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestOzoneAtRestEncryption.java
+++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestOzoneAtRestEncryption.java
@@ -317,7 +317,19 @@
     try (OzoneInputStream inputStream = bucket.readKey(keyName)) {
       assertInstanceOf(CryptoInputStream.class, inputStream.getInputStream());
       inputStream.seek(123);
-      PositionedReadTestHelper.runConcurrentPositionedReads(data, inputStream::readFully);
+      PositionedReadTestHelper.runConcurrentPositionedReads(data, (offset, buffer) -> {
+        if ((offset & 1) == 0) {
+          inputStream.readFully(offset, buffer);
+        } else {
+          byte[] bytes = new byte[buffer.remaining()];
+          inputStream.readFully(offset, bytes);
+          buffer.put(bytes);
+        }
+      });
+      int offset = DEFAULT_CRYPTO_BUFFER_SIZE - 17;
+      byte[] actual = new byte[DEFAULT_CRYPTO_BUFFER_SIZE + 37];
+      inputStream.readFully(offset, actual);
+      assertArrayEquals(Arrays.copyOfRange(data, offset, offset + actual.length), actual);
       assertEquals(123, inputStream.getPos());
       assertEquals(Byte.toUnsignedInt(data[123]), inputStream.read());
     }

@rich7420

rich7420 commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

Overall LGTM, with the non-blocking suggestions

@peterxcli
peterxcli requested a review from taklwu October 5, 2026 17:20

@yandrey321 yandrey321 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Coverage is genuinely good on the paths it covers — the EC test asserting
opened/closed reader counts, overlap via CyclicBarrier, failover and injected failure
is the right shape, and so is the encrypted-multipart overlap test. Four gaps:

  1. No GDPR coverage anywhere. That is where both blockers live. A test that opens a
    gdprEnabled bucket's key through the filesystem and asserts
    hasCapability(PREADBYTEBUFFER) matches what read(long, ByteBuffer) and
    read(long, byte[], int, int) actually do would have caught #1 and #2. Extending
    TestOzoneAtRestEncryption or the GDPR suite is preferable to a new cluster.

  2. TestMultipartInputStream.testConcurrentPositionedRead uses
    @CsvSource({"false, false", "true, false", "false, true"}) — the all-streaming cell
    is missing. It is covered indirectly by testStreamBlockConcurrentPositionedRead, so
    either add "true, true" or say so in a comment.

  3. Nothing exercises readPositioned concurrently with initialize(), which is where
    finding #3(a) lives, nor a part whose length changes after construction.

  4. positionedReadsOverlapAcrossEncryptedParts depends on both threads issuing exactly
    two underlying reads so the CyclicBarrier(2) trips evenly. It does with the current
    single-shot stub, but a stub that returned partial reads (which is what a real
    KeyInputStream does) would make the counts diverge and the barrier time out. Worth
    a comment on the stub, or a barrier that tolerates uneven arrival.

case StreamCapabilities.UNBUFFER:
case StreamCapabilities.PREADBYTEBUFFER:
return true;
case StreamCapabilities.PREADBYTEBUFFER:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This check is a tautology in production. The wrapped stream here is always an
OzoneInputStream (built in RpcClient.createInputStream and handed to
CapableOzoneFSInputStream by OzoneFileSystem:114 / RootedOzoneFileSystem:112), and
this PR makes OzoneInputStream extend ExtendedInputStream, which implements
ByteBufferPositionedReadable unconditionally. So the answer is true for every key,
including a GDPR key, where RpcClient:2682 wraps a javax.crypto.CipherInputStream and
OzoneInputStream.readPositioned:83 throws UnsupportedOperationException.

A caller that does the right thing — probe hasCapability(PREADBYTEBUFFER), then call
read(long, ByteBuffer) — gets told yes and then gets a RuntimeException.

OzoneInputStream.hasCapability already computes this correctly (it tests its own
delegate). Delegate to it instead of re-deriving from the type:

case StreamCapabilities.PREADBYTEBUFFER:
  return getWrappedInputStream() instanceof StreamCapabilities
      && ((StreamCapabilities) getWrappedInputStream()).hasCapability(capability);

Note the new unit test does not catch this because it bypasses the production
wrapping: cursorOnlyStreamsDoNotEmulatePositionedReads constructs
new CapableOzoneFSInputStream(new SeekableOnlyInputStream(source), null) — a raw
delegate rather than an OzoneInputStream — so the instanceof legitimately answers
false there. Please add a case that wraps the delegate in OzoneInputStream, which is
the only shape the filesystem ever sees.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Both filesystem adapters call bucket.readFile(key).getInputStream(), which unwraps OzoneInputStream. For GDPR keys, the delegate is the raw CipherInputStream, so this check returns false.

return inputStream;
}

private int readImpl(long position, ByteBuffer buffer) throws IOException {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Both byte-array overloads now route here — read(long, byte[], int, int):183 and
readFully(long, byte[], int, int):228 — so they throw for any delegate that cannot do
ByteBuffer positioned reads. That is a functional regression, not just a capability
question: read(long, byte[], int, int) and readFully are plain PositionedReadable
methods. Hadoop callers invoke them with no capability probe at all (FSDataInputStream
exposes them directly; ORC, Parquet and Hadoop's own IOUtils.readFully use them), and
there is no StreamCapabilities constant that gates them.

Before this patch, bestEffortRead/bestEffortReadFullySynchronized served these under
positionedReadLock via seek-read-restore. The javadoc you deleted named exactly the
case that now breaks:

"Uses read(ByteBuffer) -- which handles streams that do not implement
 ByteBufferReadable (e.g. GDPR javax.crypto.CipherInputStream)"

Reaching concurrency by deleting the only path that worked for a whole class of keys
is the wrong trade. Options, in order of preference:

a) Keep the locked seek-read-restore fallback for the byte-array overloads only,
reached when the delegate is not ByteBufferPositionedReadable. Concurrency is
unchanged for every stream that does support native pread; GDPR keys keep
working, serialized, as they do today.
b) Make the GDPR path capable at the source — RpcClient:2682 could wrap the
KeyInputStream in something that implements ByteBufferPositionedReadable with
per-call Cipher instances — which is a larger change and probably its own Jira.

Either way the PR description should say that GDPR-encrypted keys are affected; right
now it reads as purely additive.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The old fallback casts to Seekable, which CipherInputStream doesn't implement, so GDPR positioned reads already failed. Removing support for other seek-only delegates is intentional and covered by a test.

}
throw new EOFException("EOF encountered at pos: " + pos +
" for key: " + key);
if (n == 0) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

binarySearchOffsetIndex relies on Arrays.binarySearch, whose result among duplicate
offsets is unspecified. A zero-length part produces a duplicate offset, so the search
can land on it; the clamp at line 197 then sets the limit to the current position,
part.read() correctly returns 0 for an empty buffer, and this throws — for a stream
that is in no way broken. The method created the empty buffer itself, so 0 here is not
a protocol violation.

continue (after advancing past the part) rather than throwing would be correct. Same
available <= 0 guard suggested above covers it.

I could not construct a zero-length part from today's write paths, so this is
robustness rather than a live bug — but the throw is reachable from a state the method
itself produces, which is worth closing cheaply.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The write paths omit zero-length blocks when committing to OM, so we don't need an empty-part guard here. I've removed the buffer clamp and left bounds handling to each part.

}

@Override
protected int readPositioned(long position, ByteBuffer buffer) throws IOException {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The no-arg constructor at :40 leaves inputStream null, so this NPEs on
getClass() instead of throwing the intended UnsupportedOperationException. A
null-guard in the message (or in the constructor) is a one-liner.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

RpcClient always supplies a stream, and I found no production callers of the no-arg constructor. I'd leave null handling out of this PR.

}

@Override
protected int readWithStrategy(ByteReaderStrategy strategy) throws IOException {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is dead code: read(), read(byte[],int,int) and read(ByteBuffer) at :53-:71 all
delegate straight to inputStream and never reach the base class's
read(ByteReaderStrategy). Two consequences worth deciding on deliberately:

  • the base class's sequential entry points are synchronized, these overrides are
    not, so OzoneInputStream's sequential reads have different locking from every
    other ExtendedInputStream;
  • if the strategy path is genuinely unreachable, implementing it invites a future
    reader to assume it is live.

Either drop the direct overrides and let the strategy path do the work (uniform
locking, one code path), or keep them and make readWithStrategy throw
UnsupportedOperationException with a comment saying why.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The inherited public read(ByteReaderStrategy) calls this method, so it is reachable. Making it throw would break that API.

@peterxcli
peterxcli requested a review from yandrey321 October 6, 2026 05:07

@taklwu taklwu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

please fix those tests failure if it's not flaky, otherwise +1

@smengcl smengcl left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

looks good. Thanks @peterxcli . Just one nit above

@yandrey321 yandrey321 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@rich7420 rich7420 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM!

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants