Skip to content

feat: adding elapsed_compute to the writer - #5143

Open
comphead wants to merge 4 commits into
apache:mainfrom
comphead:writer_ela
Open

comphead wants to merge 4 commits into
apache:mainfrom
comphead:writer_ela

Conversation

@comphead

@comphead comphead commented Jul 30, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #.

Rationale for this change

Adding a metric to track time spent on native writes. I was trying to measure time spent on native writes and surprisingly Comet doesn't report such metric.

What changes are included in this PR?

How are these changes tested?

@comphead
comphead requested a review from andygrove July 30, 2026 16:50
Comment on lines -149 to -151
// Clear the buffer after upload
cursor.get_mut().clear();
cursor.set_position(0);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it intentional that this is removed? The PR description doesn't talk about this change

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. The tuple-to-named-struct conversion on ParquetWriter::Remote is a clear readability win. Four positional fields repeated in every match arm was hard to follow. The hdfs_writer.take() restructure in close() also drops a redundant is_none() pre-check without changing semantics.

I do think there is a blocking problem in the remote path, plus a few other things worth sorting out before merge.

into_inner() after finish() always fails

native/core/src/execution/operators/parquet_writer.rs:172-174

arrow_writer.finish()?;
let total_bytes = arrow_writer.bytes_written() as u64;
let cursor = arrow_writer.into_inner()?;

ArrowWriter::finish() calls SerializedFileWriter::finish(), which calls write_metadata(), which sets self.finished = true. ArrowWriter::into_inner() then calls SerializedFileWriter::into_inner(), and that starts with assert_previous_writer_closed(), which returns Err when finished is already true.

I checked this against parquet 58.4.0 (the version we pin) with a standalone program:

bytes_written after finish = 446
into_inner ERR: Parquet error: SerializedFileWriter already finished

So close() returns Err for every remote write. The footer never gets uploaded and the write fails with Failed to close writer: Parquet error: SerializedFileWriter already finished. The LocalFile arm is fine because it never calls into_inner().

Pulling the buffer out through inner_mut() works, since finish() has already flushed the BufWriter down into the cursor:

arrow_writer.finish()?;
let total_bytes = arrow_writer.bytes_written() as u64;
let final_data = std::mem::take(arrow_writer.inner_mut().get_mut());

I confirmed end to end that this gives a byte-exact, readable file across three incremental uploads plus the footer. total_bytes matched the uploaded length exactly and the readback row count was correct.

Nothing in CI covers the path that broke

CI is fully green here, but that does not tell us much for this change. All three HDFS tests in this file are #[ignore]d with "This test requires a running HDFS cluster", and test_write_to_hdfs_sync and test_write_to_hdfs_streaming build their own ArrowWriter rather than going through ParquetWriter::Remote at all. The variant this PR rewrites has no automated coverage, which is why the above slipped through.

Could we add a test that uses opendal::services::Memory for the Operator? That needs no cluster, so it can run as a normal unit test. Write a few batches, close, read the uploaded bytes back with ParquetRecordBatchReaderBuilder, and assert the returned count equals the uploaded length. That may need services-memory added to the opendal dev-dependency features.

elapsed_compute ends up measuring the whole subtree

native/core/src/execution/operators/parquet_writer.rs:523

let _timer = elapsed_compute.timer();

The guard is held for the entire write_task future, which includes stream.try_next().await on the input. A ScopedTimerGuard held across an await keeps accruing wall clock while the future is suspended, so this reports time from first poll to completion of the whole subtree rather than time spent writing.

The other Comet operators do the opposite and stop the timer around anything they do not own. See shuffle_scan.rs:306 and scan.rs:122, which both bind a mut timer and call timer.stop() explicitly. Would it be worth scoping the timer to just the writer.write() and writer.close() calls so the number is comparable to the rest of the plan in the Spark UI? The Scala label is "total time (in ms) spent in this operator", which points at the narrower reading too.

Worth noting the label and createNanoTimingMetric do match CometMetricNode.baselineMetrics, so the naming and plumbing are consistent.

The bytes_written change is an unmentioned bug fix

native/core/src/execution/operators/parquet_writer.rs:160-161

The old code did std::fs::metadata(local_path).map(...).unwrap_or(0), so bytes_written was silently reported as 0 for any non-local path. Moving to ArrowWriter::bytes_written() fixes that, and I verified the two agree exactly for local files (28593 both ways).

This is a real improvement, so could you call it out in the description? It changes a metric users see. CometParquetWriterSuite has 30 tests and asserts nothing about metrics, so an assertion comparing bytes_written against the summed on-disk sizes of the written part files would lock this in.

flush() does not push bytes toward the destination

native/core/src/execution/operators/parquet_writer.rs:124

The new comment says the flush hands the buffered bytes over, but ArrowWriter::flush() is documented as "the underlying writer is not flushed with this call", and TrackedWrite wraps the cursor in a BufWriter. In practice about 8KB stays behind on every batch:

remote: after flush, bytes_written=28179 cursor_len=20043

