Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -460,6 +460,7 @@ jobs:
value: |
org.apache.comet.parquet.CometParquetWriterSuite
org.apache.comet.parquet.CometEmptyRelationParquetWriterSuite
org.apache.spark.sql.comet.CometNativeWriteSuite
org.apache.comet.parquet.ParquetReadV1Suite
org.apache.comet.parquet.ParquetTimestampLtzAsNtzSuite
org.apache.spark.sql.comet.ParquetEncryptionITCase
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,7 @@ jobs:
value: |
org.apache.comet.parquet.CometParquetWriterSuite
org.apache.comet.parquet.CometEmptyRelationParquetWriterSuite
org.apache.spark.sql.comet.CometNativeWriteSuite
org.apache.comet.parquet.ParquetReadV1Suite
org.apache.comet.parquet.ParquetTimestampLtzAsNtzSuite
org.apache.spark.sql.comet.ParquetEncryptionITCase
Expand Down
7 changes: 7 additions & 0 deletions .github/workflows/spark_sql_writer_tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,13 @@ jobs:
spark-short-version: ${{ inputs.spark-version }}
skip-native-build: true

- name: Pre-compile Spark SQL test classes
run: |
cd apache-spark
rm -rf /root/.m2/repository/org/apache/parquet
NOLINT_ON_COMPILE=true build/sbt -Dsbt.log.noformat=true -mem 3072 \
'sql/Test/compile'
- name: Run Parquet writer tests
run: |
cd apache-spark
Expand Down
6 changes: 4 additions & 2 deletions docs/source/user-guide/latest/compatibility/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,8 +29,10 @@ executed.
Supported parent joins and aggregates remain eligible for native execution. Global aggregates
still return one row (`COUNT = 0`, `SUM = NULL`), and grouped aggregates return no rows. Independent
operator restrictions and aggregate buffer compatibility checks still apply.
Parquet writes whose input plans contain an empty relation use Spark's writer to preserve
readable empty output files and their schema metadata.

A native Parquet write over a native empty relation stays native. Like Spark's writer, it runs one
task for the empty input, so the output still gets a schema-only Parquet file that readers can
infer the schema from.

## In-Memory Cache

