hudi-agent commented on code in PR #19928:
URL: https://github.com/apache/hudi/pull/19928#discussion_r4127272720


##########
pom.xml:
##########
@@ -220,6 +220,12 @@
     <skipUTs>${skipTests}</skipUTs>
     <skipFTs>${skipTests}</skipFTs>
     <skipITs>${skipTests}</skipITs>
+    <!-- Surefire forks per module, expressed per CPU core the way Maven -T 
is, so a 4-core GitHub
+         runner gets two forks and a 2-core Azure agent gets one; each fork 
gets the full argLine
+         heap; override with -Dhudi.surefire.forkCount=1 to serialise a module 
(for example when
+         bisecting a test that depends on shared state), and pass an explicit 
integer on machines
+         with a single core. -->
+    <hudi.surefire.forkCount>0.5C</hudi.surefire.forkCount>

Review Comment:
   🤖 Since this default lives in the root pom, it also applies to every local 
`mvn test`/`mvn install`. On a 12-core laptop that's 6 forks, each with 
`-Xmx4g` plus its own Spark/Hive/Kafka services, and the functional profiles no 
longer pin `forkCount=1`. Would it be safer to keep the pom default at `1` and 
pass `-Dhudi.surefire.forkCount=0.5C` (or `2`) from the CI workflows? This may 
belong to the base PR (#19892) rather than this one.
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



##########
hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java:
##########
@@ -350,9 +350,26 @@ public HoodieFileFormat getFormat() {
    */
   @Override
   public List<GenericRecord> readAvroRecords(HoodieStorage storage, 
StoragePath filePath) {
+    // A schema-less read takes the schema from the file; the copy drops any 
read schema or
+    // projection a previous caller left on the shared configuration.
+    Configuration conf = storage.getConf().unwrapCopyAs(Configuration.class);
+    conf.unset("parquet.avro.read.schema");

Review Comment:
   🤖 This protects `readAvroRecords` from the leaked schema, but the leak 
itself is still there: `HDFSParquetImporterUtils#buildHoodieRecordsForImport` 
(line 185) writes `parquet.avro.read.schema` onto `jsc.hadoopConfiguration()`. 
That happens after `Job.getInstance` has already copied the conf, so the import 
job never sees it. Would it be worth dropping that write (or moving it onto 
`job.getConfiguration()`) so the procedure stops mutating the SparkContext-wide 
conf?
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



##########
hudi-sync/hudi-hive-sync/src/test/java/org/apache/hudi/hive/testutils/HiveTestService.java:
##########
@@ -63,7 +63,9 @@
 public class HiveTestService {
   private static final int CONNECTION_TIMEOUT_MS = 30000;
   private static final String BIND_HOST = "127.0.0.1";
-  private static final int HS2_THRIFT_PORT = 9999;
+  // Chosen once per JVM rather than fixed at 9999 so that surefire forks 
running this module concurrently
+  // each get their own HiveServer2 port; within one JVM every HiveTestService 
instance still shares the port, as before.
+  private static final int HS2_THRIFT_PORT = NetworkTestUtils.nextFreePort();

Review Comment:
   🤖 Choosing this once at class init keeps HS2 (and ZK, same pattern) on one 
ephemeral-range port for the whole reused fork. Between classes, while HS2 is 
stopped, the other fork's Spark/Kafka/`nextFreePort()` binds can take it and 
the next `start()` fails to bind; the metastore `nextFreePort()` in 
`configureHive` can even return this same port. Could the port be picked per 
`start()` like the metastore port, with callers using `getJdbcHive2Url()` 
instead of the static `HS2_JDBC_URL`?
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/insert/TestInsertTable4.scala:
##########
@@ -290,10 +290,22 @@ class TestInsertTable4 extends HoodieSparkSqlTestBase {
   var listenerCallCount: Int = 0
   var countDownLatch: CountDownLatch = _
 
-  // add a listener for stages for parallelism checking with stage name
-  class StageParallelismListener(var stageName: String) extends SparkListener {
+  // add a listener for stages for parallelism checking with stage name, 
restricted to the stages of
+  // the jobs this test started. Hudi overwrites the job group id and 
description on its own jobs, so
+  // the discriminator is a local property the test sets before the query.
+  class StageParallelismListener(var stageName: String, scope: String) extends 
SparkListener {
+    private val scopedStageIds = ConcurrentHashMap.newKeySet[Int]()
+
+    override def onJobStart(jobStart: SparkListenerJobStart): Unit = {
+      if (jobStart.properties != null
+        && scope == 
jobStart.properties.getProperty(TestInsertTable4.SCOPE_PROPERTY)) {
+        jobStart.stageIds.foreach(stageId => scopedStageIds.add(stageId))
+      }
+    }
+
     override def onStageSubmitted(stageSubmitted: 
SparkListenerStageSubmitted): Unit = {

Review Comment:
   🤖 nit: the SCOPE_PROPERTY/UUID scoping plus manual set/clear of the local 
property in try/finally is duplicated almost verbatim across two test methods — 
worth factoring into a small helper so future edits to one copy don't silently 
diverge from the other.
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala:
##########
@@ -71,25 +82,56 @@ class HoodieSparkSqlTestBase extends FunSuite with 
BeforeAndAfterAll {
   //       is consistent with the fixtures
   DateTimeZone.setDefault(DateTimeZone.UTC)
   TimeZone.setDefault(DateTimeUtils.getTimeZone("UTC"))
-  protected lazy val spark: SparkSession = SparkSession.builder()
-    .config("spark.sql.warehouse.dir", sparkWareHouse.getCanonicalPath)
-    .config("spark.sql.session.timeZone", "UTC")
-    .config("hoodie.insert.shuffle.parallelism", "4")
-    .config("hoodie.upsert.shuffle.parallelism", "4")
-    .config("hoodie.delete.shuffle.parallelism", "4")
-    .config(sparkConf())
-    .getOrCreate()
+  protected lazy val spark: SparkSession = if (sharedSessionEnabled) {
+    val session = HoodieSparkSqlTestBase.sharedBaseSession().newSession()
+    SparkSession.setActiveSession(session)
+    applySuiteConfToSharedSession(session)
+    session
+  } else {
+    HoodieSparkSqlTestBase.sessionBuilder(sparkWareHouse, 
sparkConf()).getOrCreate()
+  }
 
   private var tableId = new AtomicInteger(0)
 
   private var extraConf = Map[String, String]()

Review Comment:
   🤖 nit: `hadoopConfOverrides` is a `var Seq` mutated via `:+=` and later 
reversed — might read more clearly as an `ArrayBuffer` given it's accumulating 
state to undo later.
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala:
##########
@@ -71,25 +82,56 @@ class HoodieSparkSqlTestBase extends FunSuite with 
BeforeAndAfterAll {
   //       is consistent with the fixtures
   DateTimeZone.setDefault(DateTimeZone.UTC)
   TimeZone.setDefault(DateTimeUtils.getTimeZone("UTC"))
-  protected lazy val spark: SparkSession = SparkSession.builder()
-    .config("spark.sql.warehouse.dir", sparkWareHouse.getCanonicalPath)
-    .config("spark.sql.session.timeZone", "UTC")
-    .config("hoodie.insert.shuffle.parallelism", "4")
-    .config("hoodie.upsert.shuffle.parallelism", "4")
-    .config("hoodie.delete.shuffle.parallelism", "4")
-    .config(sparkConf())
-    .getOrCreate()
+  protected lazy val spark: SparkSession = if (sharedSessionEnabled) {
+    val session = HoodieSparkSqlTestBase.sharedBaseSession().newSession()
+    SparkSession.setActiveSession(session)
+    applySuiteConfToSharedSession(session)
+    session
+  } else {
+    HoodieSparkSqlTestBase.sessionBuilder(sparkWareHouse, 
sparkConf()).getOrCreate()
+  }
 
   private var tableId = new AtomicInteger(0)
 
   private var extraConf = Map[String, String]()
 
+  // Shared mode: spark.hadoop.* keys this suite set on the shared Hadoop 
conf, with the value they replaced.
+  private var hadoopConfOverrides: Seq[(String, String)] = Seq.empty
+
   def sparkConf(): SparkConf = {
     val conf = getSparkConfForTest("Hoodie SQL Test")
     conf.setAll(extraConf)
     conf
   }
 
+  /**
+   * Shared mode: the context-level SparkConf is fixed, so the deltas a suite 
adds through extraConf or a
+   * sparkConf() override go to its session conf (hoodie.* and spark.sql.* 
keys), except spark.hadoop.*
+   * keys, which the write client reads from sparkContext.hadoopConfiguration 
and which are restored in
+   * afterAll. Any other spark.* key, and any static SQL conf, is a setting a 
child session cannot change,
+   * so it is rejected before anything is mutated: a partial apply would leave 
the shared Hadoop conf
+   * changed for every later suite in the JVM.
+   */
+  private def applySuiteConfToSharedSession(session: SparkSession): Unit = {
+    val defaults = getSparkConfForTest("Hoodie SQL Test").getAll.toMap
+    val deltas = sparkConf().getAll.filterNot { case (k, v) => 
defaults.get(k).contains(v) }
+    val (hadoopKeys, sessionKeys) = deltas.partition { case (k, _) => 
k.startsWith("spark.hadoop.") }
+    val rejected = sessionKeys.collect {
+      case (k, _) if (k.startsWith("spark.") && !k.startsWith("spark.sql.")) 
|| SQLConf.isStaticConfigKey(k) => k
+    }
+    if (rejected.nonEmpty) {
+      throw new IllegalArgumentException(
+        s"${rejected.mkString(", ")}: SparkContext-level or static settings; 
shared session mode cannot apply them per suite")
+    }
+    val hadoopConf = session.sparkContext.hadoopConfiguration
+    hadoopKeys.foreach { case (k, v) =>

Review Comment:
   🤖 The save/restore here is only safe if suites that set `spark.hadoop.*` 
never overlap. If two non-exclusive suites set the same key, the first 
`afterAll` unsets the second suite's value mid-run, and the second then 
restores the first suite's value for the rest of the JVM. Could we fail fast 
when `hadoopKeys.nonEmpty` and the suite isn't an `ExclusiveSuite`, so a future 
suite can't hit this silently?
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to