andygrove commented on code in PR #6664:
URL: https://github.com/apache/datafusion-comet/pull/6664#discussion_r4201730609
##########
spark/src/main/scala/org/apache/comet/iceberg/IcebergWriteStrategy.scala:
##########
@@ -38,9 +38,13 @@ case class IcebergWriteStrategy(session: SparkSession)
extends SparkStrategy {
override def apply(plan: LogicalPlan): Seq[SparkPlan] = {
val conf = session.sessionState.conf
+ // The native write flag plans the split operator on its own. The split
flag plans it with
+ // the native writer off, which only tests do.
+ val splitEnabled = CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.get(conf)
||
+ CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.get(conf)
// Planner strategies run whether or not Comet is enabled, so check it
here too: with Comet
// off, Spark must plan its own V2 write operator.
- if (!isCometLoaded(conf) ||
!CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.get(conf)) {
+ if (!isCometLoaded(conf) || !splitEnabled) {
Review Comment:
Agreed, the split plan with iceberg-java's writer gives a scan- or
shuffle-only deployment nothing, so it shouldn't change their plans. The
strategy now also checks `spark.comet.exec.enabled`, and the
`spark.comet.enabled=false` test runs for both settings. The fallback list in
`iceberg-writes.md` names `spark.comet.exec.enabled=false` and plan-only mode
now too.
##########
docs/source/user-guide/latest/iceberg-writes.md:
##########
@@ -92,10 +96,8 @@
spark.sql.catalog.<name>=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.<name>.type=hadoop # or hive / glue
/ rest / ...
spark.sql.catalog.<name>.warehouse=...
-# Split-operator plan (experimental, off by default)
-spark.comet.write.iceberg.splitOperator.enabled=true
-
-# Native Parquet writer (experimental, off by default; requires the split plan)
+# Split-operator plan and native Parquet writer (on by default since Comet
1.2.0; false plans
+# Spark's own write operator)
spark.comet.write.iceberg.enabled=true
Review Comment:
Good point. I dropped the line from the block, and the text above it now
says the setting is on by default and that `false` plans Spark's own operator.
The local table scan setting is the only write-related line left in the example.
##########
docs/source/user-guide/latest/migration-guide.md:
##########
@@ -83,6 +83,28 @@ it. `spark.comet.convert.oneRowRelation.enabled` is on by
default, so such a que
without either setting, and logs no warning. The list remains the way to
convert other leaf
operators, such as the scan of a Data Source V2 connector.
+### Iceberg Writes
+
+`spark.comet.write.iceberg.enabled` now defaults to `true`. Comet plans an
Iceberg `INSERT INTO`,
Review Comment:
Thanks, I'd missed that the rename landed after the 1.1 branch was cut. The
entry now says 1.1.0 named the setting `spark.comet.iceberg.write.enabled`,
that it was a testing setting Comet now ignores, and that a deployment that set
it to `false` needs `spark.comet.write.iceberg.enabled=false` to keep
iceberg-java's writer.
##########
spark/src/main/scala/org/apache/comet/CometConf.scala:
##########
@@ -125,21 +125,24 @@ object CometConf extends ShimCometConf {
conf("spark.comet.write.iceberg.splitOperator.enabled")
.category(CATEGORY_TESTING)
.doc(
- "Whether to rewrite Iceberg V2 writes from Spark's combined V2
write/commit operator " +
- "into Comet's two-operator shape: a file writer exec (inside AQE)
and a committer " +
- "(outside AQE).")
+ "Whether to plan Iceberg writes as Comet's file writer and committer
even when " +
+ "`spark.comet.write.iceberg.enabled` is false, so that Iceberg's own
writer writes " +
+ "every data file inside Comet's plan. Used by tests to compare the
two writers " +
+ "under the same plan.")
.booleanConf
.createWithDefault(false)
val COMET_ICEBERG_NATIVE_WRITE_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.write.iceberg.enabled")
- .category(CATEGORY_TESTING)
+ .category(CATEGORY_EXEC)
.doc(
- "Whether to delegate the executor-side Parquet write to Comet's native
(iceberg-rust) " +
- "writer when the table's properties allow it. Requires " +
- "`spark.comet.write.iceberg.splitOperator.enabled = true`. Off by
default.")
+ "Whether Comet plans Iceberg writes and writes their data files
natively. Comet " +
+ "replaces Spark's combined V2 write operator with a file writer
(inside AQE) under a " +
+ "committer (outside AQE), and writes the data files of each eligible
write with its " +
+ "native (iceberg-rust) writer. Other writes use Iceberg's own
writer, and Iceberg " +
+ "commits every write. Set this to false to plan Spark's own V2 write
operator.")
.booleanConf
- .createWithDefault(false)
+ .createWithDefault(true)
Review Comment:
This is #5719, and #5957 fixed it by making the reverted stage restore an
`IcebergWriteExec` instead of dropping the write. This branch now includes
#5957. I added `transition-heavy fallback keeps the write when no flag is set`,
which writes from a Parquet source with transition reversion on,
`maxTransitions=0` and no write flag set, under both AQE settings. Without
#5957 it fails with the same `EOFException` at commit, and with it the write
goes through iceberg-java's writer and commits the expected rows.
--
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]