HDDS-16258. Refactor streaming block reads and fix checksum verification for variable-sized chunks - #11302
HDDS-16258. Refactor streaming block reads and fix checksum verification for variable-sized chunks#11302peterxcli wants to merge 11 commits into
Conversation
There was a problem hiding this comment.
Warning
Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.
Copilot review overview
Review effort: Lite
Findings: 1
Open (3)
This treats any short read as EOF and fails the request, but file/channel reads are not guaranteed… · New Buffer sizing can grow to the maximum checksum interval seen in chunk metadata, which may cause… · New AresponseDataSizeof 0 isn't rejected here, butBlockReadCursorrejectsresponseSize <= 0by… · New
What changed in this PR
Updates streaming ReadBlock to correctly validate checksums across variable-length chunks by sending per-response overlapping ChunkInfo metadata and using a BlockReadCursor to align reads to chunk-relative checksum boundaries.
Changes:
- Replace flattened checksum aggregation with
BlockReadCursor+ per-responsechunkInfoListchecksum validation. - Update wire protocol: reserve
checksumDatafield and addchunkInfoListtoReadBlockResponseProto. - Expand unit/integration tests to cover uneven flushed chunks, unexpected EOF, and chunk-relative checksum verification.
| File | Description |
|---|---|
| hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java | Adds integration coverage for uneven flushed chunks and range/seek verification with streaming reads. |
| hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto | Reworks ReadBlock response payload from flat checksum data to per-response chunk metadata. |
| hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java | Updates tests for new handler signature and validates returned chunkInfoList + checksum verification behavior. |
| hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestBlockReadCursor.java | New unit tests for cursor range selection and boundary alignment. |
| hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/ReadBlockComputation.java | Removes legacy checksum-aligned range computation helper. |
| hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java | Switches streaming readBlock implementation to BlockReadCursor and per-response chunk checksum validation. |
| hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/BlockReadCursor.java | New cursor for chunk-relative checksum boundary alignment and response chunk selection. |
| hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java | Updates handler interface to remove testVariableChunks plumbing. |
| hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java | Adjusts dispatcher call site for new readBlock signature. |
| hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java | Adds tests for chunk-relative validation, malformed coverage/alignment, and NONE handling. |
| hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java | Reworks validation to operate on chunk-relative boundaries and detect missing/invalid coverage. |
| hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java | Builds ReadBlock responses using chunkInfoList instead of flat checksum data. |
| hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java | Updates tests for new response format and adds coverage for seek-prefix trimming and missing metadata. |
| hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java | Verifies each streaming response using chunkInfoList instead of checksumData. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
|
|
||
| /** Tracks chunk-relative checksum boundaries for a streaming block read. */ | ||
| class BlockReadCursor { | ||
| private static final int STREAMING_BYTES_PER_CHUNK = 1024 * 64; |
There was a problem hiding this comment.
should it be configurable ?
There was a problem hiding this comment.
That value merely preserves the old behavior for chunks written with NONE and is introduced in https://github.com/apache/ozone/pull/9369/changes#diff-04096cfd662f347dcf21e9e13e714af882dab2dfc02e4e9d71307b0c62b35a24R183
| reserved 1; | ||
| reserved "checksumData"; |
There was a problem hiding this comment.
I actually want to remove and renumber this as the streaming read is still off by default(means experimental), we could probably ignore the backward compatibility.
but that will make the check against proto.lock fail
There was a problem hiding this comment.
With this change, an older client with streaming reads enabled would fail to parse responses after the datanodes are upgraded, even following a non-rolling cluster upgrade. The reverse, an upgraded client reading from older datanodes also fails when checksum verification is enabled because chunkInfoList is absent.
At least this is just a wire proto change and does not affect stored data. Users can disable streaming reads and reopen the client, or by bringing clients and datanodes onto compatible versions. Given that streaming reads are opt-in and this break is intentional for performance reason IIUC, I’m fine with the change for now, but could we document both failure scenarios and recovery steps in the release notes?
There was a problem hiding this comment.
could we document both failure scenarios and recovery steps in the release notes?
yes, that's what I want.
chungen0126
left a comment
There was a problem hiding this comment.
Thanks @peterxcli for working on this.
Two suggestions:
- Add a flag in the read request so the server skips sending chunkInfoList when verification isn't needed, saving payload size.
- The server-side refactoring is out of scope here; please move it to a separate ticket.
| if (readBlock.getOffset() < 0 || readBlock.getLength() < 0 | ||
| || responseDataSize < 0 || responseDataSize > OZONE_SCM_CHUNK_MAX_SIZE) { | ||
| return rejectReadBlock(blockFile, streamObserver, Status.INVALID_ARGUMENT.withDescription( | ||
| "Invalid ReadBlock range or response size: " + readBlock)); | ||
| } |
There was a problem hiding this comment.
Is this validation needed on the server side? The client should already be preventing these invalid requests from being sent.
There was a problem hiding this comment.
I think you're right here. I'll admit that these validations were added by AI, and I'm ok to remove them because they're apparently not in the scope of either the server refactor or
readBlock.getOffset() < 0 || readBlock.getLength() < 0 || responseDataSize < 0: though our client won't send this, it can still protect the cluster from malformed grpc requests from arbitrary custom clients.responseDataSize > OZONE_SCM_CHUNK_MAX_SIZE: this is possible with the java client, and without this check, the client will get a fatal exception. See:
| @@ -441,37 +439,39 @@ public static void verifyChecksum(List<ByteBuffer> bufferList, int startIndex, C | |||
| public static void validateChecksums(ByteBuffer data, long blockOffset, int startIndex, | |||
There was a problem hiding this comment.
Adding those check is nice, but could you clarify the motivation behind refactoring this method? From my perspective, this change seems unnecessary. Please leave this part as it is.
| ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData()); | ||
| Checksum.verifyChecksum(data, checksumData, 0); | ||
| Checksum.validateChecksums(data, readBlock.getOffset(), 0, readBlock.getChunkInfoListList()); |
There was a problem hiding this comment.
With the current approach, if a new client interacts with an older server, reading data will completely fail because chunkInfoList will be empty, which triggers an OzoneChecksumException. Even if checksum verification cannot be performed in that case, I think we should still ensure the client can fall back and read the data normally.
There was a problem hiding this comment.
Please see my earlier comment: #11302 (comment).
Since our streaming read is off by default, I'd consider it an experimental feature. Blowing up the codebase with compatibility handling for an experimental feature makes no sense to me (and even if we did preserve compatibility, new clients would still fail in the same situation).
A few things we can do easily: 1. Add a caution to the release notes or user docs, mentioning something like: "If you want to use streaming read, don't mix datanodes running version 2.2.x with clients on 2.3+."
cc @taklwu @sodonnel @amaliujia as you might have more context with real prod usage of streaming read.
peterxcli
left a comment
There was a problem hiding this comment.
The server-side refactoring is out of scope here; please move it to a separate ticket.
Sure, let me open and merge the server refactor as a separate PR first, then we'll get back to here.
Add a flag in the read request so the server skips sending chunkInfoList when verification isn't needed, saving payload size.
let me look into this, thanks for the suggestion
| ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData()); | ||
| Checksum.verifyChecksum(data, checksumData, 0); | ||
| Checksum.validateChecksums(data, readBlock.getOffset(), 0, readBlock.getChunkInfoListList()); |
There was a problem hiding this comment.
Please see my earlier comment: #11302 (comment).
Since our streaming read is off by default, I'd consider it an experimental feature. Blowing up the codebase with compatibility handling for an experimental feature makes no sense to me (and even if we did preserve compatibility, new clients would still fail in the same situation).
A few things we can do easily: 1. Add a caution to the release notes or user docs, mentioning something like: "If you want to use streaming read, don't mix datanodes running version 2.2.x with clients on 2.3+."
cc @taklwu @sodonnel @amaliujia as you might have more context with real prod usage of streaming read.
| if (readBlock.getOffset() < 0 || readBlock.getLength() < 0 | ||
| || responseDataSize < 0 || responseDataSize > OZONE_SCM_CHUNK_MAX_SIZE) { | ||
| return rejectReadBlock(blockFile, streamObserver, Status.INVALID_ARGUMENT.withDescription( | ||
| "Invalid ReadBlock range or response size: " + readBlock)); | ||
| } |
There was a problem hiding this comment.
I think you're right here. I'll admit that these validations were added by AI, and I'm ok to remove them because they're apparently not in the scope of either the server refactor or
readBlock.getOffset() < 0 || readBlock.getLength() < 0 || responseDataSize < 0: though our client won't send this, it can still protect the cluster from malformed grpc requests from arbitrary custom clients.responseDataSize > OZONE_SCM_CHUNK_MAX_SIZE: this is possible with the java client, and without this check, the client will get a fatal exception. See:


What changes were proposed in this pull request?
Streaming ReadBlock collects checksums assuming uniform chunks, but flushes can produce short chunks whose checksum intervals restart at different block offsets.
This patch:
KeyValueHandler#readBlockImplintoBlockReadCursor. The handler consumes each range from the cursor, reads and sends the data, then advances the cursor. It no longer manages the boundary calculations itself.The client sets
includeChecksumsfrom its existing verification setting. When false, the server omitschunkInfoListwithout changing response ranges or its own checksum verification. The request field defaults to true when omitted.The production path now handles variable chunks without
testVariableChunksor the legacy flattened checksum collector. Shared verification rejects missing coverage and invalid boundaries, handlesNONE, and preserves buffer positions. Unexpected EOF fails before sending partial checksum data; OUT_OF_RANGE retains stream error delivery and file cleanup.Protocol compatibility:
ReadBlockResponseProtoreplaces the flattened checksum field withchunkInfoList = 4, reserves field 1 andchecksumData, and preservesoffset = 2anddata = 3. The release compatibility baseline (proto.lock) is unchanged. This intentionally requires matching clients and datanodes for streaming reads, which are disabled by default. There is no legacy fallback. Ordinary chunk reads and the stored format are unchanged.What is the link to the Apache JIRA
HDDS-16258
How was this patch tested?
Validated on JDK 21 with regenerated protobufs, focused checksum/client/server/cursor tests, and checkstyle. Tests cover omitted metadata, unchanged response ranges, and independent datanode verification. Existing
TestStreamReadintegration coverage verifies uneven stored chunks through full and seek/range reads.