Skip to content

perf(table): run compaction groups concurrently in RewriteDataFiles #2040

Description

@KranzL

Proposed change

Transaction.RewriteDataFiles runs compaction groups one at a time. The atomic path loops over groups at table/rewrite_data_files.go:335-361 and calls ExecuteCompactionGroup at :344. The partial-progress path has the same shape per batch at :626-647 and calls it at :631. Each group is an independent read, decode, encode and write pipeline. Only intra-group read concurrency exists today, through WithCompactionScanConcurrency (table/rewrite_data_files.go:249-252); writes use the clustered single-writer path (:466-468).

I propose a MaxConcurrentGroups field on RewriteDataFilesOptions (table/rewrite_data_files.go:180-220, next to MaxCommits and PartialProgress) that bounds how many ExecuteCompactionGroup calls run at once in both paths. Zero and one keep today's sequential behavior. A negative value is rejected the way MaxCommits below zero is at :572-573.

The following hold on main at 833e10c:

  • Output names cannot collide. Each WriteRecords call mints a fresh random write UUID when args.writeUUID is nil (table/arrow_utils.go:2076-2078 and :2218-2220), so two groups never share an output name.
  • Scans are independent per call. Table.Scan has a value receiver and builds a fresh Scan per call (table/table.go:1338). The per-call scans share the snapshot manifest cache from perf(table): cache snapshot manifests across scans #1970 (:1344), which is built for concurrent scans.
  • The mechanism has an in-repo precedent. Orphan cleanup bounds parallel work with errgroup SetLimit (table/orphan_cleanup.go:562-563) with a runtime.GOMAXPROCS default (:284) set through WithCleanupMaxConcurrency (:132-138).

No open issue tracks this. gh issue list -R apache/iceberg-go --state open --search "compaction concurrent parallel groups" returns nothing. The broader "compaction parallel" search returns only #860 (an Overwrite race) and #1178 (the REST scan planning epic), neither about compaction groups.

Measured evidence

Setup: a partitioned v2 table with 8 partitions and 8 data files per partition, 30000 rows per file of (int64 id, string data, 96 hex char payload, float64 score), 222.8 MB on disk. Each group holds one partition (8 files, 240000 rows). Group count N runs RewriteDataFiles over the first N groups of the same table with default options (the atomic path) and a fresh transaction per N; the table is never committed between runs, so every N sees identical inputs. Timing covers the RewriteDataFiles call only, including the per-group ExecuteCompactionGroup calls and the in-transaction rewrite commit. User CPU is the RUSAGE_SELF utime delta around the call.

Apple M3 Pro, 11 cores, go1.25.9 darwin/arm64, go test -run TestZZGroupScaleSequential -count=5, wall in ms, five runs:

groups wall, five runs user CPU in s, five runs user CPU per wall
1 311, 232, 214, 218, 225 0.274, 0.272, 0.255, 0.266, 0.274 0.88, 1.17, 1.19, 1.22, 1.22
2 455, 436, 439, 428, 440 0.529, 0.524, 0.514, 0.504, 0.534 1.16, 1.20, 1.17, 1.18, 1.21
4 826, 856, 812, 839, 845 0.989, 1.029, 0.998, 1.005, 1.015 1.20, 1.20, 1.23, 1.20, 1.20
8 1667, 1633, 1695, 1693, 1677 1.961, 1.982, 2.026, 2.036, 2.031 1.18, 1.21, 1.20, 1.20, 1.21

Wall time is linear in group count: 8 groups cost 7.5x one group (medians 225 ms and 1677 ms) while user CPU stays near one core on an 11 core machine. Per-group Parquet encode is single threaded, so the other cores stay idle as groups are added.

Design

MaxConcurrentGroups is an int on RewriteDataFilesOptions, not a CompactionGroupOption. GroupOptions are forwarded to every ExecuteCompactionGroup call (table/rewrite_data_files.go:215-219) and tune one group's pipeline; the number of groups in flight is a property of the executor loop. MaxCommits (:186-190) is the precedent for an executor-level int on this struct.

Both executor loops run ExecuteCompactionGroup calls under an errgroup with SetLimit(MaxConcurrentGroups), following table/orphan_cleanup.go:562-563. Results are collected per index and applied in original group order before the existing staging and commit logic, so commit semantics do not change.

Memory sizing: peak record-pipeline memory is MaxConcurrentGroups times the per-group bound PR #2039 states, which is workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n rows, where workers is the scan worker count and n is the WithCompactionArrowBatchSize value. PR #2039 is open; its docs are at table/rewrite_data_files.go:249-283 on branch bounded-scan-reorder-heap. The PR that implements this issue repeats that formula times MaxConcurrentGroups in the MaxConcurrentGroups doc.

Correctness requirements

  • Group results are applied in original group order (rewrite.ApplyResult and the batch commit lists follow the input order), so manifests and RewriteResult counters are deterministic.
  • The first group error cancels the remaining groups and is returned. In the atomic path no group result is applied after a failure.
  • Partial-progress cleanup removes outputs from every group in the batch, including a group that failed mid-write. ExecuteCompactionGroup returns partial NewDataFiles with its error (:500-507), and the batch cleanup (:612-624) must see those files from every group, failed or not.
  • Commit semantics are unchanged: one snapshot for the atomic path, at most MaxCommits batch snapshots for partial progress.
  • The default is unchanged: zero and one mean sequential execution with today's behavior.

Validation plan

  • Unit tests for ordered application of concurrent group results, first-error cancellation, and partial-progress cleanup covering groups that failed mid-write.
  • go test -race over the rewrite suites and the rolling, fanout and clustered writer suites, the set the Add tunables bounding compaction pipeline memory #1993 review ran.
  • A multi-group benchmark at concurrency 1, 2, 4 and 8 on a table of the shape in the measured evidence section, showing wall scaling with group count at each level.

Related

Willingness to contribute

  • I can contribute this improvement/feature independently
  • I would be willing to contribute this improvement/feature with guidance from the Iceberg community
  • I cannot contribute this improvement/feature at this time

Specifications

  • Table
  • View
  • REST
  • Puffin
  • Encryption
  • Other

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

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions