Spoorthi Basu created FLINK-40835:
-------------------------------------

             Summary: 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: 2.1.3, 1.20.5, 2.2.1, 2.3.0, 2.0.2, 1.19.3, 1.18.1, 1.18.0
         Environment: MiniCluster, local file system, JDK 17, macOS
            Reporter: Spoorthi Basu


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