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)