Expand Down
2 changes: 1 addition & 1 deletion docs/source/user-guide/latest/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ omitted from the tables below and may be reconsidered based on demand:
| `FileSourceScanExec` | ✅ | Parquet only. Some types and configurations fall back. See [Parquet Scan Compatibility](compatibility/scans.md). |
| `BatchScanExec` | ✅ | Apache Iceberg Parquet scans run natively. Native CSV scans are experimental and disabled by default. DataSource V2 Parquet scans are not accelerated. See [Parquet Scan Compatibility](compatibility/scans.md) and the [Iceberg Guide](iceberg.md). |
| `LocalTableScanExec` | ⚠️ | Disabled by default; there is no acceleration advantage and this operator is typically only used in test code. Can be opted into via config ([#4393](https://github.com/apache/datafusion-comet/pull/4393)). |
| `EmptyRelationExec` | ✅ | Spark 4.0 and later. See [Empty Relations](compatibility/operators.md#empty-relations) for native-input support and writer fallback. |
| `EmptyRelationExec` | ✅ | Spark 4.0 and later. See [Empty Relations](compatibility/operators.md#empty-relations) for native-input support, including native Parquet writes. |
| `RangeExec` | ⚠️ | Disabled by default. Set `spark.comet.exec.range.enabled=true` to generate the rows of `spark.range` and SQL `range()` in native code, so the operators above them run natively. It can be slower than Spark when those operators are only cheap expressions, such as a filter, which Spark compiles together with the range into one loop. Ranges whose arithmetic overflows the `Long` range fall back to Spark. |
| `InMemoryTableScanExec` | ⚠️ | Experimental, disabled by default. Set `spark.comet.exec.inMemoryCache.enabled=true` before the application starts so Comet installs its Arrow cache serializer. Relations with unsupported column types stay in Spark's cache format and fall back. See [In-Memory Cache](in-memory-cache.md). |

Expand Down
72 changes: 13 additions & 59 deletions native/core/src/execution/operators/parquet_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -234,18 +234,9 @@ impl ParquetWriter {
pub struct ParquetWriterExec {
/// Input execution plan
input: Arc<dyn ExecutionPlan>,
/// Where this task writes. When `work_dir` is set (the Spark 3.x `CometNativeWriteExec`
/// path) this is the write's output directory and is unused; the file name is derived from
/// `work_dir`. Otherwise (Spark 4.0+, `CometWriteFilesExec`) it is the exact path of the file
/// to write, chosen by the JVM commit protocol and used verbatim - this operator then never
/// derives file names of its own.
/// The exact path of the file this task writes, chosen by Spark's commit protocol on the JVM
/// side and used verbatim - this operator never derives file names of its own.
output_path: String,
/// Working directory for temporary files (used by FileCommitProtocol). Spark 3.x only.
work_dir: Option<String>,
/// Job ID for tracking this write operation
job_id: Option<String>,
/// Task attempt ID for this specific task
task_attempt_id: Option<i32>,
/// Compression codec
compression: ParquetCompression,
/// Partition ID (from Spark TaskContext)
Expand All @@ -268,9 +259,6 @@ impl ParquetWriterExec {
pub fn try_new(
input: Arc<dyn ExecutionPlan>,
output_path: String,
work_dir: Option<String>,
job_id: Option<String>,
task_attempt_id: Option<i32>,
compression: ParquetCompression,
partition_id: i32,
column_names: Vec<String>,
Expand All @@ -290,9 +278,6 @@ impl ParquetWriterExec {
Ok(ParquetWriterExec {
input,
output_path,
work_dir,
job_id,
task_attempt_id,
compression,
partition_id,
column_names,
Expand Down Expand Up @@ -478,9 +463,6 @@ impl ExecutionPlan for ParquetWriterExec {
1 => Ok(Arc::new(ParquetWriterExec::try_new(
Arc::clone(&children[0]),
self.output_path.clone(),
self.work_dir.clone(),
self.job_id.clone(),
self.task_attempt_id,
self.compression.clone(),
self.partition_id,
self.column_names.clone(),
Expand Down Expand Up @@ -510,8 +492,6 @@ impl ExecutionPlan for ParquetWriterExec {
.register(context.memory_pool());
let input = self.input.execute(partition, context)?;
let input_schema = self.input.schema();
let work_dir = self.work_dir.clone();
let task_attempt_id = self.task_attempt_id;
let compression = self.compression.to_parquet()?;
let column_names = self.column_names.clone();

Expand All @@ -530,19 +510,8 @@ impl ExecutionPlan for ParquetWriterExec {
Arc::new(Schema::new(fields))
});

let part_file = match &work_dir {
// Spark 4.0+ hands over the exact file to write, chosen by the JVM commit protocol.
None => self.output_path.clone(),
// Spark 3.x hands over a working directory instead and expects the writer to name the
// file; that branch goes away with Spark 3.x support.
Some(work_dir) => match task_attempt_id {
Some(attempt_id) => format!(
"{}/part-{:05}-{:05}.parquet",
work_dir, self.partition_id, attempt_id
),
None => format!("{}/part-{:05}.parquet", work_dir, self.partition_id),
},
};
// The JVM commit protocol has already chosen the exact file to write.
let part_file = self.output_path.clone();

// Configure writer properties
let props = WriterProperties::builder()
Expand Down Expand Up @@ -675,11 +644,10 @@ mod tests {
);
}

/// Spark 4.0+ hands over the exact file to write rather than a working directory. The writer
/// must use that path verbatim - Spark's commit protocol owns naming and staging, and
/// The JVM hands over the exact file to write. The writer must use that path verbatim - Spark's commit protocol owns naming and staging, and
/// committers that track individual files depend on the name it chose.
#[tokio::test]
async fn test_parquet_writer_uses_output_path_verbatim_without_work_dir() -> Result<()> {
async fn test_parquet_writer_uses_output_path_verbatim() -> Result<()> {
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, true)]));
let batch = RecordBatch::try_new(
Arc::clone(&schema),
Expand All @@ -696,9 +664,6 @@ mod tests {
let writer = ParquetWriterExec::try_new(
input,
output_path,
None, // work_dir: Spark 4.0+ path
None,
None,
ParquetCompression::None,
// A non-zero partition id must not leak into the file name.
3,
Expand Down Expand Up @@ -763,13 +728,10 @@ mod tests {
let memory_source = MemorySourceConfig::try_new(&[vec![batch]], input_schema, None)?;
let input = Arc::new(DataSourceExec::new(Arc::new(memory_source)));
let temp_dir = tempfile::tempdir()?;
let work_dir = format!("file://{}", temp_dir.path().display());
let output_path = format!("file://{}/part-00000.parquet", temp_dir.path().display());
let writer = ParquetWriterExec::try_new(
input,
work_dir.clone(),
Some(work_dir),
None,
None,
output_path,
ParquetCompression::None,
0,
vec!["required_id".to_string(), "values".to_string()],
Expand Down Expand Up @@ -832,13 +794,10 @@ mod tests {
let memory_source = MemorySourceConfig::try_new(&[vec![batch]], input_schema, None)?;
let input = Arc::new(DataSourceExec::new(Arc::new(memory_source)));
let temp_dir = tempfile::tempdir()?;
let work_dir = format!("file://{}", temp_dir.path().display());
let output_path = format!("file://{}/part-00000.parquet", temp_dir.path().display());
let writer = ParquetWriterExec::try_new(
input,
work_dir.clone(),
Some(work_dir),
None,
None,
output_path,
ParquetCompression::None,
0,
vec!["values".to_string()],
Expand Down Expand Up @@ -903,9 +862,6 @@ mod tests {
let writer = ParquetWriterExec::try_new(
input,
format!("file://{}", file.display()),
None,
None,
None,
ParquetCompression::None,
0,
column_names,
Expand Down Expand Up @@ -1160,16 +1116,14 @@ mod tests {
let memory_exec = Arc::new(DataSourceExec::new(Arc::new(memory_source_config)));

// Create ParquetWriterExec with DataSourceExec as input
let output_path = "unused".to_string();
let work_dir = "hdfs://namenode:9000/user/test_parquet_writer_exec".to_string();
let output_path =
"hdfs://namenode:9000/user/test_parquet_writer_exec/part-00000-00123.parquet"
.to_string();
let column_names = vec!["id".to_string(), "name".to_string()];

let parquet_writer = ParquetWriterExec::try_new(
memory_exec,
output_path,
Some(work_dir),
None, // job_id
Some(123), // task_attempt_id
ParquetCompression::None,
0, // partition_id
column_names,
Expand Down
3 changes: 0 additions & 3 deletions native/core/src/execution/planner/write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -107,9 +107,6 @@ impl OperatorBuilder for ParquetWriterBuilder {
let parquet_writer = Arc::new(ParquetWriterExec::try_new(
Arc::clone(&child.native_plan),
writer.output_path.clone(),
writer.work_dir.clone(),
writer.job_id.clone(),
writer.task_attempt_id,
codec,
planner.partition(),
writer.column_names.clone(),
Expand Down
30 changes: 11 additions & 19 deletions native/proto/src/proto/operator.proto
Original file line number Diff line number Diff line change
Expand Up @@ -785,7 +785,7 @@ message IcebergWriteCommon {
}

// Single Iceberg write operator. Per-task fields are populated by the JVM exec wrapper at task
// launch (matching how `ParquetWriter.task_attempt_id` is filled in `CometNativeWriteExec`).
// launch (matching how `ParquetWriter.output_path` is filled in by the Parquet write execs).
message IcebergWrite {
IcebergWriteCommon common = 1;

Expand Down Expand Up @@ -920,27 +920,19 @@ message ShuffleWriter {
}

message ParquetWriter {
// Where this task writes. Two shapes, selected by whether `work_dir` is set:
//
// - Spark 4.0+ (CometWriteFilesExec): `work_dir` is unset and `output_path` is the
// fully-qualified path of the Parquet file to write, set per task from
// FileCommitProtocol.newTaskTempFile. Naming and staging are owned by Spark's commit
// protocol so that task-attempt isolation, speculative execution, and committers that track
// individual files (S3A magic, streaming manifest) all behave as they do for Spark's own
// writer. The native writer uses this path verbatim.
// - Spark 3.x (CometNativeWriteExec): `work_dir` is set and the native writer derives the file
// name from it, the partition id and the task attempt id. `output_path` is the write's
// output directory and is unused natively. Goes away with Spark 3.x support.
// The fully-qualified path of the Parquet file this task writes, set per task from
// FileCommitProtocol.newTaskTempFile by CometWriteFilesExec (Spark 4.0+) or
// CometNativeWriteExec (Spark 3.x). Naming and staging are owned by Spark's commit protocol so
// that task-attempt isolation, speculative execution, and committers that track individual
// files (S3A magic, streaming manifest) all behave as they do for Spark's own writer. The native
// writer uses this path verbatim.
string output_path = 1;
CompressionCodec compression = 2;
repeated string column_names = 4;
// Working directory for temporary files (used by FileCommitProtocol). Spark 3.x only; see
// output_path above.
optional string work_dir = 5;
// Job ID for tracking this write operation
optional string job_id = 6;
// Task attempt ID for this specific task
optional int32 task_attempt_id = 7;
// Formerly work_dir, job_id and task_attempt_id, which let the native writer name its own file
// under a working directory.
reserved 5, 6, 7;
reserved "work_dir", "job_id", "task_attempt_id";
// Options for configuring object stores such as AWS S3, GCS, etc. The key-value pairs are taken
// from Hadoop configuration for compatibility with Hadoop FileSystem implementations of object
// stores.
Expand Down
Loading
Loading