Ordering is preserved so nothing corrupts, and this predates the PR. Since you are rewriting these exact lines and adding a comment about them, adding arrow_writer.sync()? after the flush() would make the upload genuinely incremental and make the comment accurate.

Description and scope

The title and body cover the metric, but the diff also rewrites the remote buffer handling and changes how bytes_written is computed. There is no linked issue either. Could you fill in the description, and consider whether the buffer and close rework wants to be its own PR given the problem above?

@comphead
comphead requested a review from andygrove July 31, 2026 15:26
@comphead

comphead commented Aug 2, 2026

Copy link
Copy Markdown
Contributor Author

Thanks @andygrove for the review, addressed comments in 8937de0

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the updates. All five points from my last round are addressed, and I verified each one against parquet 58.4.0 rather than taking it on faith.

The into_inner() fix is correct. SerializedFileWriter::finish() at file/writer.rs:299 ends with self.buf.flush()?, so the footer really is sitting in the cursor by the time inner_mut() is reached. And into_inner() at file/writer.rs:441 does start with assert_previous_writer_closed(), which confirms the old code could never have worked. The explanatory comment you added is accurate.

The timer scoping is right. stop() is now called after writer.write() and after writer.close(), and stream.try_next().await sits outside both windows. Error paths are covered too, since Drop for ScopedTimerGuard calls stop().

The LocalFile switch from close() to finish() plus bytes_written() is safe. ArrowWriter::close() is literally just self.finish(), and ArrowWriter has no Drop impl, so dropping the writer at the end of the match arm is equivalent to the old behavior.

The new memory-backend test is the best part of this round. I ran it and it passes. It also actually executes in CI, which I was initially unsure about, because native/core/Cargo.toml:98 sets default = ["hdfs-opendal"]. The dev-dependency pins opendal 0.57.0, matching the optional normal dependency, so there is no duplicate-crate problem. The test has real teeth as well. Without the into_inner() fix close() returns Err and the test fails outright, and without cursor.set_position(0) the zero-fill would break the byte-length assertion.

I have left inline comments on three things. Two more that do not map to a diff line:

The metric assertion did not make it in

This is the one item from last round I did not see addressed. The PR changes bytes_written from std::fs::metadata(...).unwrap_or(0) to ArrowWriter::bytes_written(), which is a real fix, since the old code reported 0 for any non-local path. It is also a user-visible metric change with no test pinning it down.

CometParquetWriterSuite still has 30 tests and asserts nothing about metrics. It already captures the plan at line 103, so could we add assertions there that bytes_written matches the summed on-disk size of the part files, and that elapsed_compute is greater than zero? Without that, nothing catches a regression in either the new metric or the changed one.

Description and linked issue

Could you fill in the description? It is still the unedited template, including Closes #. and the instructional comments.

Worth calling out explicitly that this also fixes bytes_written reporting 0 for non-local paths, since that is a user-visible metric change that the title and body do not currently mention.

// `BufWriter`. `sync()` pushes those bytes down into our cursor so the upload
// is genuinely incremental. Then take ownership of the cursor's buffer and
// reset it to empty for the next batch (no clone, no explicit clear).
arrow_writer.flush()?;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The sync() addition is correct. I confirmed ArrowWriter::sync() is self.writer.flush(), which reaches SerializedFileWriter::flush() and then TrackedWrite::flush(), so the upload is genuinely incremental now and your comment is accurate.

Separately, I noticed flush() here creates a new row group on every batch. I confirmed this with your new test, where three batches produce exactly three row groups. The local path does not do this, since it just calls writer.write(batch) and lets ArrowWriter manage row group boundaries at its default of 1,048,576 rows. So a file written to HDFS ends up with a row group every 8192 rows, roughly 128 times more row groups than the same data written locally, which hurts compression and bloats the footer.

This predates your PR so I am not asking you to fix it here. Given that you are measuring write performance, could you file a tracking issue for it and link it from this thread?

let final_data = std::mem::take(arrow_writer.inner_mut().get_mut());

// Write any remaining data
if !final_data.is_empty() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This guard can never be false. I checked the empty-partition case, where no batches are written at all, and finish() still emits a 355 byte footer, so close() returns Ok(355) and a valid empty Parquet file lands in the store. That case works correctly.

The reason I would rather see the guard gone than left alone is the branch that never runs. If final_data were ever empty we would drop hdfs_writer without calling close() on it, and for a multipart upload that means the object is silently never committed. Silent data loss behind a currently unreachable condition is worth deleting while these lines are already being rewritten. Could we just always take or create the writer and close it unconditionally?

/// bytes back with `ParquetRecordBatchReaderBuilder`, and asserts the returned
/// row count and the reported `bytes_written` both match the upload.
#[tokio::test]
#[cfg(feature = "hdfs-opendal")]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I ran it and it passes, and I confirmed it has teeth, since it fails without the into_inner() fix and again without the set_position(0) reset. It also runs in CI by default, because hdfs-opendal is in the default feature set.

