Skip to content

Commit 35bba04

Browse files
committed
feat: support Spark HMS/Kerberos matrix and runtime configuration via Gradle
- Added support for running Spark tests with HMS/Kerberos matrix configurations. - Introduced Gradle tasks for Apache and Cloudera HMS variants with Kerberos on/off. - Enhanced `BigDataTest` to allow runtime configuration overrides via system properties. - Improved flexibility with additive and replaceable TOML configurations for testing. - Updated documentation and examples to showcase HMS/Kerberos matrix usage.
1 parent 80a41a4 commit 35bba04

14 files changed

Lines changed: 249 additions & 59 deletions

File tree

README.md

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -238,6 +238,18 @@ This binds ports such as HDFS `8020`/`9870`, HMS `9083`, Kafka `9092`, Schema Re
238238

239239
Available JUnit port fields are `kerberosKdcPort`, `hdfsNameNodePort`, `hdfsWebPort`, `hiveMetastorePort`, `kafkaPort`, `schemaRegistryPort`, `kafkaUiPort`, `localStackS3Port`, and `fakeGcsPort`. Endpoint properties still use the actual mapped host ports returned by Testcontainers.
240240

241+
Gradle tasks can drive startup config with `bigdata.test.config`. Use `bigdata.test.config.replace=true` for data-driven matrices where the task must choose mutually exclusive services such as `hiveMetastore` versus `clouderaHms`, or Kerberos enabled versus disabled.
242+
243+
The Spark example provides:
244+
245+
```bash
246+
./gradlew :example:spark:sparkApacheHmsTest
247+
./gradlew :example:spark:sparkApacheHmsKerberosTest
248+
./gradlew :example:spark:sparkClouderaHmsTest
249+
./gradlew :example:spark:sparkClouderaHmsKerberosTest
250+
./gradlew :example:spark:sparkBigDataMatrixTest
251+
```
252+
241253
Default service ports and endpoint property keys:
242254

243255
| Service | Port names | Default ports | Main endpoint properties |

doc/user-guide.adoc

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -222,6 +222,25 @@ class MyIntegrationTest
222222

223223
When both are used, files from `@BigDataExtensions` are loaded first and files from `@BigDataTest(config = ...)` are loaded after them, so direct `@BigDataTest` config wins for duplicate image keys.
224224

225+
Gradle or any other test launcher can add startup config at runtime with the `bigdata.test.config` system property. Use a comma-separated list when a task should compose common config with one matrix-specific override:
226+
227+
[source,kotlin]
228+
----
229+
tasks.register<Test>("sparkApacheHmsKerberosTest") {
230+
useJUnitPlatform()
231+
filter.includeTestsMatching("org.openprojectx.bigdata.test.example.spark.SparkBigDataTestExample")
232+
systemProperty("bigdata.test.config.replace", "true")
233+
systemProperty(
234+
"bigdata.test.config",
235+
"classpath:spark-bigdata-test-common.toml,classpath:spark-bigdata-test-apache-hms-kerberos.toml",
236+
)
237+
}
238+
----
239+
240+
When `bigdata.test.config.replace=true`, task config replaces startup config from `@BigDataExtensions` and `@BigDataTest(config = ...)`. Without it, task config is appended last and wins for scalar values. Service booleans are additive, so replacement mode is the right choice when a matrix needs to turn Kerberos or HMS implementations on and off.
241+
242+
Extension config has the same launcher hook: `bigdata.extensions.config` and `bigdata.extensions.config.replace`.
243+
225244
`hiveMetastore` is the open-source HMS image and defaults to `ghcr.io/openprojectx/hive:3.1.3-hadoop-3.4.2-gcs-4.0.4-jdk17-0.1.4`. Hive 4 images can be tested by overriding this image key, but Spark 3.x brings a Hive 2.3 metastore client and should use the Hive 3 image unless that client stack is changed. The open-source HMS path starts an external PostgreSQL support container using `hiveMetastorePostgres`.
226245

