[
https://issues.apache.org/jira/browse/FLINK-40916?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Martijn Visser updated FLINK-40916:
-----------------------------------
Description:
_*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.
was:
_*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.
> 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: Runtime / Checkpointing
> Reporter: Isaac Montenegro Jimenez
> Assignee: Isaac Montenegro Jimenez
> Priority: Minor
> Labels: flink, pull-request-available
>
> _*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.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)