dalelane opened a new pull request, #29371:
URL: https://github.com/apache/flink/pull/29371
## What is the purpose of the change
This pull request addresses scenarios where a Flink jar fetched by
the artifact fetcher might not provide the expected or needed jar.
Artifact fetchers skip a fetch when the target file already exists.
Streaming directly to the target file means that if the fetcher
process dies mid-download (e.g. connection reset, Job Manager
process terminated, etc.) then a truncated artifact is left
behind and treated as if it is a complete one on the next fetch
attempt.
Even when fetched artifacts are complete and valid jars, the way
that fetched artifacts are cached in user.artifacts.base-dir
is based on filename. There are times where this could result in
unexpected and undesired behaviour:
- Redeploys where user.artifacts.base-dir is in persistent storage
Changing the job ...artifactstore/app.jar?v=1 to
...artifactstore/app.jar?v=2 will reuse the v1 jar and never
fetch the v2 jar
- Standalone clusters that share a base-dir, where one job wants
somehost/app.jar and another job wants differenthost/app.jar
could inadvertently both run the first jar to be fetched
- Multiple artifacts with the same file name, such as a job that
uses s3://bucket-one/udf.jar and s3://bucket-two/udf.jar could
reuse the first jar to be fetched for both
- Jars that are uniquely identified by query parameters, such as
an artifact store that has .../download?id=X and
.../download?id=Y to identify unrelated jars, but where every
URI has a common last path segment such as "download"
## Brief change log
Artifacts are now fetched into .part files and only moved to the
target file name once the fetch is completed, so the target file
location never contains partial artifacts.
- Added ArtifactUtils.copyToFileWhenComplete which copies the
stream into a part file next to the target then does an atomic
move to the target file once complete
- Updated FsArtifactFetcher and HttpArtifactFetcher to use the
util instead of FileUtils.copyToFile
Artifacts are now stored in a subdirectory of
user.artifacts.base-dir, with the subdirectory given a name from a
hash of it's full source URI, including any query parameters.
## Verifying this change
This change added tests:
- Added ArtifactUtilsTest cases that check the target name is not
used while fetching is underway, that completed copies don't
leave anything else in the directory, and that no files get left
behind if the copy fails
- Added a test that simulates an HTTP server dropping a connection
during a download, checking that no target file is left and that
a retry successfully fetches the complete artifact
- Added a test that covers URIs with no file names in the path
- Added a counting fetcher that keeps track of how many times it
fetches a jar to verify when jars are reused
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency): no
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: no
- The serializers: no
- The runtime per-record code paths (performance sensitive): no
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
- The S3 file system connector: yes
## Documentation
- Does this pull request introduce a new feature? no
- If yes, how is the feature documented? not applicable
---
##### Was generative AI tooling used to co-author this PR?
- [ ] Yes (please specify the tool below)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]