227246
`clouderaHms` is the embedded-Postgres Cloudera HMS image. Use the `clouderaHms = true` annotation flag when you want this implementation. A test class must not enable both `hiveMetastore` and `clouderaHms`.
@@ -778,6 +797,19 @@ Run it with:
778797
GRADLE_USER_HOME=/data/.gradle ./gradlew :example:spark:test
779798
----
780799

800+
The Spark example also has Gradle tasks for the HMS/Kerberos matrix:
801+
802+
[source,bash]
803+
----
804+
GRADLE_USER_HOME=/data/.gradle ./gradlew :example:spark:sparkApacheHmsTest
805+
GRADLE_USER_HOME=/data/.gradle ./gradlew :example:spark:sparkApacheHmsKerberosTest
806+
GRADLE_USER_HOME=/data/.gradle ./gradlew :example:spark:sparkClouderaHmsTest
807+
GRADLE_USER_HOME=/data/.gradle ./gradlew :example:spark:sparkClouderaHmsKerberosTest
808+
GRADLE_USER_HOME=/data/.gradle ./gradlew :example:spark:sparkBigDataMatrixTest
809+
----
810+
811+
The Kerberos axis currently enables the shared KDC and Kafka Kerberos path used by the Spark Kafka source. The HMS axis switches between the open-source Hive 3 HMS image and the Cloudera HMS image.
812+
781813
== Troubleshooting
782814

783815
=== `HADOOP_HOME and hadoop.home.dir are unset`

example/spark/build.gradle.kts

Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,3 +47,65 @@ tasks.withType<Test>().configureEach {
4747
"--add-opens=java.base/java.net=ALL-UNNAMED",
4848
)
4949
}
50+
51+
val sparkBigDataTestClass = "org.openprojectx.bigdata.test.example.spark.SparkBigDataTestExample"
52+
val sparkCommonConfig = "classpath:spark-bigdata-test-common.toml"
53+
54+
fun registerSparkMatrixTest(
55+
name: String,
56+
descriptionText: String,
57+
variantConfig: String,
58+
) = tasks.register<Test>(name) {
59+
description = descriptionText
60+
group = "verification"
61+
testClassesDirs = sourceSets.test.get().output.classesDirs
62+
classpath = sourceSets.test.get().runtimeClasspath
63+
useJUnitPlatform()
64+
filter.includeTestsMatching(sparkBigDataTestClass)
65+
systemProperty("bigdata.test.config.replace", "true")
66+
systemProperty("bigdata.test.config", "$sparkCommonConfig,$variantConfig")
67+
}
68+
69+
val sparkApacheHmsTest = registerSparkMatrixTest(
70+
name = "sparkApacheHmsTest",
71+
descriptionText = "Runs the Spark example with open-source Hive 3 HMS and plaintext Kafka.",
72+
variantConfig = "classpath:spark-bigdata-test-apache-hms.toml",
73+
)
74+
75+
val sparkApacheHmsKerberosTest = registerSparkMatrixTest(
76+
name = "sparkApacheHmsKerberosTest",
77+
descriptionText = "Runs the Spark example with open-source Hive 3 HMS and Kafka Kerberos.",
78+
variantConfig = "classpath:spark-bigdata-test-apache-hms-kerberos.toml",
79+
)
80+
sparkApacheHmsKerberosTest.configure {
81+
mustRunAfter(sparkApacheHmsTest)
82+
}
83+
84+
val sparkClouderaHmsTest = registerSparkMatrixTest(
85+
name = "sparkClouderaHmsTest",
86+
descriptionText = "Runs the Spark example with Cloudera HMS and plaintext Kafka.",
87+
variantConfig = "classpath:spark-bigdata-test-cloudera-hms.toml",
88+
)
89+
sparkClouderaHmsTest.configure {
90+
mustRunAfter(sparkApacheHmsKerberosTest)
91+
}
92+
93+
val sparkClouderaHmsKerberosTest = registerSparkMatrixTest(
94+
name = "sparkClouderaHmsKerberosTest",
95+
descriptionText = "Runs the Spark example with Cloudera HMS and Kafka Kerberos.",
96+
variantConfig = "classpath:spark-bigdata-test-cloudera-hms-kerberos.toml",
97+
)
98+
sparkClouderaHmsKerberosTest.configure {
99+
mustRunAfter(sparkClouderaHmsTest)
100+
}
101+
102+
tasks.register("sparkBigDataMatrixTest") {
103+
description = "Runs all Spark HMS/Kerberos matrix combinations."
104+
group = "verification"
105+
dependsOn(
106+
sparkApacheHmsTest,
107+
sparkApacheHmsKerberosTest,
108+
sparkClouderaHmsTest,
109+
sparkClouderaHmsKerberosTest,
110+
)
111+
}