Since this is the only automated coverage the remote writer has, would you mind also comparing the read-back values against what was written rather than just the row count? create_test_record_batch already generates distinct data per batch, so it is only a few extra lines, and buffer and cursor-position juggling is exactly the kind of change that can corrupt values.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Native Parquet writes lacked operator timing, and remote bytes_written relied on local filesystem metadata, producing zero.
  • Design approach: Add DataFusion elapsed_compute, expose it through Spark’s createNanoTimingMetric, and obtain file size from the finalized Parquet writer.
  • Correctness / compatibility analysis: Checked Spark’s timing and write-statistics sources for 3.4.3, 3.5.8, 4.0.2, 4.1.2 and 4.2.0. Nanosecond units and metric transfer are consistent. The earlier remote-close failure and inclusion of upstream processing time are fixed.
  • Key design decisions: finish() flushes the footer before byte accounting. Remote buffers are taken through inner_mut(), avoiding the invalid finish() followed by into_inner() sequence. Cursor resets preserve incremental writes.
  • Implementation sketch: Replace positional remote-writer fields with named fields, synchronize and transfer buffers, time each batch and finalization, and add a memory-backend test. This uses existing abstractions without adding another framework.
  • Behavioral changes worth calling out: Timing excludes input polling but includes destination I/O during writes and finalization. Remote byte counts now reflect the complete file. Instrumentation operates per batch, and ownership transfer removes the explicit buffer clone.
  • Suggested improvements: No additional P1/P2 changes requested. No introduced P1/P2 issues found within this review, and no substantiated unresolved P1/P2 blocker remains.

Reviewed the full three-file, four-commit diff from 38c3f8bce35642479117ba49f44a02fee31e6158 to 8937de03cdbed8f00b77e647041e073aa5b70533. Confirmed non-draft status and read all supplied reviews, comments and threads.

Routed skills: review-comet-pr and review-comet-memory-pr. This historical checkout contains the general skill under .claude/skills. Also consulted the requested .ai general and memory skills and verified they match upstream.

Exact-head CI: 61 successful checks, nine skipped, no failures or unfinished checks. Inspected logs confirm the new remote-memory test and CometParquetWriterSuite passed. These head-associated checks ran GitHub’s merge commit 0f95b06, containing the reviewed head.

Validation: A disposable harness using the unchanged exact-head ParquetWriter implementation passed 70 local/memory-backend cases across five codecs, checking complete value readback and byte counts for empty files, empty batches, nulls, multiple batches and row-group boundaries. The repository’s focused offline test could not start because async-compression v0.4.42 was uncached. No local JVM suite, live HDFS cluster test or performance benchmark was run.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Main has moved a lot since the last round, so I merged it into this branch locally to see where things stand. The conflicts go beyond the text. Main bumped opendal to 0.58 (#5324), where Operator::new returns the Operator directly, so the new test's .finish() no longer compiles, and the dev-dependency should move to 0.58.1 to match the normal one. The Remote fields are boxed now (Box<ArrowWriter<...>> and Box<Operator>). #6247 also added ParquetWriter::memory_size(), whose doc comment says the remote staging buffer keeps its capacity between uploads. That stops being true with the std::mem::take. The remote writer's reservation drops to 0 after each batch, which matches what it now holds, so only the comment needs to change. With those resolved, cargo fmt, clippy with -D warnings, the native writer tests, and CometParquetWriterSuite and CometEmptyRelationParquetWriterSuite on Spark 4.1 all pass for me. #4746 touches parquet_writer.rs and CometNativeWriteExec too and is approved, so there may be one more small merge after it lands.

Three things from the last round are still open. Two are the inline threads on the if !final_data.is_empty() guard and on comparing values in the memory-backend test. The third is a JVM assertion on the metrics. CometParquetWriterTestBase.isNativeWriteExec already picks the right node on either Spark line, so a test could assert elapsed_compute > 0 on that node, and on 3.x compare bytes_written against the summed part-file sizes. I checked both by hand on the merged branch. On Spark 3.5, bytes_written was 120,230, exactly the on-disk total.

The remote path still starts a new row group for every batch. I filed #6774 for that, so it doesn't need to be part of this PR.

"bytes_written" -> SQLMetrics.createSizeMetric(sparkContext, "written data"),
"rows_written" -> SQLMetrics.createMetric(sparkContext, "number of written rows"))
"rows_written" -> SQLMetrics.createMetric(sparkContext, "number of written rows"),
"elapsed_compute" -> SQLMetrics.createNanoTimingMetric(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since #5763, Spark 4.0+ writes go through CometWriteFilesExec instead, and this node only runs on 3.4 and 3.5. CometWriteFilesExec.metrics is Map.empty, and CometMetricNode.set drops any native metric that has no JVM entry, so on the 4.x builds the new timer is computed and thrown away. On Spark 4.1 with main merged in, the CometWriteFiles node reports no metrics at all. Adding the same elapsed_compute entry to CometWriteFilesExec.metrics makes it report about 27 ms for a 1,000-row write, in line with what this node reports on 3.5. Could we add it there too? The scaladoc on that metrics val would need an update as well. It says the node has no metrics of its own, and that bytes_written comes from std::fs::metadata, which this PR changes.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:writer Native Parquet writer enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants