parthchandra commented on code in PR #6355:
URL: https://github.com/apache/datafusion-comet/pull/6355#discussion_r4190659247
##########
native/core/src/execution/operators/parquet_writer.rs:
##########
@@ -402,6 +435,25 @@ impl ParquetWriterExec {
})?;
Ok(ParquetWriter::LocalFile(writer))
}
+ } else if Self::is_s3_destination(output_file_path,
object_store_options)? {
+ // The store is resolved exactly as the Parquet scan resolves it,
including the
+ // `fs.s3a.*` endpoint, region and credential settings, but under
the write access mode.
+ let (store, path) = object_store_for_write(output_file_path,
object_store_options)
+ .map_err(|e| {
+ DataFusionError::Execution(format!(
+ "Failed to prepare object store for '{}': {}",
+ output_file_path, e
+ ))
+ })?;
+ let writer =
+ AsyncArrowWriter::try_new(BufWriter::new(store, path), schema,
Some(props))
Review Comment:
Can you confirm what the buffer size here will be? This is probably not
tracked by the Comet allocator and could potentially be quite large.
##########
spark/src/test/scala/org/apache/comet/parquet/ParquetWriteToS3Suite.scala:
##########
@@ -0,0 +1,254 @@
+/*
+ * 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.parquet
+
+import java.net.URI
+
+import org.apache.hadoop.fs.{FileStatus, Path}
+import org.apache.spark.SparkConf
+import org.apache.spark.sql.{DataFrame, Row, SaveMode}
+
+import org.apache.comet.{CometConf, CometS3TestBase}
+import org.apache.comet.CometSparkSessionExtensions.isSpark40Plus
+import org.apache.comet.cloud.s3.{CometS3AccessMode,
MinioCometS3CredentialProvider}
+
+import software.amazon.awssdk.auth.credentials.{AwsBasicCredentials,
StaticCredentialsProvider}
+import software.amazon.awssdk.regions.Region
+import software.amazon.awssdk.services.s3.S3Client
+import software.amazon.awssdk.services.s3.model.HeadObjectRequest
+
+/**
+ * Native Parquet writes to S3, against MinIO. Which S3 destinations the
native writer accepts is
+ * covered without an S3 endpoint in `CometParquetWriterSuite`; this suite
checks that the ones it
+ * accepts are written where Spark's commit protocol expects them.
+ *
+ * A manual suite, like the other MinIO suites: it needs Docker, so CI does
not run it (see
+ * `dev/ci/check-suites.py`).
+ */
+class ParquetWriteToS3Suite extends CometParquetWriterTestBase with
CometS3TestBase {
+
+ import testImplicits._
+
+ override protected val testBucketName = "native-write-bucket"
+
+ /** A bucket whose native access goes through
[[MinioCometS3CredentialProvider]]. */
+ private val credentialBucket = "native-write-credential-bucket"
+
+ override protected def sparkConf: SparkConf = {
+ val conf = super.sparkConf
+ conf.set(
+
s"spark.hadoop.fs.s3a.bucket.$credentialBucket.comet.credential.provider.class",
+ classOf[MinioCometS3CredentialProvider].getName)
+ conf
+ }
+
+ override def beforeAll(): Unit = {
+ super.beforeAll()
+ MinioCometS3CredentialProvider.installCredentials(userName, password)
+ createBucketIfNotExists(credentialBucket)
+ }
+
+ private def s3Path(key: String, bucket: String = testBucketName): String =
+ s"s3a://$bucket/$key"
+
+ /**
+ * Persist `df` to S3 with Spark's writer and return a DataFrame that reads
it back, so that the
+ * write under test has the Comet scan below it that native writes require.
+ */
+ private def cometSource(df: DataFrame, name: String): DataFrame = {
+ val path = s3Path(s"sources/$name")
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ df.write.mode(SaveMode.Overwrite).parquet(path)
+ }
+ spark.read.parquet(path)
+ }
+
+ private def someRows(name: String): DataFrame =
+ cometSource((1 to 1000).map(i => (i, s"name_$i")).toDF("id",
"name").repartition(4), name)
Review Comment:
Can we add coverage with decimal, timestamp, and nested types (nested struct
with maps, arrays).
##########
spark/src/test/scala/org/apache/comet/parquet/ParquetWriteToS3Suite.scala:
##########
@@ -0,0 +1,254 @@
+/*
+ * 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.parquet
Review Comment:
The PR says these haven't been run. Can we run the tests manually and post
the results?
##########
native/core/src/execution/operators/parquet_writer.rs:
##########
@@ -550,47 +620,67 @@ impl ExecutionPlan for ParquetWriterExec {
let mut stream = input;
let mut total_rows = 0i64;
- while let Some(batch_result) = stream.try_next().await.transpose()
{
- let batch = batch_result?;
-
- // Track row count
- total_rows += batch.num_rows() as i64;
-
- // Rename columns in the batch to match output schema
- let renamed_batch = if !column_names.is_empty() {
- // Collection field IDs exist on the target schema, not on
arrays produced by
- // the placeholder Scan. Both schemas use the same
Catalyst data types, and
- // disabling field-name matching still recursively
validates nested nullability;
- // only nested field names and metadata are ignored.
- RecordBatch::try_new_with_options(
- Arc::clone(&schema_for_write),
- batch.columns().to_vec(),
-
&RecordBatchOptions::new().with_match_field_names(false),
- )
- .map_err(|e| {
- DataFusionError::Execution(format!("Failed to rename
batch columns: {}", e))
- })?
- } else {
- batch
- };
-
- writer.write(&renamed_batch).await.map_err(|e| {
- DataFusionError::Execution(format!("Failed to write batch:
{}", e))
- })?;
+ // A failure here leaves the file unfinished, so keep it rather
than return it at once:
+ // the writer has to discard what it may already have uploaded.
+ let written: Result<()> = async {
+ while let Some(batch_result) =
stream.try_next().await.transpose() {
+ let batch = batch_result?;
+
+ // Track row count
+ total_rows += batch.num_rows() as i64;
+
+ // Rename columns in the batch to match output schema
+ let renamed_batch = if !column_names.is_empty() {
+ // Collection field IDs exist on the target schema,
not on arrays produced
+ // by the placeholder Scan. Both schemas use the same
Catalyst data types,
+ // and disabling field-name matching still recursively
validates nested
+ // nullability; only nested field names and metadata
are ignored.
+ RecordBatch::try_new_with_options(
+ Arc::clone(&schema_for_write),
+ batch.columns().to_vec(),
+
&RecordBatchOptions::new().with_match_field_names(false),
+ )
+ .map_err(|e| {
+ DataFusionError::Execution(format!(
+ "Failed to rename batch columns: {}",
+ e
+ ))
+ })?
+ } else {
+ batch
+ };
+
+ writer.write(&renamed_batch).await.map_err(|e| {
+ DataFusionError::Execution(format!("Failed to write
batch: {}", e))
+ })?;
+ }
+ Ok(())
+ }
+ .await;
+ if let Err(e) = written {
+ if let Err(abort_error) = writer.abort().await {
+ log::warn!("Failed to abort the write of '{part_file}':
{abort_error}");
+ }
+ return Err(e);
}
- writer.close().await.map_err(|e| {
+ let uploaded_size = writer.close().await.map_err(|e| {
Review Comment:
This ultimately calls `complete` which might fail too, but in that case we
are not calling `abort`
##########
spark/src/main/scala/org/apache/comet/serde/operator/NativeWriteUtils.scala:
##########
@@ -172,19 +177,215 @@ object NativeWriteUtils {
/**
* Fail the task if the native writer would not write to exactly `filePath`.
*
- * [[escapedHdfsDestination]] declines the shapes Comet can predict at
planning time, but the
+ * [[unsupportedDestination]] declines the shapes Comet can predict at
planning time, but the
* path a write actually uses comes from
`FileCommitProtocol.newTaskTempFile`, and a custom
* commit protocol can return anything. This is the backstop: it runs before
the writer opens
* anything, so the task fails with a clear message rather than committing
successfully with the
* data left somewhere else.
*/
- def checkNativeWriteDestination(filePath: String): Unit =
+ def checkNativeWriteDestination(filePath: String): Unit = {
hdfsPathDivergence(filePath).foreach { shown =>
throw new UnsupportedOperationException(
s"Comet's native Parquet writer cannot write to '$filePath': the path
it would create " +
s"on HDFS is not the one the commit protocol chose ($shown). Set " +
"spark.comet.write.parquet.enabled=false to write this table with
Spark.")
}
+ if (s3Scheme(filePath).isDefined &&
+ (s3PathDivergence(filePath).isDefined || isS3AMagicPath(filePath))) {
+ throw new UnsupportedOperationException(
+ s"Comet's native Parquet writer cannot write to '$filePath': the
object it would create " +
+ "on S3 is not the one the commit protocol expects. Set " +
+ "spark.comet.parquet.write.enabled=false to write this table with
Spark.")
Review Comment:
In the message on line 191, the config is
`spark.comet.write.parquet.enabled` whereas here it is
`spark.comet.parquet.write.enabled`.
Can you correct to the right config.
--
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]