Skip to content

HDDS-16258. Refactor streaming block reads and fix checksum verification for variable-sized chunks - #11302

Open
peterxcli wants to merge 11 commits into
apache:masterfrom
peterxcli:HDDS-16258
Open

peterxcli wants to merge 11 commits into
apache:masterfrom
peterxcli:HDDS-16258

Conversation

@peterxcli

@peterxcli peterxcli commented Sep 22, 2026 •

Copy link
Copy Markdown
Member

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:

  • Move chunk lookup, checksum alignment, response sizing, and offset tracking from KeyValueHandler#readBlockImpl into BlockReadCursor. 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.
  • Fix checksum verification across variable-sized chunks. Include the overlapping chunks’ original metadata and stored checksums in each response when requested, and verify the data using chunk-relative checksum boundaries on both the client and server.

The client sets includeChecksums from its existing verification setting. When false, the server omits chunkInfoList without changing response ranges or its own checksum verification. The request field defaults to true when omitted.

+--------------------------------------------------------------------------------------------------+
| SAME INPUT: full read [0,12), CRC32, bytesPerChecksum=4, response buffer >=12 bytes              |
|                                                                                                  |
| Chunk metadata: C1 offset=0, len=5; C2 offset=5, len=7                                           |
| Stored bytes:   [ A B C D E ] [ F G H I J K L ]                                                  |
| Stored CRCs:    C1: [ABCD] [E]    C2: [FGHI] [JKL]                                               |
| [ABCD] means the stored CRC of bytes ABCD; intervals restart in each chunk.                      |
+--------------------------------------------------------------------------------------------------+
                                                  |
                        +-------------------------+-------------------------+
                        v                                                   v
+----------------------------------------------+    +----------------------------------------------+
| BEFORE: SERVER / readBlockImpl               |    | AFTER: SERVER / readBlockImpl                |
| ReadBlockComputation + getChecksums()        |    | BlockReadCursor                              |
|                                              |    |                                              |
| Assumes every chunk is C1-sized (5 bytes).   |    | Selects [0,12) using chunk boundaries.       |
| Picks CRCs at block offsets 0, 4, 8:         |    | Includes the two overlapping chunks          |
| [ABCD] [E] [FGHI]  (drops [JKL])             |    | with their original offsets and CRCs.        |
+----------------------------------------------+    +----------------------------------------------+
                        |                                                   |
                        v                                                   v
+----------------------------------------------+    +----------------------------------------------+
| WIRE: ReadBlockResponseProto                 |    | WIRE: ReadBlockResponseProto                 |
| offset=0, data=ABCDEFGHIJKL                  |    | offset=0, data=ABCDEFGHIJKL                  |
| Flat checksumData:                           |    | chunkInfoList:                               |
|   [ABCD] [E] [FGHI]                          |    |   C1: offset=0, len=5, [ABCD] [E]            |
|                                              |    |   C2: offset=5, len=7, [FGHI] [JKL]          |
+----------------------------------------------+    +----------------------------------------------+
                        |                                                   |
                        v                                                   v
+----------------------------------------------+    +----------------------------------------------+
| BEFORE: CLIENT / StreamBlockInputStream      |    | AFTER: CLIENT / StreamBlockInputStream       |
| Splits the response into flat groups:        |    | Uses each chunk's checksum boundaries:       |
|   [ABCD] [EFGH] [IJKL]                       |    |   C1: [ABCD] [E]   C2: [FGHI] [JKL]          |
|                                              |    |                                              |
| CRC(EFGH) != stored CRC(E)                   |    | All stored CRCs match.                       |
| FAIL: response is not enqueued.              |    | PASS: enqueue ABCDEFGHIJKL.                  |
+----------------------------------------------+    +----------------------------------------------+

The production path now handles variable chunks without testVariableChunks or the legacy flattened checksum collector. Shared verification rejects missing coverage and invalid boundaries, handles NONE, and preserves buffer positions. Unexpected EOF fails before sending partial checksum data; OUT_OF_RANGE retains stream error delivery and file cleanup.

