From 96558ca625a11844e66bb172655c1909f561297f Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 23 Sep 2026 08:04:11 -0600 Subject: [PATCH 1/2] test: report which writer ran each Iceberg write in the Iceberg CI jobs The Iceberg Spark test jobs run with the native Iceberg writer enabled, but a passing job cannot tell a native write from a silent fallback: the fallback warnings never reach the job log and no upstream test asserts which writer ran. Add a test-only config, spark.comet.testing.icebergWriteReport.dir, backed by the COMET_ICEBERG_WRITE_REPORT_DIR environment variable. When it is set, the driver plugin registers IcebergWriteReportListener, which appends one JSON line per Iceberg write: native, JVM under the split operator (with Comet's fallback reasons), or Spark's own V2 write. The core and extensions jobs set the variable, so the Iceberg diffs are unchanged, and dev/ci/summarize-iceberg-writes.py renders a per-job and all-shards table on the job summary page. dev/local-ci.sh prints the same summary. Closes #6148. --- .../workflows/iceberg_spark_test_reusable.yml | 32 +++- dev/ci/summarize-iceberg-writes.py | 149 ++++++++++++++++++ dev/local-ci.sh | 11 +- .../contributor-guide/iceberg-spark-tests.md | 27 +++- .../scala/org/apache/comet/CometConf.scala | 13 ++ .../iceberg/IcebergWriteReportListener.scala | 119 ++++++++++++++ .../main/scala/org/apache/spark/Plugins.scala | 46 ++++-- .../comet/CometIcebergWriteActionSuite.scala | 56 ++++++- .../org/apache/spark/CometPluginsSuite.scala | 18 ++- 9 files changed, 448 insertions(+), 23 deletions(-) create mode 100644 dev/ci/summarize-iceberg-writes.py create mode 100644 spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala diff --git a/.github/workflows/iceberg_spark_test_reusable.yml b/.github/workflows/iceberg_spark_test_reusable.yml index ae20b6d7683..9b4a5ae3000 100644 --- a/.github/workflows/iceberg_spark_test_reusable.yml +++ b/.github/workflows/iceberg_spark_test_reusable.yml @@ -154,12 +154,20 @@ jobs: run: | cd apache-iceberg rm -rf /root/.m2/repository/org/apache/parquet # somehow parquet cache requires cleanups - ENABLE_COMET=true ENABLE_COMET_ONHEAP=true ./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \ + # COMET_ICEBERG_WRITE_REPORT_DIR records which writer ran each Iceberg + # write; see dev/ci/summarize-iceberg-writes.py. + ENABLE_COMET=true ENABLE_COMET_ONHEAP=true COMET_ICEBERG_WRITE_REPORT_DIR="$PWD/build/comet-iceberg-writes" \ + ./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \ :iceberg-spark:iceberg-spark-${{ inputs.spark-short }}_${{ inputs.scala }}:test \ --init-script ../dev/ci/iceberg-test-shards.gradle \ -PcometShardTask=:iceberg-spark:iceberg-spark-${{ inputs.spark-short }}_${{ inputs.scala }}:test \ -PcometShardIndex=${{ matrix.shard }} -PcometShardCount=${{ needs.build-native.outputs.shard-count }} \ -Pquick=true -x javadoc + - name: Summarize Iceberg writes + if: ${{ !cancelled() }} + run: | + python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark shard ${{ matrix.shard }}" \ + apache-iceberg/build/comet-iceberg-writes - name: Upload Iceberg shard inventory and test reports if: ${{ !cancelled() }} # iceberg-spark-shard-coverage downloads the inventory, so a flaky @@ -170,6 +178,7 @@ jobs: path: | apache-iceberg/**/build/comet-shards/*.json apache-iceberg/**/build/test-results/test/*.xml + apache-iceberg/build/comet-iceberg-writes/*.jsonl retention-days: 7 iceberg-spark-shard-coverage: @@ -190,6 +199,11 @@ jobs: run: | python3 dev/ci/check-iceberg-shards.py --manifests iceberg-shard-reports \ --task :iceberg-spark:iceberg-spark-${{ inputs.spark-short }}_${{ inputs.scala }}:test + - name: Summarize Iceberg writes across shards + if: ${{ !cancelled() }} + run: | + python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark, all shards" \ + iceberg-shard-reports iceberg-spark-extensions: needs: build-native @@ -222,9 +236,23 @@ jobs: run: | cd apache-iceberg rm -rf /root/.m2/repository/org/apache/parquet # somehow parquet cache requires cleanups - ENABLE_COMET=true ENABLE_COMET_ONHEAP=true ./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \ + ENABLE_COMET=true ENABLE_COMET_ONHEAP=true COMET_ICEBERG_WRITE_REPORT_DIR="$PWD/build/comet-iceberg-writes" \ + ./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \ :iceberg-spark:iceberg-spark-extensions-${{ inputs.spark-short }}_${{ inputs.scala }}:test \ -Pquick=true -x javadoc + - name: Summarize Iceberg writes + if: ${{ !cancelled() }} + run: | + python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark-extensions" \ + apache-iceberg/build/comet-iceberg-writes + - name: Upload Iceberg write report + if: ${{ !cancelled() }} + uses: ./.github/actions/upload-artifact-retry + with: + name: iceberg-spark-extensions-writes-${{ inputs.iceberg-full }}-spark-${{ inputs.spark-full }}-scala-${{ inputs.scala }}-jdk${{ inputs.java }}-attempt-${{ github.run_attempt }} + path: apache-iceberg/build/comet-iceberg-writes/*.jsonl + if-no-files-found: ignore + retention-days: 7 iceberg-spark-runtime: needs: build-native diff --git a/dev/ci/summarize-iceberg-writes.py b/dev/ci/summarize-iceberg-writes.py new file mode 100644 index 00000000000..64cc2b94516 --- /dev/null +++ b/dev/ci/summarize-iceberg-writes.py @@ -0,0 +1,149 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Summarize which writer ran the Iceberg writes of an Iceberg Spark test run. + +The Iceberg Spark test jobs set COMET_ICEBERG_WRITE_REPORT_DIR, so Comet's +IcebergWriteReportListener appends one JSON line per Iceberg write to a file +in that directory. This prints how many writes ran on Comet's native writer, +how many Comet's split operator left on Iceberg's JVM writer and why, and how +many Spark planned without Comet's split operator: + + python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark shard 1" DIR... + +Files under a directory named like a shard artifact (...-shard-N-attempt-M) +are counted only for the latest attempt of each shard, so a rerun of failed +jobs does not count a shard twice. The summary is also appended to +$GITHUB_STEP_SUMMARY when that is set. It never fails the job. +""" + +import argparse +from collections import Counter +import json +import os +from pathlib import Path +import re + + +SHARD_ATTEMPT = re.compile(r"-shard-(\d+)-attempt-(\d+)$") +WRITERS = [ + ("native", "Comet native writer"), + ("jvm", "Iceberg JVM writer under Comet's split operator"), + ("spark", "Spark V2 write, not planned by Comet's split operator"), +] +TOP_REASONS = 20 + + +def report_files(roots): + """Every report file under roots, keeping only the latest attempt of each shard.""" + latest = {} + files = [] + for root in roots: + for path in sorted(Path(root).rglob("*.jsonl")): + match = next( + (SHARD_ATTEMPT.search(part) for part in path.parts if SHARD_ATTEMPT.search(part)), + None, + ) + if match: + shard, attempt = int(match.group(1)), int(match.group(2)) + latest[shard] = max(latest.get(shard, 0), attempt) + files.append((path, (shard, attempt))) + else: + files.append((path, None)) + return [path for path, key in files if key is None or latest[key[0]] == key[1]] + + +def load(roots): + writes = [] + for path in report_files(roots): + for line in path.read_text(encoding="utf-8").splitlines(): + if line.strip(): + writes.append(json.loads(line)) + return writes + + +def cell(text): + return " ".join(text.split()).replace("|", "\\|") + + +def summarize(title, writes): + lines = [f"### Iceberg writes: {title}", ""] + if not writes: + lines.append( + "No Iceberg writes were recorded. Either the target ran none or " + "COMET_ICEBERG_WRITE_REPORT_DIR did not reach the test JVMs." + ) + return "\n".join(lines) + "\n" + + total = len(writes) + by_writer = Counter(w["writer"] for w in writes) + lines += ["| Writer | Writes | Share |", "| --- | ---: | ---: |"] + for key, label in WRITERS: + count = by_writer.get(key, 0) + lines.append(f"| {label} | {count} | {100.0 * count / total:.1f}% |") + lines += [f"| Total | {total} | |", ""] + failed = sum(1 for w in writes if w.get("failed")) + if failed: + lines += [f"{failed} of the {total} writes ran in queries that failed.", ""] + + reasons = Counter() + for w in writes: + if w["writer"] == "jvm": + for reason in w.get("reasons") or ["(no reason recorded)"]: + reasons[reason] += 1 + if reasons: + lines += [ + "#### Why the split operator kept the JVM writer", + "", + "A write can have several reasons, so the counts can add up to more than the " + "JVM writes.", + "", + "| Writes | Reason |", + "| ---: | --- |", + ] + for reason, count in reasons.most_common(TOP_REASONS): + lines.append(f"| {count} | {cell(reason)} |") + if len(reasons) > TOP_REASONS: + lines.append(f"| | and {len(reasons) - TOP_REASONS} more reasons |") + lines.append("") + + operators = Counter(w["node"] for w in writes if w["writer"] == "spark") + if operators: + lines += ["#### Spark V2 writes by operator", "", "| Writes | Operator |", "| ---: | --- |"] + for node, count in operators.most_common(): + lines.append(f"| {count} | {cell(node)} |") + lines.append("") + + return "\n".join(lines) + "\n" + + +def main(): + parser = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + parser.add_argument("--title", required=True, help="heading for the summary") + parser.add_argument("roots", nargs="+", help="directories holding the report files") + args = parser.parse_args() + + summary = summarize(args.title, load(args.roots)) + print(summary) + step_summary = os.environ.get("GITHUB_STEP_SUMMARY") + if step_summary: + with open(step_summary, "a", encoding="utf-8") as out: + out.write(summary + "\n") + + +if __name__ == "__main__": + main() diff --git a/dev/local-ci.sh b/dev/local-ci.sh index 8e3267302d9..a6a9d57d83d 100755 --- a/dev/local-ci.sh +++ b/dev/local-ci.sh @@ -389,16 +389,25 @@ run_iceberg() { gradlew ":iceberg-spark:iceberg-spark-runtime-${SPARK}_${SCALA}:integrationTest" ;; esac + case "$target" in + shard-* | extensions) + python3 "$REPO/dev/ci/summarize-iceberg-writes.py" --title "$target" \ + "$dest/build/comet-iceberg-writes/$target" + ;; + esac ok "$target took $(hms $((SECONDS - started)))" done } -# Reads $dest, $spark and $SCALA from run_iceberg. +# Reads $dest, $spark, $SCALA and $target from run_iceberg. gradlew() { ( cd "$dest" # shellcheck disable=SC2031 export SPARK_LOCAL_IP=localhost ENABLE_COMET=true ENABLE_COMET_ONHEAP=true + # One directory per target, emptied first, so a rerun reports only its own writes. + export COMET_ICEBERG_WRITE_REPORT_DIR="$dest/build/comet-iceberg-writes/$target" + rm -rf "$COMET_ICEBERG_WRITE_REPORT_DIR" ./gradlew "-DsparkVersions=$SPARK" "-DscalaVersion=$SCALA" \ -DflinkVersions= -DkafkaVersions= "$@" -Pquick=true -x javadoc ) diff --git a/docs/source/contributor-guide/iceberg-spark-tests.md b/docs/source/contributor-guide/iceberg-spark-tests.md index fb0e262e341..07358aa7b85 100644 --- a/docs/source/contributor-guide/iceberg-spark-tests.md +++ b/docs/source/contributor-guide/iceberg-spark-tests.md @@ -43,7 +43,9 @@ Here is an overview of the changes that the diffs make to Iceberg: `LocalTableScanExec`, the conversion is declined, and the write silently runs on the JVM writer. Many Iceberg suites seed their data that way, so leaving it off hides the native writer from most of the write surface. - Enable fallback logging (`spark.comet.explainFallback.enabled`) so that every operator Comet declines is - reported in the test output together with the reason it was declined. + reported in the test output together with the reason it was declined. The output goes to the JUnit XML + reports rather than the CI job log; see [Which writer ran each Iceberg write](#which-writer-ran-each-iceberg-write) + for how CI reports native write coverage. [#3739]: https://github.com/apache/datafusion-comet/pull/3739 [#5259]: https://github.com/apache/datafusion-comet/issues/5259 @@ -154,3 +156,26 @@ path, reflection code (`org.apache.comet.iceberg.IcebergReflection`), or other l can differ across Iceberg versions. The Comet test suites in the Linux build do not exercise Iceberg's own Spark tests, so without the label the first Iceberg 1.11 verdict is the merge queue's, and the first verdict on the older Iceberg versions is the nightly run's, after the change has landed. + +### Which writer ran each Iceberg write + +A passing Iceberg job does not show that Comet's native writer ran. `CometIcebergNativeWrite` falls back +to Iceberg's JVM writer without failing the write, no upstream test asserts which writer ran, and Gradle +does not copy the fallback warnings into the job log. So the core and extensions jobs set +`COMET_ICEBERG_WRITE_REPORT_DIR`, the environment variable behind the test-only config +`spark.comet.testing.icebergWriteReport.dir`. When it is set, the Comet driver plugin registers +`IcebergWriteReportListener`, which writes one JSON line for each Iceberg write the tests run. Each line +records one of three writers: + +- `native`: Comet's native writer (`CometIcebergWriteExec`). +- `jvm`: Comet's split operator planned the write but kept Iceberg's JVM writer (`IcebergWriteExec`). + The line includes the reasons Comet recorded for not converting it. +- `spark`: Spark's own V2 write operator ran the write, for example `WriteDelta` for merge-on-read, + so Comet's split operator never saw it. + +`dev/ci/summarize-iceberg-writes.py` turns these records into a table on the job's summary page. It +shows the count and share of each writer, the most common fallback reasons, and the Spark write +operators. Each shard and the extensions job gets its own table. The shard coverage job adds one for +all shards together. The raw records are uploaded with the job's other reports. The summary never +fails a job. `dev/local-ci.sh iceberg` prints the same summary after each shard and after the +extensions target. diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index d371f1ba44c..b2b50373bcd 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -877,6 +877,19 @@ object CometConf extends ShimCometConf { .booleanConf .createWithDefault(true) + val COMET_ICEBERG_WRITE_REPORT_DIR: ConfigEntry[String] = + conf("spark.comet.testing.icebergWriteReport.dir") + .internal() + .category(CATEGORY_TESTING) + .doc("Test-only. When set, the Comet driver plugin registers a query listener that " + + "records every Iceberg write the application runs, the writer that ran it (Comet's " + + "native writer, the JVM writer behind Comet's split operator, or Spark's own V2 write) " + + "and the reasons Comet recorded for not writing natively, as JSON lines in this " + + "directory. `dev/ci/summarize-iceberg-writes.py` summarizes them. The Iceberg Spark " + + "test jobs set it through the environment variable so the Iceberg diffs need no change.") + .stringConf + .createWithEnvVarOrDefault("COMET_ICEBERG_WRITE_REPORT_DIR", "") + val COMET_SPARK_TO_ARROW_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.sparkToColumnar.enabled") .category(CATEGORY_EXEC) diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala new file mode 100644 index 00000000000..c1e9538732c --- /dev/null +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala @@ -0,0 +1,119 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.iceberg + +import java.io.File +import java.nio.charset.StandardCharsets.UTF_8 +import java.nio.file.{Files, StandardOpenOption} +import java.util.UUID + +import scala.util.control.NonFatal + +import org.json4s.JsonDSL._ +import org.json4s.jackson.JsonMethods._ + +import org.apache.spark.SparkConf +import org.apache.spark.internal.Logging +import org.apache.spark.sql.comet.{CometIcebergWriteExec, IcebergWriteExec} +import org.apache.spark.sql.execution.{CommandResultExec, QueryExecution, SparkPlan} +import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanExec, QueryStageExec} +import org.apache.spark.sql.execution.datasources.v2.V2ExistingTableWriteExec +import org.apache.spark.sql.util.QueryExecutionListener + +import org.apache.comet.CometConf.COMET_ICEBERG_WRITE_REPORT_DIR +import org.apache.comet.CometExplainInfo + +/** + * Test-only listener that records which writer ran each Iceberg write, so a CI job running + * Iceberg's own Spark suites can tell a native write from a silent fallback. The Comet driver + * plugin registers it when `spark.comet.testing.icebergWriteReport.dir` is set. Each write is + * appended as one JSON line to a file of its own in that directory, which + * `dev/ci/summarize-iceberg-writes.py` reads. + */ +class IcebergWriteReportListener(conf: SparkConf) extends QueryExecutionListener with Logging { + + private val reportFile: File = { + val dir = new File( + conf + .get(COMET_ICEBERG_WRITE_REPORT_DIR.key, COMET_ICEBERG_WRITE_REPORT_DIR.defaultValue.get)) + dir.mkdirs() + // One file per listener: Gradle runs several test JVMs, each possibly with several sessions. + new File(dir, s"iceberg-writes-${UUID.randomUUID()}.jsonl") + } + + override def onSuccess(funcName: String, qe: QueryExecution, durationNs: Long): Unit = + record(qe, failed = false) + + override def onFailure(funcName: String, qe: QueryExecution, exception: Exception): Unit = + record(qe, failed = true) + + private def record(qe: QueryExecution, failed: Boolean): Unit = { + try { + val lines = IcebergWriteReportListener.writes(qe.executedPlan).map { w => + compact( + render(("writer" -> w.writer) ~ ("node" -> w.node) ~ ("reasons" -> w.reasons.toList) ~ + ("failed" -> failed))) + "\n" + } + if (lines.nonEmpty) { + synchronized { + Files.write( + reportFile.toPath, + lines.mkString.getBytes(UTF_8), + StandardOpenOption.CREATE, + StandardOpenOption.APPEND) + } + } + } catch { + // A query that failed during planning has no executed plan; nothing was written. + case NonFatal(e) => logWarning(s"Could not record Iceberg writes for a query: $e") + } + } +} + +object IcebergWriteReportListener { + + /** Comet's native (iceberg-rust) writer ran the write. */ + val Native = "native" + + /** Comet's split operator planned the write, but Iceberg's JVM writer ran it. */ + val Jvm = "jvm" + + /** Spark's own V2 write operator ran the write; Comet's split operator did not plan it. */ + val Spark = "spark" + + case class IcebergWrite(writer: String, node: String, reasons: Seq[String]) + + /** The Iceberg writes in an executed plan, with the reasons Comet did not write natively. */ + def writes(plan: SparkPlan): Seq[IcebergWrite] = plan match { + // A command's writes are reported by the command's own execution, which runs eagerly before + // the query wrapping its result. Descending here would count them twice. + case _: CommandResultExec => Nil + case a: AdaptiveSparkPlanExec => writes(a.executedPlan) + case s: QueryStageExec => writes(s.plan) + case w: CometIcebergWriteExec => Seq(IcebergWrite(Native, w.nodeName, Nil)) + case w: IcebergWriteExec => + val reasons = w.getTagValue(CometExplainInfo.FALLBACK_REASONS).getOrElse(Set.empty) + Seq(IcebergWrite(Jvm, w.nodeName, reasons.toSeq.sorted)) + case w: V2ExistingTableWriteExec + if w.write.getClass.getName.startsWith("org.apache.iceberg.") => + Seq(IcebergWrite(Spark, w.nodeName, Nil)) + case p => p.children.flatMap(writes) + } +} diff --git a/spark/src/main/scala/org/apache/spark/Plugins.scala b/spark/src/main/scala/org/apache/spark/Plugins.scala index 24b6ab3eeaa..862a064d63f 100644 --- a/spark/src/main/scala/org/apache/spark/Plugins.scala +++ b/spark/src/main/scala/org/apache/spark/Plugins.scala @@ -29,9 +29,10 @@ import org.apache.spark.sql.internal.StaticSQLConf import org.apache.comet.{COMET_VERSION, CometSparkSessionExtensions, NativeBase} import org.apache.comet.{CometConf, ConfigEntry} -import org.apache.comet.CometConf.{COMET_METRICS_ENABLED, COMET_ONHEAP_ENABLED} +import org.apache.comet.CometConf.{COMET_ICEBERG_WRITE_REPORT_DIR, COMET_METRICS_ENABLED, COMET_ONHEAP_ENABLED} import org.apache.comet.CometKryoRegistrator import org.apache.comet.annotation.Public +import org.apache.comet.iceberg.IcebergWriteReportListener /** * Comet driver plugin. This class is loaded by Spark's plugin framework. It will be instantiated @@ -71,6 +72,7 @@ class CometDriverPlugin extends DriverPlugin with Logging { // Register Comet metrics CometDriverPlugin.registerCometMetrics(sc) + CometDriverPlugin.registerIcebergWriteReport(sc.conf) CometDriverPlugin.warnIfExecutorMemoryOverheadUnset(sc.getConf) @@ -186,27 +188,39 @@ object CometDriverPlugin extends Logging { COMET_METRICS_ENABLED.key, COMET_METRICS_ENABLED.defaultValue.get)) { sc.env.metricsSystem.registerSource(CometSource) - - val listenerKey = "spark.sql.queryExecutionListeners" - val listenerClass = "org.apache.comet.CometMetricsListener" - val listeners = sc.conf.get(listenerKey, "") - if (listeners.isEmpty) { - logInfo(s"Setting $listenerKey=$listenerClass") - sc.conf.set(listenerKey, listenerClass) - } else { - val currentListeners = listeners.split(",").map(_.trim) - if (!currentListeners.contains(listenerClass)) { - val newValue = s"$listeners,$listenerClass" - logInfo(s"Setting $listenerKey=$newValue") - sc.conf.set(listenerKey, newValue) - } - } + registerQueryExecutionListener(sc.conf, "org.apache.comet.CometMetricsListener") } else { logInfo( "Comet metrics reporting is disabled. Set spark.comet.metrics.enabled=true to enable.") } } + // Test-only: see COMET_ICEBERG_WRITE_REPORT_DIR. The value may come from the environment, which + // lets the Iceberg Spark test jobs turn the report on without changing the Iceberg diffs. + def registerIcebergWriteReport(conf: SparkConf): Unit = { + if (conf + .get(COMET_ICEBERG_WRITE_REPORT_DIR.key, COMET_ICEBERG_WRITE_REPORT_DIR.defaultValue.get) + .nonEmpty) { + registerQueryExecutionListener(conf, classOf[IcebergWriteReportListener].getName) + } + } + + private def registerQueryExecutionListener(conf: SparkConf, listenerClass: String): Unit = { + val listenerKey = "spark.sql.queryExecutionListeners" + val listeners = conf.get(listenerKey, "") + if (listeners.isEmpty) { + logInfo(s"Setting $listenerKey=$listenerClass") + conf.set(listenerKey, listenerClass) + } else { + val currentListeners = listeners.split(",").map(_.trim) + if (!currentListeners.contains(listenerClass)) { + val newValue = s"$listeners,$listenerClass" + logInfo(s"Setting $listenerKey=$newValue") + conf.set(listenerKey, newValue) + } + } + } + def registerCometSessionExtension(conf: SparkConf): Unit = { val extensionKey = StaticSQLConf.SPARK_SESSION_EXTENSIONS.key val extensionClass = classOf[CometSparkSessionExtensions].getName diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index 9c082912840..c5eae7fe483 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -20,6 +20,7 @@ package org.apache.comet import java.io.File +import java.nio.file.Files import java.sql.Timestamp import java.util.concurrent.{CountDownLatch, TimeUnit} @@ -29,7 +30,10 @@ import scala.concurrent.ExecutionContext.Implicits.global import scala.concurrent.duration.DurationInt import scala.jdk.CollectionConverters._ -import org.apache.spark.{SparkConf, Success} +import org.json4s.{DefaultFormats, Formats} +import org.json4s.jackson.JsonMethods.parse + +import org.apache.spark.{CometListenerBusUtils, SparkConf, Success} import org.apache.spark.scheduler.{SparkListener, SparkListenerTaskEnd} import org.apache.spark.sql.CometTestBase import org.apache.spark.sql.DataFrame @@ -42,7 +46,7 @@ import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types.{DoubleType, IntegerType, StringType, StructField, StructType} import org.apache.comet.CometSparkSessionExtensions.{isSpark35Plus, isSpark41Plus} -import org.apache.comet.iceberg.IcebergReflection +import org.apache.comet.iceberg.{IcebergReflection, IcebergWriteReportListener} private case class WriteSnapshot(snapshotDelta: Long, plans: Seq[SparkPlan]) @@ -2513,6 +2517,54 @@ class CometIcebergWriteActionSuite } } + test("write report records which writer ran each Iceberg write") { + assumeNativeAcceleration() + withIcebergCatalog { warehouseDir => + createTable(warehouseDir, "report_parquet", partitionSpec = "PARTITIONED BY (region)") + createTable( + warehouseDir, + "report_orc", + partitionSpec = "", + properties = Some("'write.format.default'='orc'")) + + withTempIcebergDir { reportDir => + val listener = new IcebergWriteReportListener( + new SparkConf() + .set(CometConf.COMET_ICEBERG_WRITE_REPORT_DIR.key, reportDir.getAbsolutePath)) + spark.listenerManager.register(listener) + try { + withNativeEnabled { + // Collecting the result runs a second query over the command's result; the write + // must still be reported once. + spark + .sql(s"INSERT INTO $catalog.$ns.report_parquet VALUES (1, 'us-east', 1.5)") + .collect() + spark.sql(s"INSERT INTO $catalog.$ns.report_orc VALUES (2, 'eu', 2.5)") + } + withSQLConf(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "false") { + spark.sql(s"INSERT INTO $catalog.$ns.report_parquet VALUES (3, 'eu', 3.5)") + } + CometListenerBusUtils.waitUntilEmpty(spark.sparkContext) + } finally { + spark.listenerManager.unregister(listener) + } + + implicit val formats: Formats = DefaultFormats + val writes = reportDir + .listFiles() + .toSeq + .flatMap(f => Files.readAllLines(f.toPath).asScala) + .map(parse(_)) + assert( + writes.map(w => (w \ "writer").extract[String]) == Seq("native", "jvm", "spark"), + writes.mkString("\n")) + assert((writes(1) \ "reasons").extract[Seq[String]].exists(_.contains("only parquet"))) + assert((writes(2) \ "node").extract[String] == "AppendData") + assert(writes.forall(w => !(w \ "failed").extract[Boolean])) + } + } + } + private def assertNativeWriteEngages(tableName: String, expectedIds: Seq[Int])( action: => Unit): Unit = { val snapshot = withNativeEnabled { captureWrite(tableName)(action) } diff --git a/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala b/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala index 278c281b1d9..b5652c0f671 100644 --- a/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala +++ b/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala @@ -25,7 +25,7 @@ import org.apache.logging.log4j.Level import org.apache.spark.sql.{CometTestBase, SaveMode} import org.apache.spark.sql.internal.StaticSQLConf -import org.apache.comet.COMET_VERSION +import org.apache.comet.{COMET_VERSION, CometConf} class CometPluginsSuite extends CometTestBase { override protected def sparkConf: SparkConf = { @@ -86,6 +86,22 @@ class CometPluginsSuite extends CometTestBase { } } + test("Iceberg write report listener is registered only when a report directory is set") { + val listenerKey = "spark.sql.queryExecutionListeners" + val listenerClass = "org.apache.comet.iceberg.IcebergWriteReportListener" + + val unset = new SparkConf() + CometDriverPlugin.registerIcebergWriteReport(unset) + assert(!unset.contains(listenerKey)) + + val set = new SparkConf() + .set(CometConf.COMET_ICEBERG_WRITE_REPORT_DIR.key, "/tmp/report") + .set(listenerKey, "foo") + CometDriverPlugin.registerIcebergWriteReport(set) + CometDriverPlugin.registerIcebergWriteReport(set) + assert(set.get(listenerKey) == s"foo,$listenerClass") + } + test("Comet version is exposed as a Spark config") { // The driver plugin sets spark.comet.version, which is then visible both on the SparkContext // conf and through the session runtime config (SET / spark.conf.get). From b4dd3d0c2628335fa0eb60a628fd0c625b902d88 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Fri, 25 Sep 2026 16:56:21 -0600 Subject: [PATCH 2/2] test: report streaming and Spark 3.4 CTAS/RTAS Iceberg writes, and count only each shard's latest attempt The write report missed two kinds of Iceberg write that Spark runs itself: - A streaming micro-batch runs as `WriteToDataSourceV2Exec` over a `MicroBatchWrite`, not as a `V2ExistingTableWriteExec`. - On Spark 3.4, CTAS and RTAS write the table from the create or replace exec through `TableWriteExecHelper.writeWithV2`. Spark 3.5+ runs that write as a nested append or overwrite, which the listener already sees, so `IcebergTableAsSelectShim` recognizes the execs on 3.4 only. Both went unreported, so the Spark count and the total left them out and the native share read high. The summarizer took each shard's latest attempt from its report files, so a retry that recorded no writes was hidden behind the earlier attempt's records. It now takes the latest attempt from the shard artifact directories and names a shard whose latest attempt recorded nothing. `dev/ci/test-summarize-iceberg-writes.py` covers that and runs in Preflight. --- .github/workflows/ci.yml | 3 + dev/ci/summarize-iceberg-writes.py | 62 +++++--- dev/ci/test-summarize-iceberg-writes.py | 105 +++++++++++++ .../contributor-guide/iceberg-spark-tests.md | 12 +- .../iceberg/IcebergWriteReportListener.scala | 22 ++- .../iceberg/IcebergTableAsSelectShim.scala | 40 +++++ .../iceberg/IcebergTableAsSelectShim.scala | 34 +++++ .../iceberg/IcebergTableAsSelectShim.scala | 34 +++++ .../comet/CometIcebergWriteActionSuite.scala | 139 ++++++++++++++---- 9 files changed, 394 insertions(+), 57 deletions(-) create mode 100644 dev/ci/test-summarize-iceberg-writes.py create mode 100644 spark/src/main/spark-3.4/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala create mode 100644 spark/src/main/spark-3.5/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala create mode 100644 spark/src/main/spark-4.x/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 7f87d591995..d153ff3c63b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -186,6 +186,9 @@ jobs: - name: Check Iceberg shard inventory validation run: python3 dev/ci/test-iceberg-shards.py + - name: Check Iceberg write report summary + run: python3 dev/ci/test-summarize-iceberg-writes.py + - name: Check CI config invariants run: python3 dev/ci/check-ci-config.py diff --git a/dev/ci/summarize-iceberg-writes.py b/dev/ci/summarize-iceberg-writes.py index 64cc2b94516..c94638b4e12 100644 --- a/dev/ci/summarize-iceberg-writes.py +++ b/dev/ci/summarize-iceberg-writes.py @@ -27,8 +27,11 @@ Files under a directory named like a shard artifact (...-shard-N-attempt-M) are counted only for the latest attempt of each shard, so a rerun of failed -jobs does not count a shard twice. The summary is also appended to -$GITHUB_STEP_SUMMARY when that is set. It never fails the job. +jobs does not count a shard twice. The latest attempt is the newest such +directory, whether or not it holds any report files, so a rerun that recorded +no writes is reported as missing rather than replaced by an earlier attempt. +The summary is also appended to $GITHUB_STEP_SUMMARY when that is set. It +never fails the job. """ import argparse @@ -48,40 +51,59 @@ TOP_REASONS = 20 +def shard_attempt(path): + """The (shard, attempt) of the shard artifact directory holding path, or None.""" + for part in path.parts: + match = SHARD_ATTEMPT.search(part) + if match: + return int(match.group(1)), int(match.group(2)) + return None + + def report_files(roots): - """Every report file under roots, keeping only the latest attempt of each shard.""" + """Every report file under roots, keeping only the latest attempt of each shard. + + Returns the files and the (shard, attempt) pairs whose latest attempt holds no report file. + The latest attempt comes from the artifact directories rather than the report files, because + an attempt that recorded no writes still uploads its shard inventory and test reports. + """ latest = {} files = [] - for root in roots: - for path in sorted(Path(root).rglob("*.jsonl")): - match = next( - (SHARD_ATTEMPT.search(part) for part in path.parts if SHARD_ATTEMPT.search(part)), - None, - ) - if match: - shard, attempt = int(match.group(1)), int(match.group(2)) - latest[shard] = max(latest.get(shard, 0), attempt) - files.append((path, (shard, attempt))) - else: - files.append((path, None)) - return [path for path, key in files if key is None or latest[key[0]] == key[1]] + for root in map(Path, roots): + for path in [root, *sorted(root.rglob("*"))]: + key = shard_attempt(path) + if key: + latest[key[0]] = max(latest.get(key[0], 0), key[1]) + if path.suffix == ".jsonl" and path.is_file(): + files.append((path, key)) + kept = [(path, key) for path, key in files if key is None or latest[key[0]] == key[1]] + reported = {key for _, key in kept} + missing = sorted(key for key in latest.items() if key not in reported) + return [path for path, _ in kept], missing def load(roots): + files, missing = report_files(roots) writes = [] - for path in report_files(roots): + for path in files: for line in path.read_text(encoding="utf-8").splitlines(): if line.strip(): writes.append(json.loads(line)) - return writes + return writes, missing def cell(text): return " ".join(text.split()).replace("|", "\\|") -def summarize(title, writes): +def summarize(title, writes, missing=()): lines = [f"### Iceberg writes: {title}", ""] + for shard, attempt in missing: + lines += [ + f"Shard {shard} recorded no Iceberg writes in its latest attempt ({attempt}), " + "so none of its writes are counted below.", + "", + ] if not writes: lines.append( "No Iceberg writes were recorded. Either the target ran none or " @@ -137,7 +159,7 @@ def main(): parser.add_argument("roots", nargs="+", help="directories holding the report files") args = parser.parse_args() - summary = summarize(args.title, load(args.roots)) + summary = summarize(args.title, *load(args.roots)) print(summary) step_summary = os.environ.get("GITHUB_STEP_SUMMARY") if step_summary: diff --git a/dev/ci/test-summarize-iceberg-writes.py b/dev/ci/test-summarize-iceberg-writes.py new file mode 100644 index 00000000000..cce4ada6630 --- /dev/null +++ b/dev/ci/test-summarize-iceberg-writes.py @@ -0,0 +1,105 @@ +#!/usr/bin/env python3 +# +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Fast regression tests for the Iceberg write report summary.""" + +import importlib.util +import json +from pathlib import Path +import tempfile +import unittest + + +SPEC = importlib.util.spec_from_file_location( + "summarize_iceberg_writes", Path(__file__).with_name("summarize-iceberg-writes.py")) +SUMMARIZE = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(SUMMARIZE) + +ARTIFACT = "iceberg-spark-1.11.0-spark-4.1.3-scala-2.13-jdk17-shard-{}-attempt-{}" + + +def records(*writers): + return "".join( + json.dumps({"writer": w, "node": "AppendData", "reasons": [], "failed": False}) + "\n" + for w in writers) + + +class SummarizeIcebergWritesTest(unittest.TestCase): + def setUp(self): + self.temp = tempfile.TemporaryDirectory(prefix="comet-iceberg-writes-test-") + self.addCleanup(self.temp.cleanup) + self.root = Path(self.temp.name) + + def attempt(self, shard, attempt, *writers): + """A shard attempt's artifact as the coverage job downloads it. + + Every attempt uploads its test reports. Only an attempt that recorded writes has a + report file. + """ + artifact = self.root / ARTIFACT.format(shard, attempt) + reports = artifact / "build/test-results/test" + reports.mkdir(parents=True) + (reports / "TEST-org.example.TestFixture.xml").write_text("") + if writers: + writes = artifact / "build/comet-iceberg-writes" + writes.mkdir(parents=True) + (writes / "iceberg-writes-fixture.jsonl").write_text(records(*writers)) + + def load(self): + writes, missing = SUMMARIZE.load([self.root]) + return sorted(w["writer"] for w in writes), missing + + def test_every_shard_counts_once(self): + self.attempt(1, 1, "native") + self.attempt(2, 1, "jvm", "spark") + self.assertEqual(self.load(), (["jvm", "native", "spark"], [])) + + def test_latest_attempt_replaces_an_earlier_one(self): + self.attempt(1, 1, "native", "native") + self.attempt(1, 2, "jvm") + self.assertEqual(self.load(), (["jvm"], [])) + + def test_retry_without_writes_is_reported_missing_instead_of_stale(self): + self.attempt(1, 1, "native") + self.attempt(1, 2) + self.attempt(2, 1, "spark") + self.assertEqual(self.load(), (["spark"], [(1, 2)])) + summary = SUMMARIZE.summarize("fixture", *SUMMARIZE.load([self.root])) + self.assertIn("Shard 1 recorded no Iceberg writes in its latest attempt (2)", summary) + self.assertIn("| Comet native writer | 0 | 0.0% |", summary) + self.assertIn("| Spark V2 write, not planned by Comet's split operator | 1 | 100.0% |", + summary) + + def test_root_that_is_itself_a_shard_artifact(self): + self.attempt(1, 3) + writes, missing = SUMMARIZE.load([self.root / ARTIFACT.format(1, 3)]) + self.assertEqual((writes, missing), ([], [(1, 3)])) + + def test_files_outside_shard_artifacts_always_count(self): + # A shard job and dev/local-ci.sh summarize their own report directory directly. + (self.root / "iceberg-writes-local.jsonl").write_text(records("native", "jvm")) + self.assertEqual(self.load(), (["jvm", "native"], [])) + + def test_no_writes_at_all(self): + summary = SUMMARIZE.summarize("fixture", *SUMMARIZE.load([self.root / "missing"])) + self.assertIn("No Iceberg writes were recorded", summary) + + +if __name__ == "__main__": + unittest.main() diff --git a/docs/source/contributor-guide/iceberg-spark-tests.md b/docs/source/contributor-guide/iceberg-spark-tests.md index 07358aa7b85..0db54feeffe 100644 --- a/docs/source/contributor-guide/iceberg-spark-tests.md +++ b/docs/source/contributor-guide/iceberg-spark-tests.md @@ -170,12 +170,14 @@ records one of three writers: - `native`: Comet's native writer (`CometIcebergWriteExec`). - `jvm`: Comet's split operator planned the write but kept Iceberg's JVM writer (`IcebergWriteExec`). The line includes the reasons Comet recorded for not converting it. -- `spark`: Spark's own V2 write operator ran the write, for example `WriteDelta` for merge-on-read, - so Comet's split operator never saw it. +- `spark`: Spark's own V2 write operator ran the write, so Comet's split operator never saw it. Examples + are `WriteDelta` for merge-on-read, `WriteToDataSourceV2` for a streaming micro-batch, and on Spark + 3.4 the CTAS and RTAS execs, which write the table themselves. `dev/ci/summarize-iceberg-writes.py` turns these records into a table on the job's summary page. It shows the count and share of each writer, the most common fallback reasons, and the Spark write operators. Each shard and the extensions job gets its own table. The shard coverage job adds one for -all shards together. The raw records are uploaded with the job's other reports. The summary never -fails a job. `dev/local-ci.sh iceberg` prints the same summary after each shard and after the -extensions target. +all shards together, counting only the latest attempt of each shard. A shard whose latest attempt +recorded no writes is named above the table rather than counted from an earlier attempt. The raw +records are uploaded with the job's other reports. The summary never fails a job. +`dev/local-ci.sh iceberg` prints the same summary after each shard and after the extensions target. diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala index c1e9538732c..d5f3663eb1b 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteReportListener.scala @@ -32,9 +32,11 @@ import org.json4s.jackson.JsonMethods._ import org.apache.spark.SparkConf import org.apache.spark.internal.Logging import org.apache.spark.sql.comet.{CometIcebergWriteExec, IcebergWriteExec} +import org.apache.spark.sql.connector.write.BatchWrite import org.apache.spark.sql.execution.{CommandResultExec, QueryExecution, SparkPlan} import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanExec, QueryStageExec} -import org.apache.spark.sql.execution.datasources.v2.V2ExistingTableWriteExec +import org.apache.spark.sql.execution.datasources.v2.{V2ExistingTableWriteExec, WriteToDataSourceV2Exec} +import org.apache.spark.sql.execution.streaming.sources.MicroBatchWrite import org.apache.spark.sql.util.QueryExecutionListener import org.apache.comet.CometConf.COMET_ICEBERG_WRITE_REPORT_DIR @@ -111,9 +113,23 @@ object IcebergWriteReportListener { case w: IcebergWriteExec => val reasons = w.getTagValue(CometExplainInfo.FALLBACK_REASONS).getOrElse(Set.empty) Seq(IcebergWrite(Jvm, w.nodeName, reasons.toSeq.sorted)) - case w: V2ExistingTableWriteExec - if w.write.getClass.getName.startsWith("org.apache.iceberg.") => + case w: V2ExistingTableWriteExec if isIceberg(w.write) => + Seq(IcebergWrite(Spark, w.nodeName, Nil)) + // A streaming micro-batch, which Comet's split operator never plans. + case w: WriteToDataSourceV2Exec if isIcebergMicroBatch(w.batchWrite) => + Seq(IcebergWrite(Spark, w.nodeName, Nil)) + // Spark 3.4 writes a CTAS or RTAS from the create or replace exec itself. Later versions run + // that write as a nested append or overwrite, which this listener sees as a query of its own. + case w if IcebergTableAsSelectShim.writeCatalog(w).exists(isIceberg) => Seq(IcebergWrite(Spark, w.nodeName, Nil)) case p => p.children.flatMap(writes) } + + private def isIcebergMicroBatch(write: BatchWrite): Boolean = write match { + case m: MicroBatchWrite => isIceberg(m.writeSupport) + case _ => false + } + + private def isIceberg(obj: AnyRef): Boolean = + obj.getClass.getName.startsWith("org.apache.iceberg.") } diff --git a/spark/src/main/spark-3.4/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala b/spark/src/main/spark-3.4/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala new file mode 100644 index 00000000000..4d66a374663 --- /dev/null +++ b/spark/src/main/spark-3.4/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala @@ -0,0 +1,40 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.iceberg + +import org.apache.spark.sql.connector.catalog.TableCatalog +import org.apache.spark.sql.execution.SparkPlan +import org.apache.spark.sql.execution.datasources.v2.{AtomicCreateTableAsSelectExec, AtomicReplaceTableAsSelectExec, CreateTableAsSelectExec, ReplaceTableAsSelectExec} + +/** + * Spark 3.4: CTAS and RTAS write the table from the create or replace exec itself, through + * `TableWriteExecHelper.writeWithV2`, so the write never appears as a write node of its own. + */ +private[iceberg] object IcebergTableAsSelectShim { + + /** The catalog `plan` writes through, when `plan` is a CTAS or RTAS exec. */ + def writeCatalog(plan: SparkPlan): Option[TableCatalog] = plan match { + case p: CreateTableAsSelectExec => Some(p.catalog) + case p: AtomicCreateTableAsSelectExec => Some(p.catalog) + case p: ReplaceTableAsSelectExec => Some(p.catalog) + case p: AtomicReplaceTableAsSelectExec => Some(p.catalog) + case _ => None + } +} diff --git a/spark/src/main/spark-3.5/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala b/spark/src/main/spark-3.5/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala new file mode 100644 index 00000000000..ef02700a9c8 --- /dev/null +++ b/spark/src/main/spark-3.5/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala @@ -0,0 +1,34 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.iceberg + +import org.apache.spark.sql.connector.catalog.TableCatalog +import org.apache.spark.sql.execution.SparkPlan + +/** + * Spark 3.5+: CTAS and RTAS run their write as a nested `AppendData` or `OverwriteByExpression` + * query, which is planned and reported like any other write, so no create or replace exec writes + * a table itself. + */ +private[iceberg] object IcebergTableAsSelectShim { + + /** The catalog `plan` writes through, when `plan` is a CTAS or RTAS exec. */ + def writeCatalog(plan: SparkPlan): Option[TableCatalog] = None +} diff --git a/spark/src/main/spark-4.x/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala b/spark/src/main/spark-4.x/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala new file mode 100644 index 00000000000..ef02700a9c8 --- /dev/null +++ b/spark/src/main/spark-4.x/org/apache/comet/iceberg/IcebergTableAsSelectShim.scala @@ -0,0 +1,34 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.iceberg + +import org.apache.spark.sql.connector.catalog.TableCatalog +import org.apache.spark.sql.execution.SparkPlan + +/** + * Spark 3.5+: CTAS and RTAS run their write as a nested `AppendData` or `OverwriteByExpression` + * query, which is planned and reported like any other write, so no create or replace exec writes + * a table itself. + */ +private[iceberg] object IcebergTableAsSelectShim { + + /** The catalog `plan` writes through, when `plan` is a CTAS or RTAS exec. */ + def writeCatalog(plan: SparkPlan): Option[TableCatalog] = None +} diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index ff332859b3f..01a8188f29b 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -47,6 +47,7 @@ import org.apache.spark.sql.connector.write.{BatchWrite, DataWriterFactory, Phys import org.apache.spark.sql.execution.{ColumnarToRowTransition, LeafExecNode, SparkPlan} import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.streaming.Trigger import org.apache.spark.sql.types.{DoubleType, IntegerType, StringType, StructField, StructType} import org.apache.comet.CometSparkSessionExtensions.{isSpark35Plus, isSpark41Plus} @@ -54,6 +55,12 @@ import org.apache.comet.iceberg.{IcebergReflection, IcebergWriteReportListener} private case class WriteSnapshot(snapshotDelta: Long, plans: Seq[SparkPlan]) +private case class ReportedWrite( + writer: String, + node: String, + reasons: Seq[String], + failed: Boolean) + class CometIcebergWriteActionSuite extends CometTestBase with AdaptiveSparkPlanHelper @@ -2643,42 +2650,116 @@ class CometIcebergWriteActionSuite partitionSpec = "", properties = Some("'write.format.default'='orc'")) - withTempIcebergDir { reportDir => - val listener = new IcebergWriteReportListener( - new SparkConf() - .set(CometConf.COMET_ICEBERG_WRITE_REPORT_DIR.key, reportDir.getAbsolutePath)) - spark.listenerManager.register(listener) - try { + val writes = reportedWrites { + withNativeEnabled { + // Collecting the result runs a second query over the command's result; the write + // must still be reported once. + spark + .sql(s"INSERT INTO $catalog.$ns.report_parquet VALUES (1, 'us-east', 1.5)") + .collect() + spark.sql(s"INSERT INTO $catalog.$ns.report_orc VALUES (2, 'eu', 2.5)") + } + withSQLConf(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "false") { + spark.sql(s"INSERT INTO $catalog.$ns.report_parquet VALUES (3, 'eu', 3.5)") + } + } + assert(writes.map(_.writer) == Seq("native", "jvm", "spark"), writes.mkString("\n")) + assert(writes(1).reasons.exists(_.contains("only parquet"))) + assert(writes(2).node == "AppendData") + assert(writes.forall(!_.failed)) + } + } + + test("write report records a CTAS and an RTAS once each") { + assumeNativeAcceleration() + withIcebergCatalog { _ => + val writes = reportedWrites { + withNativeEnabled { + spark.sql(s""" + CREATE TABLE $catalog.$ns.report_ctas USING iceberg AS + SELECT * FROM VALUES (1, 'us', 1.0), (2, 'eu', 2.0) AS t(id, region, amount) + """) + spark.sql(s""" + REPLACE TABLE $catalog.$ns.report_ctas USING iceberg AS + SELECT * FROM VALUES (3, 'us', 3.0) AS t(id, region, amount) + """) + } + } + // Spark 3.5+ runs each write as a nested append or overwrite, which Comet's split operator + // plans. Spark 3.4 writes from the create and replace execs, which it never sees. + val expected = + if (isSpark35Plus) Seq("native" -> "CometIcebergWrite", "native" -> "CometIcebergWrite") + else Seq("spark" -> "AtomicCreateTableAsSelect", "spark" -> "AtomicReplaceTableAsSelect") + assert(writes.map(w => w.writer -> w.node) == expected, writes.mkString("\n")) + assert(writes.forall(!_.failed)) + assertRows("report_ctas", Seq(3)) + } + } + + test("write report records a streaming micro-batch write") { + assumeNativeAcceleration() + withIcebergCatalog { warehouseDir => + createTable(warehouseDir, "report_stream", partitionSpec = "") + withTempIcebergDir { dir => + val source = new File(dir, "source").getAbsolutePath + val session = spark + import session.implicits._ + Seq((1, "us", 1.0), (2, "eu", 2.0)).toDF("id", "region", "amount").write.parquet(source) + val schema = spark.table(s"$catalog.$ns.report_stream").schema + + val writes = reportedWrites { withNativeEnabled { - // Collecting the result runs a second query over the command's result; the write - // must still be reported once. - spark - .sql(s"INSERT INTO $catalog.$ns.report_parquet VALUES (1, 'us-east', 1.5)") - .collect() - spark.sql(s"INSERT INTO $catalog.$ns.report_orc VALUES (2, 'eu', 2.5)") + spark.readStream + .schema(schema) + .parquet(source) + .writeStream + .format("iceberg") + .outputMode("append") + .trigger(Trigger.AvailableNow()) + .option("checkpointLocation", new File(dir, "checkpoint").getAbsolutePath) + .toTable(s"$catalog.$ns.report_stream") + .awaitTermination() } - withSQLConf(CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "false") { - spark.sql(s"INSERT INTO $catalog.$ns.report_parquet VALUES (3, 'eu', 3.5)") - } - CometListenerBusUtils.waitUntilEmpty(spark.sparkContext) - } finally { - spark.listenerManager.unregister(listener) } - - implicit val formats: Formats = DefaultFormats - val writes = reportDir - .listFiles() - .toSeq - .flatMap(f => Files.readAllLines(f.toPath).asScala) - .map(parse(_)) assert( - writes.map(w => (w \ "writer").extract[String]) == Seq("native", "jvm", "spark"), + writes.map(w => w.writer -> w.node) == Seq("spark" -> "WriteToDataSourceV2"), writes.mkString("\n")) - assert((writes(1) \ "reasons").extract[Seq[String]].exists(_.contains("only parquet"))) - assert((writes(2) \ "node").extract[String] == "AppendData") - assert(writes.forall(w => !(w \ "failed").extract[Boolean])) + assert(writes.forall(!_.failed)) + assertRows("report_stream", Seq(1, 2)) + } + } + } + + /** The Iceberg writes `IcebergWriteReportListener` records while `action` runs, in order. */ + private def reportedWrites(action: => Unit): Seq[ReportedWrite] = { + var writes = Seq.empty[ReportedWrite] + withTempIcebergDir { reportDir => + val listener = new IcebergWriteReportListener( + new SparkConf() + .set(CometConf.COMET_ICEBERG_WRITE_REPORT_DIR.key, reportDir.getAbsolutePath)) + spark.listenerManager.register(listener) + try { + action + CometListenerBusUtils.waitUntilEmpty(spark.sparkContext) + } finally { + spark.listenerManager.unregister(listener) } + + implicit val formats: Formats = DefaultFormats + writes = reportDir + .listFiles() + .toSeq + .flatMap(f => Files.readAllLines(f.toPath).asScala) + .map { line => + val w = parse(line) + ReportedWrite( + (w \ "writer").extract[String], + (w \ "node").extract[String], + (w \ "reasons").extract[Seq[String]], + (w \ "failed").extract[Boolean]) + } } + writes } private def assertNativeWriteEngages(tableName: String, expectedIds: Seq[Int])(