example/spark/src/test/kotlin/org/openprojectx/bigdata/test/example/spark/SparkBigDataScenario.kt

Lines changed: 57 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,9 @@ abstract class SparkBigDataScenario {
4444
spark = context.spark,
4545
bootstrapServers = context.environment.kafkaBootstrapServers,
4646
topic = kafkaAvroTopic(),
47+
securityProtocol = context.environment.kafkaSecurityProtocol,
48+
kerberosServiceName = context.environment.kafkaKerberosServiceName,
49+
jaasConfig = context.environment.kafkaJaasConfig,
4750
)
4851
},
4952
)
@@ -132,17 +135,20 @@ abstract class SparkBigDataScenario {
132135
gcsIcebergDataPath = "${extensions.required("$gcsBucketExtensionId.gs.uri")}/data/demo_$runId/events_gcs",
133136
s3CredentialProviderPath = extensions.required("s3-jceks.credential-provider.path"),
134137
s3CredentialProviderHdfsPath = extensions.required("s3-jceks.hdfs.path"),
135-
kerberosClientPrincipal = extensions.required("kerberos-material.client.principal"),
136-
kerberosClientPassword = extensions.required("kerberos-material.client.password"),
137-
kerberosClientKeytab = extensions.required("kerberos-material.client.keytab"),
138-
krb5Conf = extensions.required("kerberos-material.krb5-conf"),
139-
kafkaKerberosServiceName = extensions.required("kerberos-material.kafka.service-name"),
140-
kafkaJaasConfig = kafka.property("sasl.jaas.config"),
138+
kerberosClientPrincipal = extensions.optional("kerberos-material.client.principal"),
139+
kerberosClientPassword = extensions.optional("kerberos-material.client.password"),
140+
kerberosClientKeytab = extensions.optional("kerberos-material.client.keytab"),
141+
krb5Conf = extensions.optional("kerberos-material.krb5-conf")
142+
?: kafka.properties["java.security.krb5.conf.local"],
143+
kafkaSecurityProtocol = kafka.properties["security.protocol"],
144+
kafkaKerberosServiceName = extensions.optional("kerberos-material.kafka.service-name")
145+
?: kafka.properties["sasl.kerberos.service.name"],
146+
kafkaJaasConfig = kafka.properties["sasl.jaas.config"],
141147
)
142148
}
143149