Protocol compatibility: ReadBlockResponseProto replaces the flattened checksum field with chunkInfoList = 4, reserves field 1 and checksumData, and preserves offset = 2 and data = 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 TestStreamRead integration coverage verifies uneven stored chunks through full and seek/range reads.

Copilot AI lite review requested due to automatic review settings September 22, 2026 17:15

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.

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 High severity · 2 Medium severity

Open (3)
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-response chunkInfoList checksum validation.
  • Update wire protocol: reserve checksumData field and add chunkInfoList to ReadBlockResponseProto.
  • 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.

@peterxcli
peterxcli requested a balanced review from Copilot September 22, 2026 18:51

This comment was marked as resolved.

@peterxcli
peterxcli marked this pull request as ready for review September 22, 2026 18:52

/** Tracks chunk-relative checksum boundaries for a streaming block read. */
class BlockReadCursor {
private static final int STREAMING_BYTES_PER_CHUNK = 1024 * 64;

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.

should it be configurable ?

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.

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

Comment on lines +401 to +402
reserved 1;
reserved "checksumData";

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.

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

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 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?

@peterxcli peterxcli Sep 25, 2026 •

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.

could we document both failure scenarios and recovery steps in the release notes?

yes, that's what I want.

@chungen0126 chungen0126 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.

Thanks @peterxcli for working on this.
Two suggestions:

  1. Add a flag in the read request so the server skips sending chunkInfoList when verification isn't needed, saving payload size.
  2. The server-side refactoring is out of scope here; please move it to a separate ticket.

Comment on lines +2349 to +2353
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));
}

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.

Is this validation needed on the server side? The client should already be preventing these invalid requests from being sent.

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.

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

  1. 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.
  2. 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:
    public static int limitReadSize(long len)
    throws StorageContainerException {
    if (len > OzoneConsts.OZONE_SCM_CHUNK_MAX_SIZE) {
    String err = String.format(
    "Oversize read. max: %d, actual: %d",
    OzoneConsts.OZONE_SCM_CHUNK_MAX_SIZE, len);
    LOG.error(err);
    throw new StorageContainerException(err, UNSUPPORTED_REQUEST);
    }
    return (int) len;
    }

@@ -441,37 +439,39 @@ public static void verifyChecksum(List<ByteBuffer> bufferList, int startIndex, C
public static void validateChecksums(ByteBuffer data, long blockOffset, int startIndex,

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.

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.

Comment on lines -588 to +587
ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData());
Checksum.verifyChecksum(data, checksumData, 0);
Checksum.validateChecksums(data, readBlock.getOffset(), 0, readBlock.getChunkInfoListList());

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 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.

@peterxcli peterxcli Sep 23, 2026 •

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.

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 peterxcli left a comment

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 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

Comment on lines -588 to +587
ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData());
Checksum.verifyChecksum(data, checksumData, 0);
Checksum.validateChecksums(data, readBlock.getOffset(), 0, readBlock.getChunkInfoListList());

@peterxcli peterxcli Sep 23, 2026 •

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.

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.

Comment on lines +2349 to +2353
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));
}

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.

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

  1. 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.
  2. 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:
    public static int limitReadSize(long len)
    throws StorageContainerException {
    if (len > OzoneConsts.OZONE_SCM_CHUNK_MAX_SIZE) {
    String err = String.format(
    "Oversize read. max: %d, actual: %d",
    OzoneConsts.OZONE_SCM_CHUNK_MAX_SIZE, len);
    LOG.error(err);
    throw new StorageContainerException(err, UNSUPPORTED_REQUEST);
    }
    return (int) len;
    }

@peterxcli peterxcli changed the title HDDS-16258. Verify streaming block reads using chunk boundaries HDDS-16258. Refactor streaming block reads and fix checksum verification for variable-sized chunks Sep 27, 2026
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.

5 participants