This is an automated email from the ASF dual-hosted git repository.
hong pushed a commit to branch release-1.19
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/release-1.19 by this push:
new 90a71a12477 [FLINK-35358][clients] Reintroduce recursive JAR listing
in classpath load from "usrlib"
90a71a12477 is described below
commit 90a71a124771697a0b8b2c2bbc520856d6ae9e25
Author: Ferenc Csaky <[email protected]>
AuthorDate: Wed May 15 16:39:30 2024 +0200
[FLINK-35358][clients] Reintroduce recursive JAR listing in classpath load
from "usrlib"
---
.../program/DefaultPackagedProgramRetriever.java | 10 ++++----
.../DefaultPackagedProgramRetrieverITCase.java | 27 ++++++++++++++++++++++
2 files changed, 32 insertions(+), 5 deletions(-)
diff --git
a/flink-clients/src/main/java/org/apache/flink/client/program/DefaultPackagedProgramRetriever.java
b/flink-clients/src/main/java/org/apache/flink/client/program/DefaultPackagedProgramRetriever.java
index 654db0200c1..5e286afc4c3 100644
---
a/flink-clients/src/main/java/org/apache/flink/client/program/DefaultPackagedProgramRetriever.java
+++
b/flink-clients/src/main/java/org/apache/flink/client/program/DefaultPackagedProgramRetriever.java
@@ -114,7 +114,7 @@ public class DefaultPackagedProgramRetriever implements
PackagedProgramRetriever
List<URL> userClasspaths;
try {
final List<URL> classpathsFromUserLibDir =
- getClasspathsFromUserDir(userLibDir, jarFile);
+ getClasspathsFromUserLibDir(userLibDir, jarFile);
final List<URL> classpathsFromUserArtifactDir =
getClasspathsFromArtifacts(userArtifacts, jarFile);
final List<URL> classpathsFromConfiguration =
@@ -231,13 +231,13 @@ public class DefaultPackagedProgramRetriever implements
PackagedProgramRetriever
}
}
- private static List<URL> getClasspathsFromUserDir(
- @Nullable File userDir, @Nullable File jarFile) throws IOException
{
- if (userDir == null) {
+ private static List<URL> getClasspathsFromUserLibDir(
+ @Nullable File userLibDir, @Nullable File jarFile) throws
IOException {
+ if (userLibDir == null) {
return Collections.emptyList();
}
- try (Stream<Path> files = Files.list(userDir.toPath())) {
+ try (Stream<Path> files = Files.walk(userLibDir.toPath())) {
return getClasspathsFromArtifacts(files, jarFile);
}
}
diff --git
a/flink-clients/src/test/java/org/apache/flink/client/program/DefaultPackagedProgramRetrieverITCase.java
b/flink-clients/src/test/java/org/apache/flink/client/program/DefaultPackagedProgramRetrieverITCase.java
index 431bba61613..6fbe706aaaf 100644
---
a/flink-clients/src/test/java/org/apache/flink/client/program/DefaultPackagedProgramRetrieverITCase.java
+++
b/flink-clients/src/test/java/org/apache/flink/client/program/DefaultPackagedProgramRetrieverITCase.java
@@ -495,6 +495,33 @@ class DefaultPackagedProgramRetrieverITCase {
assertThat(actualClasspath).isEqualTo(expectedClasspath);
}
+ @Test
+ void testRetrieveFromJarFileWithNonRootUserLib()
+ throws IOException, FlinkException, ProgramInvocationException {
+ final PackagedProgramRetriever retrieverUnderTest =
+ DefaultPackagedProgramRetriever.create(
+
singleEntryClassClasspathProvider.getDirectory().getParentFile(),
+ // the testJob jar is not on the user classpath
+ testJobEntryClassClasspathProvider.getJobJar(),
+ null,
+ null,
+
ClasspathProviderExtension.parametersForTestJob("suffix"),
+ new Configuration());
+ final JobGraph jobGraph = retrieveJobGraph(retrieverUnderTest, new
Configuration());
+
+ assertThat(jobGraph.getUserJars())
+ .contains(
+ new org.apache.flink.core.fs.Path(
+
testJobEntryClassClasspathProvider.getJobJar().toURI()));
+ final List<String> actualClasspath =
+
jobGraph.getClasspaths().stream().map(URL::toString).collect(Collectors.toList());
+ final List<String> expectedClasspath =
+ extractRelativizedURLsForJarsFromDirectory(
+ singleEntryClassClasspathProvider.getDirectory());
+
+ assertThat(actualClasspath).isEqualTo(expectedClasspath);
+ }
+
@Test
void testRetrieveFromJarFileWithArtifacts()
throws IOException, FlinkException, ProgramInvocationException {