andygrove commented on code in PR #6805:
URL: https://github.com/apache/datafusion-comet/pull/6805#discussion_r4232361656


##########
spark/src/main/java/org/apache/comet/CometShuffleBlockIterator.java:
##########
@@ -77,17 +81,15 @@ public int hasNext() throws IOException {
     }
 
     // Read 16-byte header: clear() resets position=0, limit=capacity,
-    // preparing the buffer for channel.read() to fill it
+    // preparing the buffer for readFully() to fill it
     headerBuf.clear();
-    while (headerBuf.hasRemaining()) {
-      int bytesRead = channel.read(headerBuf);
-      if (bytesRead < 0) {
-        if (headerBuf.position() == 0) {
-          close();
-          return -1;
-        }
-        throw new EOFException("Data corrupt: unexpected EOF while reading 
batch header");
+    readFully(inputStream, headerBuf, readBufferSize);

Review Comment:
   Dropping `Channels.newChannel` also drops the only thing that stopped a 
killed task on the direct read path. The JDK channel checks the thread's 
interrupt status on every read and throws `ClosedByInterruptException`, and 
native calls `hasNext()` on the task thread. `readAsRawStream()` in 
`CometBlockStoreShuffleReader` returns a bare `SequenceInputStream`, so this 
path has no `InterruptibleIterator`. With this change a killed task keeps 
reading until the shuffle block it is on runs out.
   
   I measured it locally with a 20M row shuffle into a single block, read by a 
native project, and killed the reduce task 1 s into its 5.6 s read. With an 
interrupting kill, which is what speculation and Spark Connect and Thrift 
server cancellation use, the task stopped 11 ms after the kill on main and kept 
going for 4.6 s on this branch. The skewed partition this PR is aimed at is 
also the task most likely to get a speculative copy, and the losing copy now 
holds its core and memory until it finishes reading its current block.
   
   Could `readAsRawStream()` wrap the stream so that each read calls 
`context.killTaskIfInterrupted()`, the way the Celeborn reader's stream already 
does? I tried that on top of this branch. Kills with and without an interrupt 
both stopped within about 10 ms, and a task that isn't killed took the same 
time. A test that marks the task context interrupted and checks that the next 
read throws `TaskKilledException` would keep this from regressing again.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to