Isaac Montenegro Jimenez created FLINK-40916:
------------------------------------------------
Summary: Optionally skip eager initialisation of checkpoint
locations
Key: FLINK-40916
URL: https://issues.apache.org/jira/browse/FLINK-40916
Project: Flink
Issue Type: Improvement
Components: flink-contrib
Reporter: Isaac Montenegro Jimenez
_*Summary*_
In the context of [Apache Flink|https://github.com/apache/flink]
During startup of the CheckpointCoordinator component there is an eager file
system initialisation that creates directories used for checkpoint locations.
This eager initialisation is blocking for startup, from internal measurements
we’ve seen it can take from ~180ms to ~600ms depending on the file system.
There is a fallback creation that happens later on once the first checkpoint
gets triggered.
We are interested in skipping this initial initialisation as we want customers
to start processing their data faster. The tradeoff is twofold here:
# This eager initialisation provides fail-fast so that if the directories
can’t be created, the job will fail right away.
## If initialisation fails at checkpoint time, that checkpoint fails with an
IO_EXCEPTION, which counts toward
execution.checkpointing.tolerable-failed-checkpoints. Once exceeded, the job
fails over according to its restart strategy. (also see *Related bug* below)
# We are moving directory creation time from startup to the first checkpoint
time.
The intention would be to add a configuration option to allow Flink users to
skip this eager initialisation.
*Detailed code paths*
During construction, if periodic checkpoints are configured, there is a call to
[“initializeBaseLocationsForCheckpoint()"|https://github.com/apache/flink/blob/769d9fe34d5615801605a64aa80498bb23b7cb33/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java#L354-L357]
where the directories for the checkpoint locations get initialized. With
FsCheckpointStorageAccess we call mkdir in the [underlying
filesystem|https://github.com/apache/flink/blob/769d9fe/flink-runtime/src/main/java/org/apache/flink/runtime/state/filesystem/FsCheckpointStorageAccess.java#L147].
In some cases, like with the
[GoogleHadoopFileSystem|https://github.com/GoogleCloudDataproc/hadoop-connectors/blob/v3.1.11/gcs/src/main/java/com/google/cloud/hadoop/fs/gcs/GoogleHadoopFileSystem.java#L916]
this ends up being several network calls.
Later on, when the first checkpoint is triggered, if the base locations have
not been created yet, they will be created as part of the
[startTriggeringCheckpoint
function|https://github.com/apache/flink/blob/769d9fe34d5615801605a64aa80498bb23b7cb33/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java#L625],
meaning skipping the early initialisation just moves this directory creation
time.
*Related bug*
While investigating this, I found there is a bug hidden by this eager
initialisation that we might as well fix right away.
In the [startTriggeringCheckpoint
function|https://github.com/apache/flink/blob/769d9fe34d5615801605a64aa80498bb23b7cb33/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java#L641C13-L641C50]
if the checkpoint locations have not been initialised yet, they get
initialised as part of the [initializeCheckpointLocation
call|https://github.com/apache/flink/blob/769d9fe34d5615801605a64aa80498bb23b7cb33/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java#L679].
This is controlled by a flag (baseLocationsForCheckpointInitialized,
introduced in [PR 17278|https://github.com/apache/flink/pull/17278] for
[FLINK-24280|https://issues.apache.org/jira/browse/FLINK-24280] ). This flag
gets flipped to true eagerly, before it’s confirmed that the operation was
successful. There’s also an edge case where if startTriggeringCheckpoint gets
run for a savepoint (not a checkpoint) the [initialisation code gets
skipped|https://github.com/apache/flink/blob/769d9fe34d5615801605a64aa80498bb23b7cb33/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java#L881],
but the flag has been turned on. As part of this proposed change I’d propose
also moving the flag change to the initializeCheckpointLocation.
This bug might be manifesting right now if periodic checkpointing is disabled,
a savepoint is taken first, and a checkpoint is later triggered manually (e.g.
via the REST API). Impact would depend on the filesystem.
*Testing done*
I did testing in the internal Confluent fork
([https://github.com/confluentinc/flink] ) against jobs in several filesystems
(AWS S3 via Presto, ADLS Gen2 via an internal native filesystem (not the
upstream flink-azure-fs-hadoop), GCS via gcs-connector) with a straightforward
change that removes the eager filesystem initialisation and applies the flag
fix. We saw the expected startup reduction and directory creation moved to
first checkpoint creation time.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)