andygrove commented on code in PR #6117:
URL: https://github.com/apache/datafusion-comet/pull/6117#discussion_r4109345875
##########
spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala:
##########
@@ -1709,6 +1710,86 @@ class CometIcebergWriteActionSuite
}
}
+ test("native acceleration: a post-native handoff failure cleans up task
files") {
+ assumeNativeAcceleration()
+ withIcebergCatalog { warehouseDir =>
+ createTable(warehouseDir, "handoff_target", partitionSpec = "")
+ coalesceInsert("handoff_target", Seq((0, "seed", 0.0)))
+ val before = countSnapshots("handoff_target")
+ val root = dataDir("handoff_target").toPath.toAbsolutePath
+
+ def relativePath(location: String): String = {
+ val uri = new java.net.URI(location)
+ val file = if (uri.getScheme == null) new File(location) else new
File(uri)
+ root.relativize(file.toPath.toAbsolutePath).toString
+ }
+
+ def metadataFiles: Set[String] = spark
+ .sql(s"SELECT file_path FROM $catalog.$ns.handoff_target.files")
+ .collect()
+ .map(row => relativePath(row.getString(0)))
+ .toSet
+
+ val committed = metadataFiles
+ assert(committed.nonEmpty, "seed write did not create a data file")
+ assert(parquetFiles(root.toFile) == committed)
+
+ val session = spark
+ import session.implicits._
+ (1 to 1000)
+ .map(i => (i, s"r$i", i.toDouble))
+ .toDF("id", "region", "amount")
+ .coalesce(1)
+ .createOrReplaceTempView("handoff_src")
+
+ val attempts = new AtomicReference[Vector[(Int, Int,
Vector[String])]](Vector.empty)
+ val (failedPlans, error) = withNativeEnabled {
+ CometIcebergWriteExec.withPostNativeHandoffFailpoint { locations =>
+ val tc = TaskContext.get()
+ attempts.getAndUpdate(_ :+ ((tc.partitionId(), tc.attemptNumber(),
locations.toVector)))
+ throw new RuntimeException("post-native handoff injected failure")
+ } {
+ captureFailedPlans(spark) {
+ spark.sql(s"INSERT INTO $catalog.$ns.handoff_target " +
+ "SELECT id, region, amount FROM handoff_src")
+ }
+ }
+ }
+ assert(
+ error.toSeq
+ .flatMap(exceptionChain)
+ .exists(t =>
+ Option(t.getMessage).exists(_.contains("post-native handoff
injected failure"))),
+ s"expected the handoff failure to reach Spark, got $error")
+ assert(
+ failedPlans.exists(p =>
+ collectWithSubqueries(p) { case w: CometIcebergWriteExec => w
}.nonEmpty),
+ s"failed write did not run
natively:\n${failedPlans.mkString("\n--\n")}")
+ val handoffs = attempts.get()
+ // CometTestBase starts local[5]. Spark gives that master one allowed
task failure, so
+ // the job aborts without a retry and the handoff runs once.
+ assert(handoffs.map(_._1) == Vector(0), s"expected one task: $handoffs")
Review Comment:
This ties the test to `local[5]`, and #6111 changes this suite's master to
`local[5,2]`. With both merged the handoff runs once per attempt, and this
fails with `Vector(0, 0) did not equal Vector(0)`. Could it accept every
attempt of partition 0 instead? `handoffs.forall(_._1 == 0)` together with
`handoffs.map(_._2) == handoffs.indices.toVector` holds under either master. I
tried that with all three of your #5646 PRs merged into `main`, and the whole
suite passed on Spark 3.4, 4.0 and 4.1. Under `local[5,2]` it also checks that
both attempts clean up their own files.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]