Repository navigation
Conversation
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Resolve the positioned-read client leak and correct capability delegation in OzoneInputStream.
Review effort: Lite
Findings: 2
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.
8f3291b to
0a5b8af
Compare
| int offset = (int) (((long) threadId * source.length / THREAD_COUNT + i * 17L) | ||
| % (source.length - BUFFER_SIZE)); |
There was a problem hiding this comment.
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 ─────────────────────────────┘
There was a problem hiding this comment.
thanks for improving this test helper
|
cc @yandrey321 @taklwu ptal! |
| int offset = (int) (((long) threadId * source.length / THREAD_COUNT + i * 17L) | ||
| % (source.length - BUFFER_SIZE)); |
There was a problem hiding this comment.
thanks for improving this test helper
| */ | ||
| public interface PartInputStream | ||
| extends CanUnbuffer, Seekable { | ||
| extends CanUnbuffer, Seekable, ByteBufferPositionedReadable { |
There was a problem hiding this comment.
with ByteBufferPositionedReadable, that looks better.
| assertInstanceOf(MultipartInputStream.class, inputStream.getInputStream()); | ||
| if (numParts == 2) { | ||
| inputStream.seek(123); | ||
| PositionedReadTestHelper.runConcurrentPositionedReads(inputData, inputStream::readFully); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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());
}|
Overall LGTM, with the non-blocking suggestions |
yandrey321
left a comment
There was a problem hiding this comment.
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:
-
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. -
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. -
Nothing exercises readPositioned concurrently with initialize(), which is where
finding #3(a) lives, nor a part whose length changes after construction. -
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: |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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) { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
The inherited public read(ByteReaderStrategy) calls this method, so it is reachable. Making it throw would break that API.
taklwu
left a comment
There was a problem hiding this comment.
please fix those tests failure if it's not flaky, otherwise +1
smengcl
left a comment
There was a problem hiding this comment.
looks good. Thanks @peterxcli . Just one nit above

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.
OzoneInputStreamnow extendsExtendedInputStream, and consequently Hadoop'sFSInputStream, 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)