[ 
https://issues.apache.org/jira/browse/FLINK-40835?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18123552#comment-18123552
 ] 

Amrit Sarkar edited comment on FLINK-40835 at 10/5/26 9:43 PM:
---------------------------------------------------------------

Hi, I'd like to work on this. The fix looks like overriding {{read(byte[], int, 
int)}} in {{CompressibleFSDataInputStream}} to delegate to the wrapped stream, 
plus a test that checks that bulk reads reach the delegate. Could I assign it 
to myself? Thanks.


was (Author: [email protected]):
Hi, I'd like to work on this. The fix looks like overriding {{read(byte[], int, 
int)}} in {{CompressibleFSDataInputStream}} to delegate to the wrapped stream, 
plus a test that checks bulk reads reach the delegate. Could someone assign it 
to me? Thanks.

> Operator state restore reads state files one byte at a time
> -----------------------------------------------------------
>
>                 Key: FLINK-40835
>                 URL: https://issues.apache.org/jira/browse/FLINK-40835
>             Project: Flink
>          Issue Type: Bug
>          Components: Runtime / State Backends
>    Affects Versions: 1.18.0, 1.18.1, 1.19.3, 2.0.2, 2.3.0, 2.2.1, 1.20.5, 
> 2.1.3
>         Environment: MiniCluster, local file system, JDK 17, macOS
>            Reporter: Spoorthi Basu
>            Priority: Major
>
> {{CompressibleFSDataInputStream}} overrides {{read()}} but not {{read(byte[], 
> int, int)}}. Bulk reads therefore fall back to 
> {{java.io.InputStream#read(byte[], int, int)}}, which loops over {{read()}} 
> one byte at a time. The stream it wraps, {{ForwardingInputStream}}, does 
> implement bulk reads, but the call never reaches it.
> The class is used by {{OperatorStateRestoreOperation}} for every operator 
> state restore (list, union and broadcast), which includes the split state of 
> every source reader. On a local file system each byte becomes a separate 
> native read.
> Restoring a single list-state element on Flink 2.2.0 (MiniCluster, local file 
> system, snapshot compression off, which is the default):
> ||State size||Today||With bulk read||
> |30 MB|11,077 ms|140 ms|
> |60 MB|22,136 ms|131 ms|
> |120 MB|44,484 ms|136 ms|
> The same restore on current master (1690c6bed94) took 11,855 ms at 30 MB and 
> 44,323 ms at 120 MB. The read path is the same in every release since 1.18.0 
> (checked at every release tag).
> Every stack sample of the task thread during the slow restore on 2.2.0 was:
> {noformat}
> java.io.FileInputStream.read0(Native Method)
> java.io.FileInputStream.read(FileInputStream.java:228)
> org.apache.flink.core.fs.local.LocalDataInputStream.read(LocalDataInputStream.java:70)
> org.apache.flink.core.fs.FSDataInputStreamWrapper.read(FSDataInputStreamWrapper.java:50)
> org.apache.flink.runtime.util.ForwardingInputStream.read(ForwardingInputStream.java:42)
> org.apache.flink.runtime.state.CompressibleFSDataInputStream.read(CompressibleFSDataInputStream.java:62)
> java.io.InputStream.read(InputStream.java:293)
> java.io.DataInputStream.readFully(DataInputStream.java:201)
> org.apache.flink.api.common.typeutils.base.array.BytePrimitiveArraySerializer.deserialize(BytePrimitiveArraySerializer.java:82)
> org.apache.flink.runtime.state.OperatorStateRestoreOperation.deserializeOperatorStateValues(OperatorStateRestoreOperation.java:236)
> org.apache.flink.runtime.state.OperatorStateRestoreOperation.restore(OperatorStateRestoreOperation.java:207)
> {noformat}
> The class was introduced by FLINK-30113. FLINK-34063 and FLINK-36530 later 
> changed {{seek()}} in the same class and left the read path as it was.
> Delegating the bulk read fixes it. The "with bulk read" column was measured 
> with this change:
> {code:java}
> @Override
> public int read(byte[] b, int off, int len) throws IOException {
>     return compressingDelegate.read(b, off, len);
> }
> {code}
> These measurements are on the local file system. Object store and HDFS 
> clients buffer network reads, so there each byte costs a pass through the 
> wrappers and the client's single-byte {{read()}} rather than a system call. I 
> have not measured that case.
> Related: FLINK-39000 (redundant seeks in the same restore path, fixed in 
> 2.3.0).



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to