Skip to content

Commit 85a3983

Browse files
committed
feat: extend Spark dependency matrix with Cloudera support and task aliases
- Added Cloudera Spark/Hadoop dependency line to test matrix with resolvable runtime classpaths. - Introduced compatibility task aliases for legacy HMS test tasks. - Updated documentation to reflect matrix structure, task usage, and runtime configuration details. - Enhanced test setups with uncompressed Parquet default and improved HiveMetastoreClient compatibility.
1 parent 450255f commit 85a3983

2 files changed

Lines changed: 49 additions & 20 deletions

File tree

example/spark/build.gradle.kts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,8 @@ import org.gradle.api.attributes.Usage
77
plugins {
88
id("buildsrc.convention.kotlin-jvm")
99
id("org.openprojectx.spark.platform") version "0.1.38-SNAPSHOT"
10+
id("org.openprojectx.hadoop-native-loader") version "0.1.1"
11+
1012
}
1113

1214
description = "Spark JUnit 5 example for bigdata-test"

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

Lines changed: 47 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -155,30 +155,41 @@ abstract class SparkBigDataScenario {
155155
.appName("bigdata-test-spark-example")
156156
.master("local[2]")
157157
.config("spark.ui.enabled", "false")
158-
.config("spark.driver.extraJavaOptions", "--add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED")
159-
.config("spark.executor.extraJavaOptions", "--add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED")
158+
.config(
159+
"spark.driver.extraJavaOptions",
160+
"--add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED"
161+
)
162+
.config(
163+
"spark.executor.extraJavaOptions",
164+
"--add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED"
165+
)
160166
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
161167
.config("spark.sql.catalog.s3", "org.apache.iceberg.spark.SparkCatalog")
162168
.config("spark.sql.catalog.s3.type", "hadoop")
163169
.config("spark.sql.catalog.s3.warehouse", "s3a://${environment.s3Bucket}/warehouse")
164170
.config("spark.sql.catalog.gcs_local", "org.apache.iceberg.spark.SparkCatalog")
165171
.config("spark.sql.catalog.gcs_local.type", "hadoop")
166-
.config("spark.sql.catalog.gcs_local.warehouse", "file:${Files.createTempDirectory("bigdata-test-gcs-iceberg-warehouse-")}")
172+
.config(
173+
"spark.sql.catalog.gcs_local.warehouse",
174+
"file:${Files.createTempDirectory("bigdata-test-gcs-iceberg-warehouse-")}"
175+
)
167176
.config("spark.sql.catalog.hms", "org.apache.iceberg.spark.SparkCatalog")
168177
.config("spark.sql.catalog.hms.type", "hive")
169178
.config("spark.sql.catalog.hms.uri", environment.hiveMetastoreUri)
170-
.config("spark.sql.catalog.hms.warehouse", "file:${Files.createTempDirectory("bigdata-test-hms-warehouse-")}")
179+
.config(
180+
"spark.sql.catalog.hms.warehouse",
181+
"file:${Files.createTempDirectory("bigdata-test-hms-warehouse-")}"
182+
)
171183
.config("spark.sql.warehouse.dir", "file:${Files.createTempDirectory("bigdata-test-spark-warehouse-")}")
172184
.config("spark.sql.statistics.size.autoUpdate.enabled", "false")
173-
.config("spark.sql.parquet.compression.codec", "uncompressed")
174-
.config("spark.sql.catalog.s3.write.parquet.compression-codec", "uncompressed")
175-
.config("spark.sql.catalog.gcs_local.write.parquet.compression-codec", "uncompressed")
176-
.config("spark.sql.catalog.hms.write.parquet.compression-codec", "uncompressed")
177185
.config("hive.metastore.uris", environment.hiveMetastoreUri)
178186
.config("spark.hadoop.fs.defaultFS", environment.hdfsUri)
179187
.config("spark.hadoop.hadoop.security.credential.provider.path", environment.s3CredentialProviderPath)
180188
.config("spark.hadoop.fs.s3a.endpoint", environment.s3Endpoint)
181-
.config("spark.hadoop.fs.s3a.aws.credentials.provider", "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider")
189+
.config(
190+
"spark.hadoop.fs.s3a.aws.credentials.provider",
191+
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
192+
)
182193
.config("spark.hadoop.fs.s3a.access.key", "test")
183194
.config("spark.hadoop.fs.s3a.secret.key", "test")
184195
.config("spark.hadoop.fs.s3a.path.style.access", "true")
@@ -219,8 +230,9 @@ abstract class SparkBigDataScenario {
219230
}
220231

221232
protected fun assertHdfsConfigStore(spark: SparkSession, hdfsUri: String, hdfsPath: String) {
222-
val exists = org.apache.hadoop.fs.FileSystem.get(URI.create(hdfsUri), spark.sparkContext().hadoopConfiguration())
223-
.use { fs -> fs.exists(org.apache.hadoop.fs.Path(hdfsPath)) }
233+
val exists =
234+
org.apache.hadoop.fs.FileSystem.get(URI.create(hdfsUri), spark.sparkContext().hadoopConfiguration())
235+
.use { fs -> fs.exists(org.apache.hadoop.fs.Path(hdfsPath)) }
224236
check(exists) { "Expected S3 JCEKS file in HDFS for ${spark.sparkContext().appName()}" }
225237
}
226238

@@ -261,11 +273,9 @@ abstract class SparkBigDataScenario {
261273
) {
262274
val identifier = "$catalog.$namespace.$table"
263275
spark.sql("CREATE NAMESPACE IF NOT EXISTS $catalog.$namespace")
264-
val properties = listOfNotNull(
265-
"'write.parquet.compression-codec'='uncompressed'",
266-
dataPath?.let { "'write.data.path'='$it'" },
267-
).joinToString(", ")
268-
val tableProperties = " TBLPROPERTIES ($properties)"
276+
val tableProperties = icebergTableProperties(
277+
"write.data.path" to dataPath,
278+
)
269279
spark.sql(
270280
"""
271281
CREATE TABLE $identifier (
@@ -318,7 +328,7 @@ abstract class SparkBigDataScenario {
318328
UNION ALL
319329
SELECT 2 AS id, 'beta' AS name, '$storageName' AS storage
320330
""".trimIndent(),
321-
).write().mode("overwrite").option("compression", "uncompressed").parquet(location)
331+
).write().mode("overwrite").parquet(location)
322332
}
323333
spark.sql(
324334
"""
@@ -337,12 +347,18 @@ abstract class SparkBigDataScenario {
337347
}
338348

339349
private fun assertParquetFiles(spark: SparkSession, location: String) {
340-
val files = org.apache.hadoop.fs.FileSystem.get(URI.create(location), spark.sparkContext().hadoopConfiguration())
341-
.use { fs -> fs.listStatus(org.apache.hadoop.fs.Path(location)).map { it.path.name } }
350+
val files =
351+
org.apache.hadoop.fs.FileSystem.get(URI.create(location), spark.sparkContext().hadoopConfiguration())
352+
.use { fs -> fs.listStatus(org.apache.hadoop.fs.Path(location)).map { it.path.name } }
342353
check(files.any { it.endsWith(".parquet") }) { "Expected Parquet files under $location, got $files" }
343354
}
344355

345-
private fun assertHiveMetastoreExternalParquetTable(hiveMetastoreUri: String, database: String, table: String, location: String) {
356+
private fun assertHiveMetastoreExternalParquetTable(
357+
hiveMetastoreUri: String,
358+
database: String,
359+
table: String,
360+
location: String
361+
) {
346362
val conf = HiveConf()
347363
conf.setVar(HiveConf.ConfVars.METASTOREURIS, hiveMetastoreUri)
348364
val client = hiveMetastoreClient(conf)
@@ -371,6 +387,17 @@ abstract class SparkBigDataScenario {
371387
?: error("No supported HiveMetaStoreClient constructor found")
372388
return constructor.newInstance(conf) as IMetaStoreClient
373389
}
390+
391+
private fun icebergTableProperties(vararg properties: Pair<String, String?>): String {
392+
val entries = properties.mapNotNull { (key, value) ->
393+
value?.let { "'$key'='$it'" }
394+
}
395+
return if (entries.isEmpty()) {
396+
""
397+
} else {
398+
" TBLPROPERTIES (${entries.joinToString(", ")})"
399+
}
400+
}
374401
}
375402

376403
fun interface SparkScenarioCheck {

0 commit comments

Comments
 (0)