144150
private fun createSparkSession(environment: SparkScenarioEnvironment): SparkSession {
145-
System.setProperty("java.security.krb5.conf", environment.krb5Conf)
151+
environment.krb5Conf?.let { System.setProperty("java.security.krb5.conf", it) }
146152
return configureSpark(
147153
SparkSession.builder()
148154
.appName("bigdata-test-spark-example")
@@ -184,34 +190,59 @@ abstract class SparkBigDataScenario {
184190
.config("spark.hadoop.fs.gs.create.items.conflict.check.enable", "false")
185191
.config("spark.hadoop.fs.gs.implicit.dir.repair.enable", "false")
186192
.config("spark.hadoop.fs.gs.hierarchical.namespace.folders.enable", "false")
187-
.config("spark.hadoop.java.security.krb5.conf", environment.krb5Conf)
188-
.config("spark.hadoop.bigdata.test.kerberos.client.principal", environment.kerberosClientPrincipal)
189-
.config("spark.hadoop.bigdata.test.kerberos.client.keytab", environment.kerberosClientKeytab)
190-
.config("spark.hadoop.bigdata.test.kafka.service.name", environment.kafkaKerberosServiceName)
191-
.config("spark.hadoop.bigdata.test.kafka.jaas.config", environment.kafkaJaasConfig),
193+
.configureKerberos(environment),
192194
environment,
193195
).enableHiveSupport()
194196
.getOrCreate()
195197
}
196198

199+
private fun SparkSession.Builder.configureKerberos(environment: SparkScenarioEnvironment): SparkSession.Builder {
200+
environment.krb5Conf?.let { config("spark.hadoop.java.security.krb5.conf", it) }
201+
environment.kerberosClientPrincipal?.let {
202+
config("spark.hadoop.bigdata.test.kerberos.client.principal", it)
203+
}
204+
environment.kerberosClientKeytab?.let {
205+
config("spark.hadoop.bigdata.test.kerberos.client.keytab", it)
206+
}
207+
environment.kafkaKerberosServiceName?.let {
208+
config("spark.hadoop.bigdata.test.kafka.service.name", it)
209+
}
210+
environment.kafkaJaasConfig?.let {
211+
config("spark.hadoop.bigdata.test.kafka.jaas.config", it)
212+
}
213+
return this
214+
}
215+
197216
protected fun assertHdfsConfigStore(spark: SparkSession, hdfsUri: String, hdfsPath: String) {
198217
val exists = org.apache.hadoop.fs.FileSystem.get(URI.create(hdfsUri), spark.sparkContext().hadoopConfiguration())
199218
.use { fs -> fs.exists(org.apache.hadoop.fs.Path(hdfsPath)) }
200219
check(exists) { "Expected S3 JCEKS file in HDFS for ${spark.sparkContext().appName()}" }
201220
}
202221

203-
protected fun assertKafkaAvroInput(spark: SparkSession, bootstrapServers: String, topic: String) {
204-
val rows = spark.read()
222+
protected fun assertKafkaAvroInput(
223+
spark: SparkSession,
224+
bootstrapServers: String,
225+
topic: String,
226+
securityProtocol: String?,
227+
kerberosServiceName: String?,
228+
jaasConfig: String?,
229+
) {
230+
val reader = spark.read()
205231
.format("kafka")
206232
.option("kafka.bootstrap.servers", bootstrapServers)
207233
.option("subscribe", topic)
208234
.option("startingOffsets", "earliest")
209235
.option("endingOffsets", "latest")
210-
.option("kafka.security.protocol", "SASL_PLAINTEXT")
211-
.option("kafka.sasl.mechanism", "GSSAPI")
212-
.option("kafka.sasl.kerberos.service.name", spark.sparkContext().hadoopConfiguration().get("bigdata.test.kafka.service.name", "kafka"))
213-
.option("kafka.sasl.jaas.config", spark.sparkContext().hadoopConfiguration().get("bigdata.test.kafka.jaas.config", ""))
214-
.load()
236+
237+
if (securityProtocol == "SASL_PLAINTEXT") {
238+
reader
239+
.option("kafka.security.protocol", securityProtocol)
240+
.option("kafka.sasl.mechanism", "GSSAPI")
241+
.option("kafka.sasl.kerberos.service.name", kerberosServiceName ?: "kafka")
242+
.option("kafka.sasl.jaas.config", jaasConfig ?: "")
243+
}
244+
245+
val rows = reader.load()
215246
check(rows.count() == 2L) { "Expected two Avro Kafka records in $topic" }
216247
}
217248

@@ -343,10 +374,11 @@ data class SparkScenarioEnvironment(
343374
val gcsIcebergDataPath: String,
344375
val s3CredentialProviderPath: String,
345376
val s3CredentialProviderHdfsPath: String,
346-
val kerberosClientPrincipal: String,
347-
val kerberosClientPassword: String,
348-
val kerberosClientKeytab: String,
349-
val krb5Conf: String,
350-
val kafkaKerberosServiceName: String,
351-
val kafkaJaasConfig: String,
377+
val kerberosClientPrincipal: String?,
378+
val kerberosClientPassword: String?,
379+
val kerberosClientKeytab: String?,
380+
val krb5Conf: String?,
381+
val kafkaSecurityProtocol: String?,
382+
val kafkaKerberosServiceName: String?,
383+
val kafkaJaasConfig: String?,
352384
)

example/spark/src/test/kotlin/org/openprojectx/bigdata/test/example/spark/SparkBigDataTestExample.kt

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,12 @@ import org.openprojectx.bigdata.test.extensions.junit5.BigDataExtensions
66
import org.openprojectx.bigdata.test.junit5.BigDataTest
77

88
@BigDataExtensions("classpath:spark-bigdata-extensions.toml")
9-
@BigDataTest
9+
@BigDataTest(
10+
config = [
11+
"classpath:spark-bigdata-test-common.toml",
12+
"classpath:spark-bigdata-test-apache-hms-kerberos.toml",
13+
],
14+
)
1015
class SparkBigDataTestExample : SparkBigDataScenario() {
1116
override val runId: String get() = scenarioRunId
1217
override val s3BucketExtensionId: String get() = S3_BUCKET_ID

example/spark/src/test/resources/spark-bigdata-extensions.toml

Lines changed: 0 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -1,31 +1,3 @@
1-
[images]
2-
kerberos = "openprojectx/kerby-kdc:latest"
3-
hdfs = "apache/hadoop:3.5.0"
4-
hiveMetastore = "ghcr.io/openprojectx/hive:3.1.3-hadoop-3.4.2-gcs-4.0.4-jdk17-0.1.4"
5-
hiveMetastorePostgres = "postgres:16-alpine"
6-
clouderaHms = "ghcr.io/openprojectx/cloudera-hms:0.1.16"
7-
kafka = "apache/kafka:4.1.2"
8-
schemaRegistry = "confluentinc/cp-schema-registry:7.8.0"
9-
kafkaUi = "ghcr.io/kafbat/kafka-ui:latest"
10-
localStackS3 = "localstack/localstack:4.14.0"
11-
fakeGcs = "fsouza/fake-gcs-server:1.54"
12-
13-
[services]
14-
kerberos = true
15-
hdfs = true
16-
hiveMetastore = true
17-
kafka = true
18-
kafkaKerberos = true
19-
schemaRegistry = true
20-
localStackS3 = true
21-
fakeGcs = true
22-
23-
[ports]
24-
kafka = 19092
25-
26-
[containerLogs]
27-
mode = "FILE"
28-
291
[s3Jceks]
302
enabled = true
313
hdfsDir = "/bigdata-test/spark"
@@ -34,9 +6,6 @@ fileName = "s3.jceks"
346
[kafkaAvro]
357
enabled = true
368

37-
[kerberosMaterial]
38-
enabled = true
39-
409
[[kafkaAvro.topics]]
4110
name = "spark-avro-events"
4211
schema = "classpath:schemas/spark-event.avsc"
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
[services]
2+
kerberos = true
3+
hiveMetastore = true
4+
kafkaKerberos = true
5+
6+
[ports]
7+
kafka = 19092
Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
[services]
2+
hiveMetastore = true
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
[services]
2+
kerberos = true
3+
clouderaHms = true
4+
kafkaKerberos = true
5+
6+
[ports]
7+
kafka = 19092
Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
[services]
2+
clouderaHms = true

0 commit comments

Comments
 (0)