diff --git a/.github/trigger_files/beam_PostCommit_Java_Nexmark_Flink.json b/.github/trigger_files/beam_PostCommit_Java_Nexmark_Flink.json index e3d6056a5de9..b26833333238 100644 --- a/.github/trigger_files/beam_PostCommit_Java_Nexmark_Flink.json +++ b/.github/trigger_files/beam_PostCommit_Java_Nexmark_Flink.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run", - "modification": 1 + "modification": 2 } diff --git a/.github/trigger_files/beam_PostCommit_Java_Tpcds_Flink.json b/.github/trigger_files/beam_PostCommit_Java_Tpcds_Flink.json new file mode 100644 index 000000000000..e3d6056a5de9 --- /dev/null +++ b/.github/trigger_files/beam_PostCommit_Java_Tpcds_Flink.json @@ -0,0 +1,4 @@ +{ + "comment": "Modify this file in a trivial way to cause this test suite to run", + "modification": 1 +} diff --git a/.github/workflows/beam_PostCommit_Java_Nexmark_Flink.yml b/.github/workflows/beam_PostCommit_Java_Nexmark_Flink.yml index 4a3574906a28..f702c2e7a82d 100644 --- a/.github/workflows/beam_PostCommit_Java_Nexmark_Flink.yml +++ b/.github/workflows/beam_PostCommit_Java_Nexmark_Flink.yml @@ -109,6 +109,7 @@ jobs: uses: ./.github/actions/gradle-command-self-hosted-action with: gradle-command: :sdks:java:testing:nexmark:run + # TODO: Fix query 6,9 for Flink 2 arguments: | - -Pnexmark.runner=:runners:flink:2.2 \ + -Pnexmark.runner=:runners:flink:1.20 \ "${{ env.GRADLE_COMMAND_ARGUMENTS }}--streaming=${{ matrix.streaming }}" diff --git a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy index 344cf66e01fc..591e72d5944c 100644 --- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy +++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy @@ -954,6 +954,18 @@ class BeamModulePlugin implements Plugin { /** ***********************************************************************************************/ + // Resolves a capability conflict for a given configuration by selecting the preferred capability provider. + project.ext.resolveCapabilitiesConflict = { Configuration configuration, String capability, String preferred -> + configuration.resolutionStrategy.capabilitiesResolution.withCapability(capability) { + def candidate = candidates.find { it.id.toString().contains(preferred) } + if (candidate != null) { + select(candidate) + } else { + selectHighestVersion() + } + } + } + // Returns a string representing the relocated path to be used with the shadow plugin when // given a suffix such as "com.google.common". project.ext.getJavaRelocatedPath = { String suffix -> diff --git a/examples/java/common.gradle b/examples/java/common.gradle index 0e800e7cf987..f32667d733fb 100644 --- a/examples/java/common.gradle +++ b/examples/java/common.gradle @@ -35,16 +35,8 @@ configurations.sparkRunnerPreCommit { exclude group: "org.slf4j", module: "jul-to-slf4j" exclude group: "org.slf4j", module: "slf4j-jdk14" } -configurations.flinkRunnerPreCommit { - resolutionStrategy.capabilitiesResolution.withCapability("org.lz4:lz4-java") { - def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') } - if (candidate != null) { - select(candidate) - } else { - selectHighestVersion() - } - } -} +resolveCapabilitiesConflict(configurations.flinkRunnerPreCommit, 'org.lz4:lz4-java', 'at.yawk.lz4') + dependencies { directRunnerPreCommit project(path: ":runners:direct-java", configuration: "shadow") diff --git a/runners/flink/2.1/build.gradle b/runners/flink/2.1/build.gradle index e9092c2977f7..1e0d565b50d9 100644 --- a/runners/flink/2.1/build.gradle +++ b/runners/flink/2.1/build.gradle @@ -45,12 +45,6 @@ apply from: "../flink_runner.gradle" // Flink 2.1 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java // Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict configurations.all { - resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') { - def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') } - if (candidate != null) { - select(candidate) - } else { - selectHighestVersion() - } - } + resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4') } + diff --git a/runners/flink/2.1/job-server/build.gradle b/runners/flink/2.1/job-server/build.gradle index 277ddc07fdaf..0910fef11208 100644 --- a/runners/flink/2.1/job-server/build.gradle +++ b/runners/flink/2.1/job-server/build.gradle @@ -33,12 +33,6 @@ apply from: "$basePath/flink_job_server.gradle" // Flink 2.1 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java // Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict configurations.all { - resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') { - def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') } - if (candidate != null) { - select(candidate) - } else { - selectHighestVersion() - } - } + resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4') } + diff --git a/runners/flink/2.2/build.gradle b/runners/flink/2.2/build.gradle index 3856d20a3512..0321dcf42d10 100644 --- a/runners/flink/2.2/build.gradle +++ b/runners/flink/2.2/build.gradle @@ -60,12 +60,6 @@ apply from: "../flink_runner.gradle" // Flink 2.2 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java // Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict configurations.all { - resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') { - def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') } - if (candidate != null) { - select(candidate) - } else { - selectHighestVersion() - } - } + resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4') } + diff --git a/runners/flink/2.2/job-server/build.gradle b/runners/flink/2.2/job-server/build.gradle index 3f6b184e3fca..f116f5a1dcbf 100644 --- a/runners/flink/2.2/job-server/build.gradle +++ b/runners/flink/2.2/job-server/build.gradle @@ -33,12 +33,6 @@ apply from: "$basePath/flink_job_server.gradle" // Flink 2.2 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java // Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict configurations.all { - resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') { - def candidate = candidates.find { it.id.toString().contains('at.yawk.lz4') } - if (candidate != null) { - select(candidate) - } else { - selectHighestVersion() - } - } + resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4') } + diff --git a/sdks/java/testing/nexmark/build.gradle b/sdks/java/testing/nexmark/build.gradle index ec3447579613..bd917a3935ad 100644 --- a/sdks/java/testing/nexmark/build.gradle +++ b/sdks/java/testing/nexmark/build.gradle @@ -39,6 +39,17 @@ def nexmarkRunnerDependency = project.findProperty(nexmarkRunnerProperty) def nexmarkRunnerVersionProperty = "nexmark.runner.version" def nexmarkRunnerVersion = project.findProperty(nexmarkRunnerVersionProperty) def isSparkRunner = nexmarkRunnerDependency.startsWith(":runners:spark:") +def isFlink2 = nexmarkRunnerDependency.startsWith(":runners:flink:") && { + def version = nexmarkRunnerDependency.substring(":runners:flink:".length()) + try { + def parts = version.split('\\.') + def major = parts[0].toInteger() + def minor = parts.length > 1 ? parts[1].toInteger() : 0 + return (major > 2) || (major == 2 && minor >= 1) + } catch (Exception e) { + return false + } +}() def isDataflowRunner = ":runners:google-cloud-dataflow-java".equals(nexmarkRunnerDependency) def isDataflowRunnerV2 = isDataflowRunner && "V2".equals(nexmarkRunnerVersion) def runnerConfiguration = ":runners:direct-java".equals(nexmarkRunnerDependency) ? "shadow" : null @@ -92,6 +103,10 @@ dependencies { gradleRun project(path: nexmarkRunnerDependency, configuration: runnerConfiguration) } +if (isFlink2) { + resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 'at.yawk.lz4') +} + if (isSparkRunner) { configurations.gradleRun { // Using Spark runner causes a StackOverflowError if slf4j-jdk14 is on the classpath diff --git a/sdks/java/testing/tpcds/build.gradle b/sdks/java/testing/tpcds/build.gradle index 1714e6e61a80..6a05df96ab20 100644 --- a/sdks/java/testing/tpcds/build.gradle +++ b/sdks/java/testing/tpcds/build.gradle @@ -34,6 +34,17 @@ def tpcdsRunnerProperty = "tpcds.runner" def tpcdsRunnerDependency = project.findProperty(tpcdsRunnerProperty) ?: ":runners:direct-java" def isSpark = tpcdsRunnerDependency.startsWith(":runners:spark:") +def isFlink2 = tpcdsRunnerDependency.startsWith(":runners:flink:") && { + def version = tpcdsRunnerDependency.substring(":runners:flink:".length()) + try { + def parts = version.split('\\.') + def major = parts[0].toInteger() + def minor = parts.length > 1 ? parts[1].toInteger() : 0 + return (major > 2) || (major == 2 && minor >= 1) + } catch (Exception e) { + return false + } +}() def isDataflowRunner = ":runners:google-cloud-dataflow-java".equals(tpcdsRunnerDependency) def runnerConfiguration = ":runners:direct-java".equals(tpcdsRunnerDependency) ? "shadow" : null @@ -83,6 +94,10 @@ dependencies { gradleRun project(path: tpcdsRunnerDependency, configuration: runnerConfiguration) } +if (isFlink2) { + resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 'at.yawk.lz4') +} + if (isSpark) { configurations.gradleRun { exclude group: "org.slf4j", module: "slf4j-jdk14"