Skip to content

Commit 983a0be

Browse files
aboitreauddevflow.devflow-routing-intake
andauthored
Add Spark 4.0.0 test suite on JDK 17 to spark_2.13 (#12049)
Add Spark 4.0.0 test suite on JDK 17 to spark_2.13 Mirror Spark's full JDK 17 module opens in test_spark4 Co-authored-by: devflow.devflow-routing-intake <devflow.devflow-routing-intake@kubernetes.us1.ddbuild.io>
1 parent 731b447 commit 983a0be

3 files changed

Lines changed: 85 additions & 17 deletions

File tree

dd-java-agent/instrumentation/spark/spark-common/src/testFixtures/groovy/datadog/trace/instrumentation/spark/AbstractSparkListenerTest.groovy

Lines changed: 21 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,15 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
6060
)
6161
}
6262

63+
protected applicationEndEvent(Long time) {
64+
// Constructor of SparkListenerApplicationEnd gained an exitCode parameter starting spark 4.0
65+
if (TestSparkComputation.getSparkVersion() < "4") {
66+
return new SparkListenerApplicationEnd(time)
67+
}
68+
69+
return new SparkListenerApplicationEnd(time, Option.empty())
70+
}
71+
6372
protected jobStartEvent(Integer jobId, Long time, ArrayList<Integer> stageIds) {
6473
def stageInfos = stageIds.collect { stageId ->
6574
createStageInfo(stageId)
@@ -214,7 +223,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
214223
listener.onExecutorAdded(executorAddedEvent(3000L, "executor-2", 5))
215224
listener.onExecutorRemoved(executorRemovedEvent(4000L, "executor-2"))
216225

217-
listener.onApplicationEnd(new SparkListenerApplicationEnd(5000L))
226+
listener.onApplicationEnd(applicationEndEvent(5000L))
218227

219228
expect:
220229
def expectedExecutorTime = (5000 - 2000) * 4 + (4000 - 3000) * 5
@@ -256,7 +265,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
256265
listener.onStageCompleted(stageCompletedEvent(3, 3100L))
257266

258267
listener.onJobEnd(jobEndEvent(1, 3100L))
259-
listener.onApplicationEnd(new SparkListenerApplicationEnd(3100L))
268+
listener.onApplicationEnd(applicationEndEvent(3100L))
260269

261270
expect:
262271
/*
@@ -339,7 +348,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
339348
listener.onStageCompleted(stageCompletedEvent(2, 3000L))
340349

341350
listener.onJobEnd(jobEndEvent(1, 3000L))
342-
listener.onApplicationEnd(new SparkListenerApplicationEnd(3000L))
351+
listener.onApplicationEnd(applicationEndEvent(3000L))
343352

344353
expect:
345354
assertTraces(1) {
@@ -384,7 +393,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
384393
listener.onTaskEnd(taskEndEvent(1,1900L, 300))
385394
listener.onStageCompleted(stageCompletedEvent(1, 2200L))
386395
listener.onJobEnd(jobEndEvent(1, 17200L))
387-
listener.onApplicationEnd(new SparkListenerApplicationEnd(3100L))
396+
listener.onApplicationEnd(applicationEndEvent(3100L))
388397

389398
assertTraces(1) {
390399
trace(3) {
@@ -430,7 +439,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
430439

431440
listener.onStageCompleted(stageCompletedEvent(2, 17200L))
432441
listener.onJobEnd(jobEndEvent(1, 17200L))
433-
listener.onApplicationEnd(new SparkListenerApplicationEnd(3100L))
442+
listener.onApplicationEnd(applicationEndEvent(3100L))
434443

435444
expect:
436445
def relativeAccuracy = 1/32D
@@ -480,7 +489,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
480489

481490
when:
482491
listener.onApplicationStart(applicationStartEvent(1000L))
483-
listener.onApplicationEnd(new SparkListenerApplicationEnd(2000L))
492+
listener.onApplicationEnd(applicationEndEvent(2000L))
484493

485494
then:
486495
assertTraces(1) {
@@ -500,7 +509,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
500509
listener.onApplicationStart(applicationStartEvent(1000L))
501510
listener.onJobStart(jobStartEvent(1, 1900L, [1]))
502511
listener.onJobEnd(jobFailedEvent(1, 2200L, "Job was cancelled by user"))
503-
listener.onApplicationEnd(new SparkListenerApplicationEnd(2300L))
512+
listener.onApplicationEnd(applicationEndEvent(2300L))
504513

505514
expect:
506515
assertTraces(1) {
@@ -531,7 +540,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
531540
def analysisException = new RuntimeException("[TABLE_OR_VIEW_NOT_FOUND] The table or view `missing_table` cannot be found.")
532541
listener.onSqlFailure(analysisException)
533542

534-
listener.onApplicationEnd(new SparkListenerApplicationEnd(2000L))
543+
listener.onApplicationEnd(applicationEndEvent(2000L))
535544

536545
expect:
537546
assertTraces(1) {
@@ -563,7 +572,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
563572
listener.onStageCompleted(stageCompletedEvent(1, 1800L))
564573
listener.onJobEnd(jobFailedEvent(1, 2000L, "Job aborted due to NullPointerException"))
565574

566-
listener.onApplicationEnd(new SparkListenerApplicationEnd(3000L))
575+
listener.onApplicationEnd(applicationEndEvent(3000L))
567576

568577
expect:
569578
assertTraces(1) {
@@ -663,7 +672,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
663672
injectSysConfig("dd.tags", ddTags)
664673
def listener = getTestDatadogSparkListener()
665674
listener.onApplicationStart(applicationStartEvent(1000L))
666-
listener.onApplicationEnd(new SparkListenerApplicationEnd(5000L))
675+
listener.onApplicationEnd(applicationEndEvent(5000L))
667676

668677
expect:
669678
assertTraces(1) {
@@ -729,7 +738,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
729738

730739
when:
731740
listener.onApplicationStart(applicationStartEvent(1000L))
732-
listener.onApplicationEnd(new SparkListenerApplicationEnd(2000L))
741+
listener.onApplicationEnd(applicationEndEvent(2000L))
733742

734743
then:
735744
assertTraces(1) {
@@ -749,7 +758,7 @@ abstract class AbstractSparkListenerTest extends InstrumentationSpecification {
749758

750759
when:
751760
listener.onApplicationStart(applicationStartEvent(1000L))
752-
listener.onApplicationEnd(new SparkListenerApplicationEnd(2000L))
761+
listener.onApplicationEnd(applicationEndEvent(2000L))
753762

754763
then:
755764
assertTraces(1) {

dd-java-agent/instrumentation/spark/spark-common/src/testFixtures/groovy/datadog/trace/instrumentation/spark/AbstractSparkTest.groovy

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -157,7 +157,8 @@ abstract class AbstractSparkTest extends InstrumentationSpecification {
157157
errored true
158158
parent()
159159
assert span.tags["error.type"] == "Spark Application Failed"
160-
assert span.tags["error.message"] =~ /^Job aborted due to stage failure.*java.lang.NullPointerException$/
160+
// JDK 15+'s helpful NullPointerException messages (JEP 358) append extra detail after the exception class name.
161+
assert span.tags["error.message"] =~ /^Job aborted due to stage failure.*java.lang.NullPointerException.*$/
161162
assert span.tags["error.stack"] =~ /(?s)^org.apache.spark.SparkException.*Caused by: java.lang.NullPointerException.*$/
162163
}
163164
span {
@@ -167,7 +168,8 @@ abstract class AbstractSparkTest extends InstrumentationSpecification {
167168
errored true
168169
childOf(span(0))
169170
assert span.tags["error.type"] == "Spark Job Failed"
170-
assert span.tags["error.message"] =~ /^Job aborted due to stage failure.*java.lang.NullPointerException$/
171+
// JDK 15+'s helpful NullPointerException messages (JEP 358) append extra detail after the exception class name.
172+
assert span.tags["error.message"] =~ /^Job aborted due to stage failure.*java.lang.NullPointerException.*$/
171173
assert span.tags["error.stack"] =~ /(?s)^org.apache.spark.SparkException.*Caused by: java.lang.NullPointerException.*$/
172174
}
173175
span {
@@ -177,7 +179,8 @@ abstract class AbstractSparkTest extends InstrumentationSpecification {
177179
errored true
178180
childOf(span(1))
179181
assert span.tags["error.type"] == "Spark Stage Failed"
180-
assert span.tags["error.message"] =~ /^Job aborted due to stage failure.*java.lang.NullPointerException$/
182+
// JDK 15+'s helpful NullPointerException messages (JEP 358) append extra detail after the exception class name.
183+
assert span.tags["error.message"] =~ /^Job aborted due to stage failure.*java.lang.NullPointerException.*$/
181184
assert span.tags["error.stack"] =~ /(?s).*\n\tat datadog.trace.instrumentation.spark.TestSparkComputation.{500,}\$/
182185
}
183186
span {
@@ -186,8 +189,9 @@ abstract class AbstractSparkTest extends InstrumentationSpecification {
186189
errored true
187190
childOf(span(2))
188191
assert span.tags["error.type"] == "Spark Task Failed"
189-
assert span.tags["error.message"] == "java.lang.NullPointerException: null"
190-
assert span.tags["error.stack"] =~ /(?s)^java.lang.NullPointerException\n\tat datadog.trace.instrumentation.spark.TestSparkComputation.{500,}\$/
192+
// JDK 15+'s helpful NullPointerException messages (JEP 358) append extra detail after the exception class name.
193+
assert span.tags["error.message"] =~ /^java\.lang\.NullPointerException: .*$/
194+
assert span.tags["error.stack"] =~ /(?s)^java.lang.NullPointerException.*\n\tat datadog.trace.instrumentation.spark.TestSparkComputation.{500,}\$/
191195
}
192196
}
193197
}
@@ -829,6 +833,9 @@ abstract class AbstractSparkTest extends InstrumentationSpecification {
829833
def sparkSession = SparkSession.builder()
830834
.config("spark.master", "local[2]")
831835
.config("spark.sql.shuffle.partitions", "2")
836+
// Under ANSI mode (default since Spark 4.0), casting the non-numeric "value" column
837+
// content to a number for the filter predicate below throws instead of yielding null.
838+
.config("spark.sql.ansi.enabled", "false")
832839
.getOrCreate()
833840

834841
def df = generateSampleDataframe(sparkSession)

dd-java-agent/instrumentation/spark/spark_2.13/build.gradle

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ apply from: "$rootDir/gradle/java.gradle"
2424

2525
addTestSuiteForDir('latestDepTest', 'test')
2626
addTestSuite('test_spark32')
27+
addTestSuiteForDir('test_spark4', 'test')
2728

2829
testJvmConstraints {
2930
// Hadoop does not behave correctly with OpenJ9 https://issues.apache.org/jira/browse/HADOOP-18174
@@ -53,6 +54,10 @@ dependencies {
5354
test_spark32Implementation group: 'org.apache.spark', name: "spark-sql_$scalaVersion", version: "3.2.4"
5455
test_spark32Implementation group: 'org.apache.spark', name: "spark-yarn_$scalaVersion", version: "3.2.4"
5556

57+
test_spark4Implementation group: 'org.apache.spark', name: "spark-core_$scalaVersion", version: "4.0.0"
58+
test_spark4Implementation group: 'org.apache.spark', name: "spark-sql_$scalaVersion", version: "4.0.0"
59+
test_spark4Implementation group: 'org.apache.spark', name: "spark-yarn_$scalaVersion", version: "4.0.0"
60+
5661
// FIXME: Currently not working on Spark 4.0.0 preview releases.
5762
latestDepTestImplementation group: 'org.apache.spark', name: "spark-core_$scalaVersion", version: '3.+'
5863
latestDepTestImplementation group: 'org.apache.spark', name: "spark-sql_$scalaVersion", version: '3.+'
@@ -61,4 +66,51 @@ dependencies {
6166

6267
tasks.named("test", Test) {
6368
dependsOn "test_spark32"
69+
dependsOn "test_spark4"
70+
}
71+
72+
tasks.named("test_spark4", Test) {
73+
testJvmConstraints {
74+
// Override the module-wide maxJavaVersion = 11 cap (set above for the Spark 3.2 test JVM
75+
// constraints) since Spark 4.0.0 requires JDK 17+.
76+
minJavaVersion = JavaVersion.VERSION_17
77+
// Hadoop's UserGroupInformation calls javax.security.auth.Subject.getSubject(), which
78+
// throws UnsupportedOperationException once the SecurityManager is fully removed in JDK 24+
79+
// (JEP 486).
80+
maxJavaVersion = JavaVersion.VERSION_23
81+
}
82+
83+
// Spark itself only opens these modules when launched through its own scripts
84+
// (org.apache.spark.launcher.JavaModuleOptions); since this suite starts a SparkSession
85+
// directly in-process, mirror that same list so Spark's own reflective access (Kryo shuffle
86+
// serializers, Tungsten unsafe memory, Hadoop UGI, etc.) doesn't fail under JDK 16+ module encapsulation.
87+
conditionalJvmArgs(
88+
it,
89+
JavaVersion.VERSION_16,
90+
[
91+
"--add-opens=java.base/java.lang=ALL-UNNAMED",
92+
"--add-opens=java.base/java.lang.invoke=ALL-UNNAMED",
93+
"--add-opens=java.base/java.lang.reflect=ALL-UNNAMED",
94+
"--add-opens=java.base/java.io=ALL-UNNAMED",
95+
"--add-opens=java.base/java.net=ALL-UNNAMED",
96+
"--add-opens=java.base/java.nio=ALL-UNNAMED",
97+
"--add-opens=java.base/java.util=ALL-UNNAMED",
98+
"--add-opens=java.base/java.util.concurrent=ALL-UNNAMED",
99+
"--add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED",
100+
"--add-opens=java.base/jdk.internal.ref=ALL-UNNAMED",
101+
"--add-opens=java.base/sun.nio.ch=ALL-UNNAMED",
102+
"--add-opens=java.base/sun.nio.cs=ALL-UNNAMED",
103+
"--add-opens=java.base/sun.security.action=ALL-UNNAMED",
104+
"--add-opens=java.base/sun.util.calendar=ALL-UNNAMED",
105+
"--add-opens=java.security.jgss/sun.security.krb5=ALL-UNNAMED"
106+
]
107+
)
108+
}
109+
110+
// Spark 4.0.0 transitively pulls a newer protobuf-java that's binary-incompatible with the
111+
// pre-generated com.datadoghq.sketch.ddsketch.proto.DDSketch classes used by the test fixtures.
112+
['test_spark4CompileClasspath', 'test_spark4RuntimeClasspath'].each {
113+
configurations.named(it) {
114+
resolutionStrategy.force 'com.google.protobuf:protobuf-java:3.14.0'
115+
}
64116
}

0 commit comments

Comments
 (0)