diff --git a/.test-infra/jenkins/job_PostCommit_Java_PortableValidatesRunner_Spark_Batch.groovy b/.test-infra/jenkins/job_PostCommit_Java_PortableValidatesRunner_Spark_Batch.groovy index 3da5fed831b2..ce7340f64da8 100644 --- a/.test-infra/jenkins/job_PostCommit_Java_PortableValidatesRunner_Spark_Batch.groovy +++ b/.test-infra/jenkins/job_PostCommit_Java_PortableValidatesRunner_Spark_Batch.groovy @@ -38,6 +38,8 @@ PostcommitJobBuilder.postCommitJob('beam_PostCommit_Java_PVR_Spark_Batch', rootBuildScriptDir(commonJobProperties.checkoutDir) tasks(':runners:spark:2:job-server:validatesPortableRunnerBatch') tasks(':runners:spark:3:job-server:validatesPortableRunnerBatch') + tasks(':runners:spark:2:job-server:validatesPortableRunnerDocker') + tasks(':runners:spark:3:job-server:validatesPortableRunnerDocker') commonJobProperties.setGradleSwitches(delegate) } } diff --git a/.test-infra/jenkins/job_PreCommit_Java_PortableValidatesRunner_Flink_Docker.groovy b/.test-infra/jenkins/job_PreCommit_Java_PortableValidatesRunner_Flink_Docker.groovy new file mode 100644 index 000000000000..bb14a792291c --- /dev/null +++ b/.test-infra/jenkins/job_PreCommit_Java_PortableValidatesRunner_Flink_Docker.groovy @@ -0,0 +1,41 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import CommonTestProperties +import PrecommitJobBuilder + +// This job runs a limited subset of ValidatesRunner tests against the Flink runner in the docker environment. +PrecommitJobBuilder builder = new PrecommitJobBuilder( + scope: this, + nameBase: 'Java_PVR_Flink_Docker', + gradleTask: ":runners:flink:${CommonTestProperties.getFlinkVersion()}:job-server:validatesPortableRunnerDocker", + timeoutMins: 240, + triggerPathPatterns: [ + '^sdks/java/core/src/test/java/org/apache/beam/sdk/.*$', + '^sdks/java/container/.*$', + '^sdks/java/harness/.*$', + '^runners/flink/.*$', + '^runners/java-fn-execution/.*$', + ], + ) +builder.build { + // Publish all test results to Jenkins. + publishers { + archiveJunit('**/build/test-results/**/*.xml') + } +} 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 71c65eb1f31f..b957bc61e894 100644 --- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy +++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy @@ -620,6 +620,10 @@ class BeamModulePlugin implements Plugin { jackson_datatype_joda : "com.fasterxml.jackson.datatype:jackson-datatype-joda:$jackson_version", jackson_module_scala_2_11 : "com.fasterxml.jackson.module:jackson-module-scala_2.11:$jackson_version", jackson_module_scala_2_12 : "com.fasterxml.jackson.module:jackson-module-scala_2.12:$jackson_version", + // Swap to use the officially published version of 0.4.x once available + // instead of relying on a community published copy. See + // https://github.com/jbellis/jamm/issues/44 for additional details. + jamm : 'io.github.stephankoelle:jamm:0.4.1', jaxb_api : "jakarta.xml.bind:jakarta.xml.bind-api:$jaxb_api_version", jaxb_impl : "com.sun.xml.bind:jaxb-impl:$jaxb_api_version", joda_time : "joda-time:joda-time:2.10.10", diff --git a/runners/direct-java/build.gradle b/runners/direct-java/build.gradle index ad8b36083a5d..da9351cb797f 100644 --- a/runners/direct-java/build.gradle +++ b/runners/direct-java/build.gradle @@ -127,6 +127,8 @@ task needsRunnerTests(type: Test) { testClassesDirs += files(project(":sdks:java:core").sourceSets.test.output.classesDirs) useJUnit { includeCategories "org.apache.beam.sdk.testing.NeedsRunner" + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories "org.apache.beam.sdk.testing.LargeKeys\$Above100MB" // MetricsPusher isn't implemented in direct runner excludeCategories "org.apache.beam.sdk.testing.UsesMetricsPusher" @@ -160,6 +162,8 @@ task validatesRunner(type: Test) { testClassesDirs += files(project(":sdks:java:core").sourceSets.test.output.classesDirs) useJUnit { includeCategories "org.apache.beam.sdk.testing.ValidatesRunner" + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories "org.apache.beam.sdk.testing.LargeKeys\$Above100MB" excludeCategories 'org.apache.beam.sdk.testing.UsesMetricsPusher' excludeCategories "org.apache.beam.sdk.testing.UsesCrossLanguageTransforms" diff --git a/runners/flink/flink_runner.gradle b/runners/flink/flink_runner.gradle index 608a24e11dd9..5c33345e82db 100644 --- a/runners/flink/flink_runner.gradle +++ b/runners/flink/flink_runner.gradle @@ -219,6 +219,8 @@ def createValidatesRunnerTask(Map m) { excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' } else { includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories 'org.apache.beam.sdk.testing.UsesBundleFinalizer' } excludeCategories 'org.apache.beam.sdk.testing.FlattenWithHeterogeneousCoders' diff --git a/runners/flink/job-server/flink_job_server.gradle b/runners/flink/job-server/flink_job_server.gradle index 1dcb837d7715..a6c98b379050 100644 --- a/runners/flink/job-server/flink_job_server.gradle +++ b/runners/flink/job-server/flink_job_server.gradle @@ -126,7 +126,7 @@ runShadow { jvmArgs += ["-Dorg.slf4j.simpleLogger.defaultLogLevel=${project.property('logLevel')}"] } -def portableValidatesRunnerTask(String name, Boolean streaming, Boolean checkpointing) { +def portableValidatesRunnerTask(String name, boolean streaming, boolean checkpointing, boolean docker) { def pipelineOptions = [ // Limit resource consumption via parallelism "--parallelism=2", @@ -145,15 +145,24 @@ def portableValidatesRunnerTask(String name, Boolean streaming, Boolean checkpoi testClasspathConfiguration: configurations.validatesPortableRunner, numParallelTests: 1, pipelineOpts: pipelineOptions, - environment: BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.EMBEDDED, + environment: docker ? BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.DOCKER : BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.EMBEDDED, testCategories: { - if (streaming && checkpointing) { - includeCategories 'org.apache.beam.sdk.testing.UsesBundleFinalizer' - excludeCategories 'org.apache.beam.sdk.testing.UsesBoundedSplittableParDo' - // TestStreamSource does not support checkpointing - excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' - } else { + if (docker) { + includeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' + return + } + + if (streaming && checkpointing) { + includeCategories 'org.apache.beam.sdk.testing.UsesBundleFinalizer' + excludeCategories 'org.apache.beam.sdk.testing.UsesBoundedSplittableParDo' + // TestStreamSource does not support checkpointing + excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' + return + } + includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories 'org.apache.beam.sdk.testing.FlattenWithHeterogeneousCoders' // Larger keys are possible, but they require more memory. excludeCategories 'org.apache.beam.sdk.testing.LargeKeys$Above10MB' @@ -176,14 +185,14 @@ def portableValidatesRunnerTask(String name, Boolean streaming, Boolean checkpoi excludeCategories 'org.apache.beam.sdk.testing.UsesTestStreamWithProcessingTime' excludeCategories 'org.apache.beam.sdk.testing.UsesTestStreamWithMultipleStages' excludeCategories 'org.apache.beam.sdk.testing.UsesTestStreamWithOutputTimestamp' - } else { - excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo' - excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedPCollections' - excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' - excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderedDelivery' - excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderInBundle' + return } - } + + excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo' + excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedPCollections' + excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' + excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderedDelivery' + excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderInBundle' }, testFilter: { // TODO(BEAM-10016) @@ -200,11 +209,13 @@ def portableValidatesRunnerTask(String name, Boolean streaming, Boolean checkpoi ) } -project.ext.validatesPortableRunnerBatch = portableValidatesRunnerTask("Batch", false, false) -project.ext.validatesPortableRunnerStreaming = portableValidatesRunnerTask("Streaming", true, false) -project.ext.validatesPortableRunnerStreamingCheckpoint = portableValidatesRunnerTask("StreamingCheckpointing", true, true) +project.ext.validatesPortableRunnerDocker = portableValidatesRunnerTask("Docker", false, false, true) +project.ext.validatesPortableRunnerBatch = portableValidatesRunnerTask("Batch", false, false, false) +project.ext.validatesPortableRunnerStreaming = portableValidatesRunnerTask("Streaming", true, false, false) +project.ext.validatesPortableRunnerStreamingCheckpoint = portableValidatesRunnerTask("StreamingCheckpointing", true, true, false) task validatesPortableRunner() { + dependsOn validatesPortableRunnerDocker dependsOn validatesPortableRunnerBatch dependsOn validatesPortableRunnerStreaming dependsOn validatesPortableRunnerStreamingCheckpoint diff --git a/runners/google-cloud-dataflow-java/build.gradle b/runners/google-cloud-dataflow-java/build.gradle index cbcbceb59adb..35bfc8bfc2ea 100644 --- a/runners/google-cloud-dataflow-java/build.gradle +++ b/runners/google-cloud-dataflow-java/build.gradle @@ -166,6 +166,8 @@ def runnerV2PipelineOptions = [ ] def commonLegacyExcludeCategories = [ + // Should be run only in a properly configured SDK harness environment + 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment', 'org.apache.beam.sdk.testing.LargeKeys$Above10MB', 'org.apache.beam.sdk.testing.UsesAttemptedMetrics', 'org.apache.beam.sdk.testing.UsesCrossLanguageTransforms', diff --git a/runners/jet/build.gradle b/runners/jet/build.gradle index f8808a0a07ea..75b60b18ebaf 100644 --- a/runners/jet/build.gradle +++ b/runners/jet/build.gradle @@ -75,6 +75,8 @@ task validatesRunnerBatch(type: Test) { testClassesDirs = files(project(":sdks:java:core").sourceSets.test.output.classesDirs) useJUnit { includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories "org.apache.beam.sdk.testing.LargeKeys\$Above100MB" excludeCategories 'org.apache.beam.sdk.testing.UsesTimerMap' excludeCategories 'org.apache.beam.sdk.testing.UsesOnWindowExpiration' diff --git a/runners/portability/java/build.gradle b/runners/portability/java/build.gradle index 21c3c44902b2..8fb935441406 100644 --- a/runners/portability/java/build.gradle +++ b/runners/portability/java/build.gradle @@ -137,6 +137,8 @@ def createUlrValidatesRunnerTask = { name, environmentType, dockerImageTask = "" testClassesDirs = files(project(":sdks:java:core").sourceSets.test.output.classesDirs) useJUnit { includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories 'org.apache.beam.sdk.testing.UsesGaugeMetrics' excludeCategories 'org.apache.beam.sdk.testing.UsesOnWindowExpiration' excludeCategories 'org.apache.beam.sdk.testing.UsesMapState' diff --git a/runners/samza/build.gradle b/runners/samza/build.gradle index bb215c686cec..c244f84ad0b3 100644 --- a/runners/samza/build.gradle +++ b/runners/samza/build.gradle @@ -98,6 +98,8 @@ task validatesRunner(type: Test) { useJUnit { includeCategories 'org.apache.beam.sdk.testing.NeedsRunner' includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo' excludeCategories 'org.apache.beam.sdk.testing.UsesSchema' excludeCategories 'org.apache.beam.sdk.testing.LargeKeys$Above100MB' diff --git a/runners/samza/job-server/build.gradle b/runners/samza/job-server/build.gradle index fc0198b5d961..bdc5722a8544 100644 --- a/runners/samza/job-server/build.gradle +++ b/runners/samza/job-server/build.gradle @@ -58,21 +58,29 @@ runShadow { args = [] } -def tempDir = File.createTempDir() -def pipelineOptions = [ - "--configOverride={\"job.non-logged.store.base.dir\":\"" + tempDir + "\"}" -] -createPortableValidatesRunnerTask( - name: "validatesPortableRunner", - jobServerDriver: "org.apache.beam.runners.samza.SamzaJobServerDriver", - jobServerConfig: "--job-host=localhost,--job-port=0,--artifact-port=0,--expansion-port=0", - testClasspathConfiguration: configurations.validatesPortableRunner, - numParallelTests: 1, - pipelineOpts: pipelineOptions, - environment: BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.EMBEDDED, - testCategories: { +def portableValidatesRunnerTask(String name, boolean docker) { + def tempDir = File.createTempDir() + def pipelineOptions = [ + "--configOverride={\"job.non-logged.store.base.dir\":\"" + tempDir + "\"}" + ] + createPortableValidatesRunnerTask( + name: "validatesPortableRunner${name}", + jobServerDriver: "org.apache.beam.runners.samza.SamzaJobServerDriver", + jobServerConfig: "--job-host=localhost,--job-port=0,--artifact-port=0,--expansion-port=0", + testClasspathConfiguration: configurations.validatesPortableRunner, + numParallelTests: 1, + pipelineOpts: pipelineOptions, + environment: docker ? BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.DOCKER : BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.EMBEDDED, + testCategories: { + if (docker) { + includeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' + return + } + includeCategories 'org.apache.beam.sdk.testing.NeedsRunner' includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' // TODO: BEAM-12350 excludeCategories 'org.apache.beam.sdk.testing.UsesAttemptedMetrics' // TODO: BEAM-12681 @@ -173,7 +181,16 @@ createPortableValidatesRunnerTask( // TODO(BEAM-13498) excludeTestsMatching 'org.apache.beam.sdk.transforms.ParDoTest$TimestampTests.testProcessElementSkew' } -) + ) +} + +project.ext.validatesPortableRunnerDocker = portableValidatesRunnerTask("Docker", true) +project.ext.validatesPortableRunnerEmbedded = portableValidatesRunnerTask("Embedded", false) + +task validatesPortableRunner() { + dependsOn validatesPortableRunnerDocker + dependsOn validatesPortableRunnerEmbedded +} def jobPort = BeamModulePlugin.getRandomPort() def artifactPort = BeamModulePlugin.getRandomPort() diff --git a/runners/spark/job-server/spark_job_server.gradle b/runners/spark/job-server/spark_job_server.gradle index ef23855217e1..19c37f3bce7f 100644 --- a/runners/spark/job-server/spark_job_server.gradle +++ b/runners/spark/job-server/spark_job_server.gradle @@ -86,98 +86,109 @@ runShadow { jvmArgs += ["-Dorg.slf4j.simpleLogger.defaultLogLevel=${project.property('logLevel')}"] } -def portableValidatesRunnerTask(String name, Boolean streaming) { +def portableValidatesRunnerTask(String name, boolean streaming, boolean docker) { def pipelineOptions = [] def testCategories def testFilter - if (streaming) { - pipelineOptions += "--streaming" - pipelineOptions += "--streamingTimeoutMs=30000" - - testCategories = { - includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' - excludeCategories 'org.apache.beam.sdk.testing.FlattenWithHeterogeneousCoders' - excludeCategories 'org.apache.beam.sdk.testing.LargeKeys$Above100MB' - excludeCategories 'org.apache.beam.sdk.testing.UsesCommittedMetrics' - excludeCategories 'org.apache.beam.sdk.testing.UsesCustomWindowMerging' - excludeCategories 'org.apache.beam.sdk.testing.UsesFailureMessage' - excludeCategories 'org.apache.beam.sdk.testing.UsesGaugeMetrics' - excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderedDelivery' - excludeCategories 'org.apache.beam.sdk.testing.UsesParDoLifecycle' - excludeCategories 'org.apache.beam.sdk.testing.UsesMapState' - excludeCategories 'org.apache.beam.sdk.testing.UsesSetState' - excludeCategories 'org.apache.beam.sdk.testing.UsesOrderedListState' - excludeCategories 'org.apache.beam.sdk.testing.UsesTimerMap' - excludeCategories 'org.apache.beam.sdk.testing.UsesKeyInParDo' - excludeCategories 'org.apache.beam.sdk.testing.UsesOnWindowExpiration' - excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' - // TODO (BEAM-7222) SplittableDoFnTests - excludeCategories 'org.apache.beam.sdk.testing.UsesBoundedSplittableParDo' - excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo' - excludeCategories 'org.apache.beam.sdk.testing.UsesStrictTimerOrdering' - excludeCategories 'org.apache.beam.sdk.testing.UsesBundleFinalizer' - // Currently unsupported in portable streaming: - // TODO (BEAM-10712) - excludeCategories 'org.apache.beam.sdk.testing.UsesSideInputs' - // TODO (BEAM-10754) - excludeCategories 'org.apache.beam.sdk.testing.UsesStatefulParDo' - // TODO (BEAM-10755) - excludeCategories 'org.apache.beam.sdk.testing.UsesTimersInParDo' - } + if (docker) { + // Run the limited set of tests that need to validate the environment + // that contains the SDK is configured properly. + testCategories = { includeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' } + testFilter = { } + } else { + if (streaming) { + pipelineOptions += "--streaming" + pipelineOptions += "--streamingTimeoutMs=30000" - testFilter = { - // TODO (BEAM-10094) - excludeTestsMatching 'org.apache.beam.sdk.transforms.FlattenTest.testFlattenWithDifferentInputAndOutputCoders2' - // TODO (BEAM-10784) Currently unsupported in portable streaming: - // // Timeout error - excludeTestsMatching 'org.apache.beam.sdk.testing.PAssertTest.testWindowedContainsInAnyOrder' - excludeTestsMatching 'org.apache.beam.sdk.testing.PAssertTest.testWindowedSerializablePredicate' - excludeTestsMatching 'org.apache.beam.sdk.transforms.windowing.WindowTest.testNoWindowFnDoesNotReassignWindows' - // // Assertion error: empty iterable output - excludeTestsMatching 'org.apache.beam.sdk.transforms.CombineTest$WindowingTests.testFixedWindowsCombine' - excludeTestsMatching 'org.apache.beam.sdk.transforms.CombineTest$WindowingTests.testSessionsCombine' - excludeTestsMatching 'org.apache.beam.sdk.transforms.GroupByKeyTest$WindowTests' - excludeTestsMatching 'org.apache.beam.sdk.transforms.ReshuffleTest.testReshuffleAfterFixedWindowsAndGroupByKey' - excludeTestsMatching 'org.apache.beam.sdk.transforms.ReshuffleTest.testReshuffleAfterSessionsAndGroupByKey' - excludeTestsMatching 'org.apache.beam.sdk.transforms.ReshuffleTest.testReshuffleAfterSlidingWindowsAndGroupByKey' - excludeTestsMatching 'org.apache.beam.sdk.transforms.join.CoGroupByKeyTest.testCoGroupByKeyWithWindowing' - excludeTestsMatching 'org.apache.beam.sdk.transforms.windowing.WindowingTest' - // // Assertion error: incorrect output - excludeTestsMatching 'CombineTest$BasicTests.testHotKeyCombining' - } - } - else { - testCategories = { - includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' - excludeCategories 'org.apache.beam.sdk.testing.FlattenWithHeterogeneousCoders' - excludeCategories 'org.apache.beam.sdk.testing.LargeKeys$Above100MB' - excludeCategories 'org.apache.beam.sdk.testing.UsesCommittedMetrics' - excludeCategories 'org.apache.beam.sdk.testing.UsesCustomWindowMerging' - excludeCategories 'org.apache.beam.sdk.testing.UsesFailureMessage' - excludeCategories 'org.apache.beam.sdk.testing.UsesGaugeMetrics' - excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderedDelivery' - excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderInBundle' - excludeCategories 'org.apache.beam.sdk.testing.UsesParDoLifecycle' - excludeCategories 'org.apache.beam.sdk.testing.UsesMapState' - excludeCategories 'org.apache.beam.sdk.testing.UsesSetState' - excludeCategories 'org.apache.beam.sdk.testing.UsesOrderedListState' - excludeCategories 'org.apache.beam.sdk.testing.UsesTimerMap' - excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedPCollections' - excludeCategories 'org.apache.beam.sdk.testing.UsesKeyInParDo' - excludeCategories 'org.apache.beam.sdk.testing.UsesOnWindowExpiration' - excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' - // TODO (BEAM-7222) SplittableDoFnTests - excludeCategories 'org.apache.beam.sdk.testing.UsesBoundedSplittableParDo' - excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo' - excludeCategories 'org.apache.beam.sdk.testing.UsesStrictTimerOrdering' - excludeCategories 'org.apache.beam.sdk.testing.UsesBundleFinalizer' + testCategories = { + includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' + excludeCategories 'org.apache.beam.sdk.testing.FlattenWithHeterogeneousCoders' + excludeCategories 'org.apache.beam.sdk.testing.LargeKeys$Above100MB' + excludeCategories 'org.apache.beam.sdk.testing.UsesCommittedMetrics' + excludeCategories 'org.apache.beam.sdk.testing.UsesCustomWindowMerging' + excludeCategories 'org.apache.beam.sdk.testing.UsesFailureMessage' + excludeCategories 'org.apache.beam.sdk.testing.UsesGaugeMetrics' + excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderedDelivery' + excludeCategories 'org.apache.beam.sdk.testing.UsesParDoLifecycle' + excludeCategories 'org.apache.beam.sdk.testing.UsesMapState' + excludeCategories 'org.apache.beam.sdk.testing.UsesSetState' + excludeCategories 'org.apache.beam.sdk.testing.UsesOrderedListState' + excludeCategories 'org.apache.beam.sdk.testing.UsesTimerMap' + excludeCategories 'org.apache.beam.sdk.testing.UsesKeyInParDo' + excludeCategories 'org.apache.beam.sdk.testing.UsesOnWindowExpiration' + excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' + // TODO (BEAM-7222) SplittableDoFnTests + excludeCategories 'org.apache.beam.sdk.testing.UsesBoundedSplittableParDo' + excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo' + excludeCategories 'org.apache.beam.sdk.testing.UsesStrictTimerOrdering' + excludeCategories 'org.apache.beam.sdk.testing.UsesBundleFinalizer' + // Currently unsupported in portable streaming: + // TODO (BEAM-10712) + excludeCategories 'org.apache.beam.sdk.testing.UsesSideInputs' + // TODO (BEAM-10754) + excludeCategories 'org.apache.beam.sdk.testing.UsesStatefulParDo' + // TODO (BEAM-10755) + excludeCategories 'org.apache.beam.sdk.testing.UsesTimersInParDo' + } + + testFilter = { + // TODO (BEAM-10094) + excludeTestsMatching 'org.apache.beam.sdk.transforms.FlattenTest.testFlattenWithDifferentInputAndOutputCoders2' + // TODO (BEAM-10784) Currently unsupported in portable streaming: + // // Timeout error + excludeTestsMatching 'org.apache.beam.sdk.testing.PAssertTest.testWindowedContainsInAnyOrder' + excludeTestsMatching 'org.apache.beam.sdk.testing.PAssertTest.testWindowedSerializablePredicate' + excludeTestsMatching 'org.apache.beam.sdk.transforms.windowing.WindowTest.testNoWindowFnDoesNotReassignWindows' + // // Assertion error: empty iterable output + excludeTestsMatching 'org.apache.beam.sdk.transforms.CombineTest$WindowingTests.testFixedWindowsCombine' + excludeTestsMatching 'org.apache.beam.sdk.transforms.CombineTest$WindowingTests.testSessionsCombine' + excludeTestsMatching 'org.apache.beam.sdk.transforms.GroupByKeyTest$WindowTests' + excludeTestsMatching 'org.apache.beam.sdk.transforms.ReshuffleTest.testReshuffleAfterFixedWindowsAndGroupByKey' + excludeTestsMatching 'org.apache.beam.sdk.transforms.ReshuffleTest.testReshuffleAfterSessionsAndGroupByKey' + excludeTestsMatching 'org.apache.beam.sdk.transforms.ReshuffleTest.testReshuffleAfterSlidingWindowsAndGroupByKey' + excludeTestsMatching 'org.apache.beam.sdk.transforms.join.CoGroupByKeyTest.testCoGroupByKeyWithWindowing' + excludeTestsMatching 'org.apache.beam.sdk.transforms.windowing.WindowingTest' + // // Assertion error: incorrect output + excludeTestsMatching 'CombineTest$BasicTests.testHotKeyCombining' + } } - testFilter = { - // TODO (BEAM-10094) - excludeTestsMatching 'org.apache.beam.sdk.transforms.FlattenTest.testFlattenWithDifferentInputAndOutputCoders2' - // TODO(BEAM-13498) - excludeTestsMatching 'org.apache.beam.sdk.transforms.ParDoTest$TimestampTests.testProcessElementSkew' + else { + testCategories = { + includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' + excludeCategories 'org.apache.beam.sdk.testing.FlattenWithHeterogeneousCoders' + excludeCategories 'org.apache.beam.sdk.testing.LargeKeys$Above100MB' + excludeCategories 'org.apache.beam.sdk.testing.UsesCommittedMetrics' + excludeCategories 'org.apache.beam.sdk.testing.UsesCustomWindowMerging' + excludeCategories 'org.apache.beam.sdk.testing.UsesFailureMessage' + excludeCategories 'org.apache.beam.sdk.testing.UsesGaugeMetrics' + excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderedDelivery' + excludeCategories 'org.apache.beam.sdk.testing.UsesPerKeyOrderInBundle' + excludeCategories 'org.apache.beam.sdk.testing.UsesParDoLifecycle' + excludeCategories 'org.apache.beam.sdk.testing.UsesMapState' + excludeCategories 'org.apache.beam.sdk.testing.UsesSetState' + excludeCategories 'org.apache.beam.sdk.testing.UsesOrderedListState' + excludeCategories 'org.apache.beam.sdk.testing.UsesTimerMap' + excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedPCollections' + excludeCategories 'org.apache.beam.sdk.testing.UsesKeyInParDo' + excludeCategories 'org.apache.beam.sdk.testing.UsesOnWindowExpiration' + excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' + // TODO (BEAM-7222) SplittableDoFnTests + excludeCategories 'org.apache.beam.sdk.testing.UsesBoundedSplittableParDo' + excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo' + excludeCategories 'org.apache.beam.sdk.testing.UsesStrictTimerOrdering' + excludeCategories 'org.apache.beam.sdk.testing.UsesBundleFinalizer' + } + testFilter = { + // TODO (BEAM-10094) + excludeTestsMatching 'org.apache.beam.sdk.transforms.FlattenTest.testFlattenWithDifferentInputAndOutputCoders2' + // TODO(BEAM-13498) + excludeTestsMatching 'org.apache.beam.sdk.transforms.ParDoTest$TimestampTests.testProcessElementSkew' + } } } @@ -188,7 +199,7 @@ def portableValidatesRunnerTask(String name, Boolean streaming) { testClasspathConfiguration: configurations.validatesPortableRunner, numParallelTests: 4, pipelineOpts: pipelineOptions, - environment: BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.EMBEDDED, + environment: docker ? BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.DOCKER : BeamModulePlugin.PortableValidatesRunnerConfiguration.Environment.EMBEDDED, systemProperties: [ "beam.spark.test.reuseSparkContext": "false", "spark.ui.enabled": "false", @@ -199,10 +210,12 @@ def portableValidatesRunnerTask(String name, Boolean streaming) { ) } -project.ext.validatesPortableRunnerBatch = portableValidatesRunnerTask("Batch", false) -project.ext.validatesPortableRunnerStreaming = portableValidatesRunnerTask("Streaming", true) +project.ext.validatesPortableRunnerDocker= portableValidatesRunnerTask("Docker", false, true) +project.ext.validatesPortableRunnerBatch = portableValidatesRunnerTask("Batch", false, false) +project.ext.validatesPortableRunnerStreaming = portableValidatesRunnerTask("Streaming", true, false) task validatesPortableRunner() { + dependsOn validatesPortableRunnerDocker dependsOn validatesPortableRunnerBatch dependsOn validatesPortableRunnerStreaming } diff --git a/runners/spark/spark_runner.gradle b/runners/spark/spark_runner.gradle index 5a84b93f55f1..43f8b4739db1 100644 --- a/runners/spark/spark_runner.gradle +++ b/runners/spark/spark_runner.gradle @@ -265,6 +265,8 @@ task validatesRunnerBatch(type: Test) { useJUnit { includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' includeCategories 'org.apache.beam.runners.spark.UsesCheckpointRecovery' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories 'org.apache.beam.sdk.testing.UsesCustomWindowMerging' excludeCategories 'org.apache.beam.sdk.testing.UsesTimerMap' excludeCategories 'org.apache.beam.sdk.testing.UsesOnWindowExpiration' @@ -339,6 +341,8 @@ task validatesStructuredStreamingRunnerBatch(type: Test) { jvmArgs '-Xmx7g' useJUnit { includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' // Unbounded excludeCategories 'org.apache.beam.sdk.testing.UsesUnboundedPCollections' excludeCategories 'org.apache.beam.sdk.testing.UsesTestStream' diff --git a/runners/twister2/build.gradle b/runners/twister2/build.gradle index de7b62e37224..25e177f8f486 100644 --- a/runners/twister2/build.gradle +++ b/runners/twister2/build.gradle @@ -81,6 +81,8 @@ task validatesRunnerBatch(type: Test) { forkEvery 1 useJUnit { includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner' + // Should be run only in a properly configured SDK harness environment + excludeCategories 'org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment' excludeCategories 'org.apache.beam.sdk.testing.FlattenWithHeterogeneousCoders' excludeCategories 'org.apache.beam.sdk.testing.UsesStatefulParDo' excludeCategories 'org.apache.beam.sdk.testing.UsesTimersInParDo' diff --git a/sdks/java/container/Dockerfile b/sdks/java/container/Dockerfile index b5baf8fa7d4f..6e12ff1dc5c1 100644 --- a/sdks/java/container/Dockerfile +++ b/sdks/java/container/Dockerfile @@ -30,6 +30,9 @@ ADD target/beam-sdks-java-harness.jar /opt/apache/beam/jars/ ADD target/beam-sdks-java-io-kafka.jar /opt/apache/beam/jars/ ADD target/kafka-clients.jar /opt/apache/beam/jars/ +# Required to use jamm as a javaagent to get accurate object size measuring +ADD target/jamm.jar /opt/apache/beam/jars/ + ADD target/linux_amd64/boot /opt/apache/beam/ COPY target/LICENSE /opt/apache/beam/ diff --git a/sdks/java/container/boot.go b/sdks/java/container/boot.go index 5d120fc77729..14f18c568269 100644 --- a/sdks/java/container/boot.go +++ b/sdks/java/container/boot.go @@ -49,10 +49,12 @@ var ( ) const ( + disableJammAgentOption = "disable_jamm_agent" enableGoogleCloudProfilerOption = "enable_google_cloud_profiler" enableGoogleCloudHeapSamplingOption = "enable_google_cloud_heap_sampling" googleCloudProfilerAgentBaseArgs = "-agentpath:/opt/google_cloud_profiler/profiler_java_agent.so=-logtostderr,-cprof_service=%s,-cprof_service_version=%s" googleCloudProfilerAgentHeapArgs = googleCloudProfilerAgentBaseArgs + ",-cprof_enable_heap_sampling,-cprof_heap_sampling_interval=2097152" + jammAgentArgs = "-javaagent:/opt/apache/beam/jars/jamm.jar" ) func main() { @@ -185,6 +187,13 @@ func main() { } } + disableJammAgent := strings.Contains(options, disableJammAgentOption) + if disableJammAgent { + log.Printf("Disabling Jamm agent. Measuring object size will be inaccurate.") + } else { + args = append(args, jammAgentArgs) + } + args = append(args, "org.apache.beam.fn.harness.FnHarness") log.Printf("Executing: java %v", strings.Join(args, " ")) diff --git a/sdks/java/container/build.gradle b/sdks/java/container/build.gradle index 142e4aaefe08..2314c9153eba 100644 --- a/sdks/java/container/build.gradle +++ b/sdks/java/container/build.gradle @@ -42,6 +42,7 @@ dependencies { dockerDependency project(":sdks:java:io:kafka") // This dependency is set to 'compileOnly' scope in :sdks:java:io:kafka dockerDependency library.java.kafka_clients + dockerDependency library.java.jamm } goBuild { diff --git a/sdks/java/container/common.gradle b/sdks/java/container/common.gradle index 6f828ece9259..f48172e52c73 100644 --- a/sdks/java/container/common.gradle +++ b/sdks/java/container/common.gradle @@ -51,6 +51,7 @@ task copyDockerfileDependencies(type: Copy) { rename 'beam-sdks-java-harness-.*.jar', 'beam-sdks-java-harness.jar' rename 'beam-sdks-java-io-kafka.*.jar', 'beam-sdks-java-io-kafka.jar' rename 'kafka-clients.*.jar', 'kafka-clients.jar' + rename "jamm.*.jar", "jamm.jar" into "build/target" } diff --git a/sdks/java/core/build.gradle b/sdks/java/core/build.gradle index 9939bf9c2a38..a34e5e1435f6 100644 --- a/sdks/java/core/build.gradle +++ b/sdks/java/core/build.gradle @@ -103,6 +103,7 @@ dependencies { shadowTest library.java.quickcheck_core shadowTest library.java.avro_tests shadowTest library.java.zstd_jni + shadowTest library.java.jamm testRuntimeOnly library.java.slf4j_jdk14 testImplementation 'io.airlift:aircompressor:0.18' testImplementation 'com.facebook.presto.hadoop:hadoop-apache2:3.2.0-1' diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/UsesSdkHarnessEnvironment.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/UsesSdkHarnessEnvironment.java new file mode 100644 index 000000000000..d459e32980bd --- /dev/null +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/UsesSdkHarnessEnvironment.java @@ -0,0 +1,26 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk.testing; + +import org.apache.beam.sdk.annotations.Internal; + +/** + * Category tag for tests which validate that the SDK harness executes in a well formed environment. + */ +@Internal +public interface UsesSdkHarnessEnvironment {} diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/SdkHarnessEnvironmentTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/SdkHarnessEnvironmentTest.java new file mode 100644 index 000000000000..dd2d469fd4be --- /dev/null +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/SdkHarnessEnvironmentTest.java @@ -0,0 +1,69 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.greaterThan; + +import org.apache.beam.sdk.coders.StringUtf8Coder; +import org.apache.beam.sdk.testing.PAssert; +import org.apache.beam.sdk.testing.TestPipeline; +import org.apache.beam.sdk.testing.UsesSdkHarnessEnvironment; +import org.apache.beam.sdk.testing.ValidatesRunner; +import org.apache.beam.sdk.transforms.Create; +import org.apache.beam.sdk.transforms.DoFn; +import org.apache.beam.sdk.transforms.ParDo; +import org.apache.beam.sdk.values.PCollection; +import org.github.jamm.MemoryMeter; +import org.github.jamm.MemoryMeter.Guess; +import org.junit.Rule; +import org.junit.Test; +import org.junit.experimental.categories.Category; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +/** Tests that validate the SDK harness is configured correctly for a runner. */ +@RunWith(JUnit4.class) +public class SdkHarnessEnvironmentTest { + + @Rule public final TestPipeline p = TestPipeline.create(); + + /** + * {@link DoFn} used to validate that Jamm was setup as a java agent to get accurate measuring. + */ + private static class JammDoFn extends DoFn { + @ProcessElement + public void processElement(ProcessContext c) { + MemoryMeter memoryMeter = + MemoryMeter.builder().withGuessing(Guess.ALWAYS_INSTRUMENTATION).build(); + assertThat(memoryMeter.measureDeep(c.element()), greaterThan(0L)); + c.output("measured"); + } + } + + @Test + @Category({ValidatesRunner.class, UsesSdkHarnessEnvironment.class}) + public void testJammAgentAvailable() throws Exception { + PCollection input = p.apply(Create.of("jamm").withCoder(StringUtf8Coder.of())); + + PCollection output = input.apply(ParDo.of(new JammDoFn())); + + PAssert.that(output).containsInAnyOrder("measured"); + p.run().waitUntilFinish(); + } +} diff --git a/sdks/java/harness/build.gradle b/sdks/java/harness/build.gradle index e0559be4fb8c..04463356c66d 100644 --- a/sdks/java/harness/build.gradle +++ b/sdks/java/harness/build.gradle @@ -68,11 +68,7 @@ dependencies { implementation library.java.joda_time implementation library.java.slf4j_api implementation library.java.vendored_grpc_1_36_0 - - // Swap to use the officially published version of 0.4.x once available - // instead of relying on a community published copy. See - // https://github.com/jbellis/jamm/issues/44 for additional details. - implementation 'io.github.stephankoelle:jamm:0.4.1' + implementation library.java.jamm testImplementation library.java.junit testImplementation library.java.mockito_core testImplementation project(path: ":sdks:java:core", configuration: "shadowTest") diff --git a/sdks/java/io/hadoop-format/build.gradle b/sdks/java/io/hadoop-format/build.gradle index 1406893bc82e..69ce5de768ff 100644 --- a/sdks/java/io/hadoop-format/build.gradle +++ b/sdks/java/io/hadoop-format/build.gradle @@ -65,6 +65,9 @@ dependencies { permitUnusedDeclared library.java.hadoop_hdfs compileOnly library.java.hadoop_hdfs_client compileOnly library.java.hadoop_mapreduce_client_core + // Ensure that the older version of JAMM that cassandra relies on appears + // on the classpath before the one provided by :sdks:java:core shadowTest. + testImplementation "com.github.jbellis:jamm:0.3.0" testImplementation project(path: ":sdks:java:core", configuration: "shadowTest") testImplementation project(path: ":sdks:java:io:common", configuration: "testRuntimeMigration") testImplementation project(":sdks:java:testing:test-utils")