Skip to content

Native Parquet writer starts a new row group for every batch when writing to HDFS #6774

Description

@andygrove

Describe the bug

ParquetWriter::Remote, the HDFS path in native/core/src/execution/operators/parquet_writer.rs, calls ArrowWriter::flush() after every input batch so that it can upload the encoded bytes incrementally. flush() closes the in-progress row group, so every batch becomes its own row group. With the default batch size of 8192 rows, a file written natively to HDFS gets a row group every 8192 rows.

The local path only calls write() and lets ArrowWriter cut row groups at its default of 1,048,576 rows. Spark's own writer cuts them at parquet.block.size (128 MB). Small row groups mean small column chunks, worse compression and dictionary encoding, a bigger footer, and more work for everything that reads the file.

Steps to reproduce

Write the same five 1,000-row batches through ParquetWriter::Remote with an in-memory opendal Operator, and through ParquetWriter::LocalFile. On main at f80042f the remote file has 5 row groups and the local file has 1.

Expected behavior

The remote path cuts row groups where the local path does. It could drop the per-batch flush(), call sync() so that completed row groups reach the cursor, and upload whatever the cursor holds. Since #6247, memory_size() charges the in-progress row group to the task's memory pool, so holding it until the row group is full is accounted for.

Additional context

Found during the #5143 review on 2026-08-04. It predates that PR. Related to #5304, which is about honoring parquet.block.size. Fixing #5304 alone would not help here, because the remote path closes a row group on every batch whatever the configured size is.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions