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)

Reply via email to