Skip to content

Commit 9a9f48b

Browse files
authored
Fix Flink 2.1+ dependency resolution (#39061)
1 parent 117bb89 commit 9a9f48b

11 files changed

Lines changed: 59 additions & 44 deletions

File tree

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
11
{
22
"comment": "Modify this file in a trivial way to cause this test suite to run",
3-
"modification": 1
3+
"modification": 2
44
}
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
{
2+
"comment": "Modify this file in a trivial way to cause this test suite to run",
3+
"modification": 1
4+
}

.github/workflows/beam_PostCommit_Java_Nexmark_Flink.yml

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,7 @@ jobs:
109109
uses: ./.github/actions/gradle-command-self-hosted-action
110110
with:
111111
gradle-command: :sdks:java:testing:nexmark:run
112+
# TODO: Fix query 6,9 for Flink 2
112113
arguments: |
113-
-Pnexmark.runner=:runners:flink:2.2 \
114+
-Pnexmark.runner=:runners:flink:1.20 \
114115
"${{ env.GRADLE_COMMAND_ARGUMENTS }}--streaming=${{ matrix.streaming }}"

buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -954,6 +954,18 @@ class BeamModulePlugin implements Plugin<Project> {
954954

955955
/** ***********************************************************************************************/
956956

957+
// Resolves a capability conflict for a given configuration by selecting the preferred capability provider.
958+
project.ext.resolveCapabilitiesConflict = { Configuration configuration, String capability, String preferred ->
959+
configuration.resolutionStrategy.capabilitiesResolution.withCapability(capability) {
960+
def candidate = candidates.find { it.id.toString().contains(preferred) }
961+
if (candidate != null) {
962+
select(candidate)
963+
} else {
964+
selectHighestVersion()
965+
}
966+
}
967+
}
968+
957969
// Returns a string representing the relocated path to be used with the shadow plugin when
958970
// given a suffix such as "com.google.common".
959971
project.ext.getJavaRelocatedPath = { String suffix ->

examples/java/common.gradle

Lines changed: 2 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -35,16 +35,8 @@ configurations.sparkRunnerPreCommit {
3535
exclude group: "org.slf4j", module: "jul-to-slf4j"
3636
exclude group: "org.slf4j", module: "slf4j-jdk14"
3737
}
38-
configurations.flinkRunnerPreCommit {
39-
resolutionStrategy.capabilitiesResolution.withCapability("org.lz4:lz4-java") {
40-
def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') }
41-
if (candidate != null) {
42-
select(candidate)
43-
} else {
44-
selectHighestVersion()
45-
}
46-
}
47-
}
38+
resolveCapabilitiesConflict(configurations.flinkRunnerPreCommit, 'org.lz4:lz4-java', 'at.yawk.lz4')
39+
4840

4941
dependencies {
5042
directRunnerPreCommit project(path: ":runners:direct-java", configuration: "shadow")

runners/flink/2.1/build.gradle

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -45,12 +45,6 @@ apply from: "../flink_runner.gradle"
4545
// Flink 2.1 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
4646
// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
4747
configurations.all {
48-
resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') {
49-
def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') }
50-
if (candidate != null) {
51-
select(candidate)
52-
} else {
53-
selectHighestVersion()
54-
}
55-
}
48+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
5649
}
50+

runners/flink/2.1/job-server/build.gradle

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -33,12 +33,6 @@ apply from: "$basePath/flink_job_server.gradle"
3333
// Flink 2.1 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
3434
// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
3535
configurations.all {
36-
resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') {
37-
def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') }
38-
if (candidate != null) {
39-
select(candidate)
40-
} else {
41-
selectHighestVersion()
42-
}
43-
}
36+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
4437
}
38+

runners/flink/2.2/build.gradle

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -60,12 +60,6 @@ apply from: "../flink_runner.gradle"
6060
// Flink 2.2 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
6161
// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
6262
configurations.all {
63-
resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') {
64-
def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') }
65-
if (candidate != null) {
66-
select(candidate)
67-
} else {
68-
selectHighestVersion()
69-
}
70-
}
63+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
7164
}
65+

runners/flink/2.2/job-server/build.gradle

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -33,12 +33,6 @@ apply from: "$basePath/flink_job_server.gradle"
3333
// Flink 2.2 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
3434
// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
3535
configurations.all {
36-
resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') {
37-
def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') }
38-
if (candidate != null) {
39-
select(candidate)
40-
} else {
41-
selectHighestVersion()
42-
}
43-
}
36+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
4437
}
38+

sdks/java/testing/nexmark/build.gradle

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,17 @@ def nexmarkRunnerDependency = project.findProperty(nexmarkRunnerProperty)
3939
def nexmarkRunnerVersionProperty = "nexmark.runner.version"
4040
def nexmarkRunnerVersion = project.findProperty(nexmarkRunnerVersionProperty)
4141
def isSparkRunner = nexmarkRunnerDependency.startsWith(":runners:spark:")
42+
def isFlink2 = nexmarkRunnerDependency.startsWith(":runners:flink:") && {
43+
def version = nexmarkRunnerDependency.substring(":runners:flink:".length())
44+
try {
45+
def parts = version.split('\\.')
46+
def major = parts[0].toInteger()
47+
def minor = parts.length > 1 ? parts[1].toInteger() : 0
48+
return (major > 2) || (major == 2 && minor >= 1)
49+
} catch (Exception e) {
50+
return false
51+
}
52+
}()
4253
def isDataflowRunner = ":runners:google-cloud-dataflow-java".equals(nexmarkRunnerDependency)
4354
def isDataflowRunnerV2 = isDataflowRunner && "V2".equals(nexmarkRunnerVersion)
4455
def runnerConfiguration = ":runners:direct-java".equals(nexmarkRunnerDependency) ? "shadow" : null
@@ -92,6 +103,10 @@ dependencies {
92103
gradleRun project(path: nexmarkRunnerDependency, configuration: runnerConfiguration)
93104
}
94105

106+
if (isFlink2) {
107+
resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 'at.yawk.lz4')
108+
}
109+
95110
if (isSparkRunner) {
96111
configurations.gradleRun {
97112
// Using Spark runner causes a StackOverflowError if slf4j-jdk14 is on the classpath

0 commit comments

Comments
 (0)