Skip to content

[Ingest] Optimize the number of parquet files created in the Shift task #120

Description

@s-prosvirnin

Problem

  • Parquet files in a single Iceberg snapshot overlap in time (3 shift tasks cover one time window with 2-3 second overlaps). File-level pruning doesn't work — queries are forced to read all files.
    • ParquetQueueReader::plan_segments groups WAL row groups only by tenant_id and splits into shift tasks by volume in WAL offset order, without considering sort key and row group time boundaries.
  • Each shift task writes two files: "large + tail" (65 MB + 0.35 MB). Source: RollingFileWriterBuilder with max_file_size_mb=64 cuts within a single task when output slightly exceeds max_file_size_mb.
  • IcebergStorage has an independent second level of volume control via max_file_size_mb. With planner target 64 MB and failover writer 64 MB, rollover triggers on every "random +1 MB" above the target.

Target Behavior

  • Planner starts accounting for Iceberg table sorting: considers RowGroupBoundaryRange from WAL row group metadata, sorts and clusters row groups within tenant.
  • Writer rollover becomes failover-only (writer_max_parquet_bytes, 2x upper_bound_input_bytes_per_task). Normally 1 shift task = 1 parquet file.
  • If small tail (last target parquet file in shift job) fits in previous parquet file — add to previous.
    • If last_chunk.bytes + prev_chunk.bytes ≤ (upper_bound_input_bytes_per_task) — merge into prev. Otherwise add as new last file.
    • For this, last_chunk.bytes must be <= lower_bound_input_bytes_per_task, otherwise it qualifies for separate last file.
    • The goal is to avoid making the last file < lower_bound_input_bytes_per_task.
  • Core rules for accumulating row groups to target parquet file size:
    • Iceberg parquet file size should be in range [lower_bound_input_bytes_per_task, upper_bound_input_bytes_per_task].
    • Parquet file size >= lower_bound_input_bytes_per_task (default 64 MB) if packed from non-overlapping clusters; <= upper_bound_input_bytes_per_task (default 128 MB) if single cluster.
    • For non-overlapping clusters, parquet file size stays at lower target threshold: once chunk >= lower_bound_input_bytes_per_task, can close if next cluster isn't needed to eliminate small tail or for other packing rule.
    • upper_bound_input_bytes_per_task is upper limit for stuffing additional data if it reduces overlaps/tails, and hard cap for oversized/split/tail-merge.
    • If total size of all non-overlapping clusters for tenant < lower_bound_input_bytes_per_task (e.g., 30 MB with lower_bound_input_bytes_per_task=64 MB), bin-pack yields one chunk of 30 MB. No alternative (nothing to pack from).
    • If atomic cluster (overlapping row groups) > upper_bound_input_bytes_per_task (e.g., 200 MB with upper_bound_input_bytes_per_task=128 MB), doesn't fit in one chunk. split_oversized cuts into sub-chunks ≤ upper_bound_input_bytes_per_task. Get multiple chunks. Each formally "from one source cluster", but disjointness between sub-chunks is lost.
    • Planner must account for day(timestamp) boundary during clustering (temporal partitioning). One logical chunk at day boundary → 2 physical files. Otherwise ingest around midnight produces wider, less predictable output.
  • Implementation must use resources optimally (CPU, memory, S3 requests). Means minimum operations, minimum memory consumption. But must maintain code readability balance.

Scheme Benefits

  • Output size is preserved. Volume-based split remains, just in sorted order. Files stay in target 64-128 MB range.
  • In practice outputs are disjoint or near-disjoint. Overlapping row groups of one service form atomic cluster, not split by bin-pack. Disjointness is lost only when cluster itself > upper_bound_input_bytes_per_task (see Case 2 / split_oversized)

Context

  • row_groups_merger writes in sort key order (verify)
  • Backward compatibility of configs and schemas not required (app not in production).

Refactoring

  • ParquetQueueReader should not take policy decisions like group_by_column_name = "tenant_id". Reader shouldn't know about tenant_id, RowGroupBoundaryRange, or planner.
    • Need to add to reader a pass-through of fields list to extract from parquet metadata. Grouping field should also be included and group_key_from_row_group removed. Return result in conditional structure HashMap<FieldId, ExtractedValue>.
    • Need to add to Reader API logic:
      • from column statistics extract field X as utf8 singleton
      • from file key-value metadata extract key Y as string payload
    • Then planner-side adapter interprets field X as tenant_id, payload Y as RowGroupBoundaryRange.
    • Target schema:
      • Reader should know only about extraction mechanisms: column stats, key-value metadata, row_group_idx binding. Reader shouldn't know specific fields, since Queue doesn't know schema.
      • Planner or intermediate adapter should know the meaning: this is tenant, this is boundary payload, parse into domain type. Planner shouldn't know storage quirks in Queue.

New Concepts

Concept Definition Size Formed In
Cluster Set of row groups of one tenant that have overlaps in composite sort key From 1 row group to infinity (depends on overlaps) swept_line_cluster in planner
Chunk Set of clusters (or sub-split of one oversized cluster) that become one shift task and then one parquet file >= lower_bound_input_bytes_per_task if packed from non-overlapping clusters; <= upper_bound_input_bytes_per_task if single cluster bin_pack in planner
lower_bound_input_bytes_per_task Lower bound of target parquet file. via config in planner
upper_bound_input_bytes_per_task Upper bound of target parquet file. via config in planner
writer_max_parquet_bytes Limit for max Iceberg parquet file. Needed because file size can exceed lower_bound_input_bytes_per_task to make fewer overlapping files. writer_max_parquet_bytes works as final limiter for max parquet file size. constant in IcebergStorage
  • Chunk size = sum of row_group_bytes for all row groups in it.
  • 15 row groups of one service that overlap each other — this is 120 MB cluster (with default 8MB row group in WAL). Happens when:
    • Several pods (instances) of one service write in parallel, or
    • Clients retry failed ingests, or
    • One service runs multithreaded and time in logs is significantly shuffled.
  • Current parameter max_input_bytes_per_task becomes lower_bound_input_bytes_per_task.
  • Add upper_bound_input_bytes_per_task (hard_cap) as upper bound of target parquet file.
  • Add writer_max_parquet_bytes (WRITER_FILE_SIZE_FAILOVER_FACTOR) as failover in case of unplanned exceeding upper_bound_input_bytes_per_task.

Nuances

  • Overlaps don't disappear completely, they shrink. Overlapping row groups of one service form atomic cluster, not split by bin-pack. Disjointness is lost only when cluster itself > upper_bound_input_bytes_per_task (see Case 2 / split_oversized)
  • Unbalanced shift tasks. If one tenant in commit has one "fat" service (say 300MB) and five "thin" ones (10MB each) — after sorting, fat service becomes 5 shift tasks on its own, thin ones merge into one. Task-time balance suffers. Same problem exists now — this is neutral, not worse.
  • Choice of row group sort key. To sort for disjointness in sequential cut, better by min_bound — then row groups with smaller keys land in earlier shifts. This gives monotonic sort-key space coverage left to right. Use min boundary key in lex sense, not min timestamp.
  • Bin-packing clusters: clusters (disjoint by construction) pack into shift tasks greedily by size, in sort key order. Concatenation of disjoint clusters stays disjoint relative to adjacent shift tasks.
  • Different services never overlap in composite boundary space, even if their timestamps overlap. Because key = (acct, service_name, ts DESC): if RG-A service_name="service_name-01", RG-B service_name="service_name-02", then max(RG-A) < min(RG-B) always, regardless of ts. Clusters naturally form per service. 20 services form 20 disjoint clusters (one per service), bin-pack packs into one chunk → one file (if fits in upper_bound_input_bytes_per_task).
  • Small files arise in exactly one case: one service_name has heavily overlapping row groups (retries, multi-replica), plus everything else in this tenant is empty. Then this service_name cluster is atomic, and if it's smaller than lower_bound_input_bytes_per_task, shift task comes out small. But that's fine, just low log throughput, only subsequent compaction helps.

Cases

Case 1. "Fat" service, monotonic in time (one producer)

Conditions: one service generates much more than upper_bound_input_bytes_per_task (e.g., 200 MB with upper_bound_input_bytes_per_task = 128 MB). Row groups arrive sequentially, time is monotonic (one pod, no retries).

Input: ~30 WAL row groups (~6.7 MB each), tenant=default, service_name=A, time ranges without overlaps: [t0..t10], [t10..t20], ..., [t290..t300]. Total 200 MB.

Output: sort by min_key (order already monotonic), swept_line_cluster gives ~30 disjoint clusters, bin_pack packs disjoint clusters to lower target threshold:

  • chunk 1: 10 row groups × 6.7 MB = ~67 MB. Chunk already >= lower_bound_input_bytes_per_task, next disjoint cluster not needed to eliminate small tail → flush.
  • chunk 2: next 10 row groups × 6.7 MB = ~67 MB. Chunk already >= lower_bound_input_bytes_per_task, next disjoint cluster not needed to eliminate small tail → flush.
  • chunk 3: remaining 10 row groups × 6.7 MB = ~67 MB.

tail_merge: chunk 3 (~67 MB) >= lower_bound_input_bytes_per_task (64 MB) → don't merge, file already in target range.

Output: 3 parquet files: ~67 MB, ~67 MB, ~67 MB, disjoint by timestamp. Writer doesn't do rollover (writer_max_parquet_bytes = 256 MB).

Downsides: 3 target files remain, because for disjoint clusters lower threshold is normal chunk close point. But strictly better than old behavior: files are disjoint by sort key, and writer no longer creates tail files from rollover at 64 MB.

Case 2. "Fat" service with overlapping (multi-pod)

Conditions: one service, many row groups, all with overlapping time-ranges (several replicas write in parallel or client mass-retries).

Input: 30 WAL row groups (~6.7 MB each), all service_name=A, time-ranges [10..50], [15..55], [20..60], .... Total 200 MB.

Output: swept_line_cluster merges all into one atomic cluster 200 MB. 200 MB > upper_bound_input_bytes_per_task (128 MB) → split_oversized cuts greedily:

  • sub-chunk 1: 19 row groups × 6.7 MB = ~127 MB (next row group exceeds upper_bound_input_bytes_per_task → flush).
  • sub-chunk 2: remaining 11 row groups × 6.7 MB = ~74 MB.
  • After sorting by min_bound, 30 RG spread by min's [10, 15, 20, ..., 155]. sub1 (first 19) ≈ bounds [10, ~140], sub2 (rest 11) ≈ [105, 200]. Overlap zone [105, 140] — narrow relative to file width. Pruning works: query ts > 150 cuts sub1, query ts < 100 cuts sub2.

tail_merge: sub-chunk 2 (74 MB) >= lower_bound_input_bytes_per_task (64 MB) → don't merge.

Output: 2 parquet files ~127 MB and ~74 MB, both in range [lower_bound_input_bytes_per_task, upper_bound_input_bytes_per_task], but with timestamp overlap between them (disjointness lost between sub-chunks of one source cluster).

Downsides:

  • File-level pruning by timestamp doesn't work for this window (same as old behavior).
  • Metric shift_planner_oversized_clusters_total increments → signal for investigation root cause (why retry storm).

Case 3. Bunch of thin services

Conditions: data from many services, each service's data much smaller than lower_bound_input_bytes_per_task.

Input: 20 row groups, service_name-01..service_name-20, 3 MB each. Total 60 MB.

Output: sort by min_key orders row groups by service_name. swept_line_cluster: 20 disjoint clusters (different services don't overlap in composite boundary space). bin_pack: 60 MB ≤ upper_bound_input_bytes_per_task → all in one chunk. tail_merge not applied (no previous chunk). 1 parquet file 60 MB with bounds [service_name-01..service_name-20].

Downsides: file (60 MB) slightly smaller than lower_bound_input_bytes_per_task (64 MB), but not a regression — just no more data in this commit. Identical to old behavior in file count, plus guarantee of disjointness within file by composite sort key.

Case 4. "Fat" + "thin" services together

Conditions: one large service and several small ones in one shift job.

Input: service_name-fat 200 MB (~30 row groups × ~6.7 MB), service_name-skinny-1..5 × 2 MB = 10 MB. Total 210 MB.

Output: sort by min_key: first all row groups of service_name-fat (by ts DESC), then service_name-skinny-1..5. swept_line_cluster gives 30 + 5 = 35 disjoint clusters (all monotonic). bin_pack packs disjoint clusters to lower target threshold:

  • chunk 1: 10 row groups of service_name-fat × 6.7 MB = ~67 MB. Chunk already >= lower_bound_input_bytes_per_task, next disjoint cluster not needed to eliminate small tail → flush.
  • chunk 2: next 10 row groups of service_name-fat × 6.7 MB = ~67 MB. Chunk already >= lower_bound_input_bytes_per_task, next disjoint cluster not needed to eliminate small tail → flush.
  • chunk 3: remaining 10 row groups of service_name-fat (~67 MB).
  • chunk 4: 5 clusters of service_name-skinny-* = 10 MB.

tail_merge: chunk 4 (10 MB) < lower_bound_input_bytes_per_task (64 MB), chunk 3 + chunk 4 = ~77 MB ≤ upper_bound_input_bytes_per_task (128 MB) → merge chunk 4 into chunk 3.

Output: 3 parquet files: ~67 MB (service_name-fat), ~67 MB (service_name-fat), and ~77 MB (service_name-fat remainder + all service_name-skinny-*).

Downsides:

  • Bounds of third file [service_name-fat..service_name-skinny-5] wider than first two. Service-level pruning for service_name-skinny-* doesn't cut third file.
  • Not worse than old behavior — previously more files and writer could create additional tail files. Here small skinny tail absorbed by previous parquet file via tail_merge.

Case 5. Overlapping row groups within service, fit in one task

Conditions: one service with row group time overlaps, but total size ≤ upper_bound_input_bytes_per_task.

Input: 4 row groups for service_name-A, time-ranges [10..30], [20..40], [50..70], [60..80]. Total 50 MB.

Output: swept_line_cluster gives 2 clusters (rg1+rg2 overlapping; rg3+rg4 overlapping; cluster A and cluster B disjoint). bin_pack: 50 MB ≤ upper_bound_input_bytes_per_task → 1 chunk. tail_merge not applied (no previous chunk). 1 parquet file 50 MB with timestamp bounds [10..80].

Downsides: file (50 MB) smaller than lower_bound_input_bytes_per_task (64 MB), but this is only chunk — tail_merge not applicable. Better than "dumb sort" — clustering protects from accidental cluster cut at upper_bound_input_bytes_per_task boundary.

Case 6. Several disjoint time-clusters of one service

Conditions: one service, two time "gaps" (e.g., idle and active periods).

Input: cluster A (service_name-A, time-range [10..40], 50 MB), cluster B (service_name-A, time-range [200..230], 50 MB). Total 100 MB.

Output: swept_line_cluster: 2 disjoint clusters (gap [40..200] between them). bin_pack: 50 + 50 = 100 MB ≤ upper_bound_input_bytes_per_task (128 MB) → both merge into one chunk. 1 parquet file 100 MB with timestamp bounds [10..230] (but no data in [40..200]).

Downsides:

  • File bounds don't reflect time gap [40..200] — file-level pruning might incorrectly decide file contains data in this period. Row-group-level pruning within file (by Parquet ColumnIndex) correctly cuts empty range.
  • File (100 MB) in target range [lower_bound_input_bytes_per_task, upper_bound_input_bytes_per_task].

Case 7. Multi-tenant in one commit

Conditions: 2 different tenants in one WAL window.

Input: tenant=A 60 MB, tenant=B 30 MB.

Output: group_by_tenant → 2 buckets. Each independently goes through sort → swept_line_cluster → bin_pack → tail_merge. Since each tenant's sum ≤ upper_bound_input_bytes_per_task — each gives 1 chunk. Output: 2 shift tasks, 2 parquet files: 60 MB (tenant=A) and 30 MB (tenant=B).

Downsides: both files smaller than lower_bound_input_bytes_per_task (64 MB), but tail_merge between tenants impossible — different partition values in Iceberg. Not a regression, behavior identical to current tenant split.

Case 8. Empty / null tenant_id

Conditions: WAL row group has tenant_id = null or empty string.

Input: row groups with tenant_id="".

Output: group_by_tenant returns hard error. Plan task fails entirely, retried by runtime after input fix.

Downsides:

  • Plan iteration on topic blocked until input fixed. This is intentional fail-fast cost: prefer noticeable breakage over silent corruption.

Case 9. Trickle ingest (one small row group)

Conditions: very low throughput, commit gets 1 row group with few thousand rows.

Input: 1 row group, 1 MB.

Output: swept_line_cluster → 1 cluster. bin_pack: 1 MB ≤ upper_bound_input_bytes_per_task → 1 chunk. tail_merge not applicable (no previous chunk). 1 parquet file 1 MB.

Downsides:

  • Small file (1 MB << lower_bound_input_bytes_per_task). Not a regression — old behavior gave same.
  • Long-term small files problem requires separate file compaction (out of scope for this task).

Case 10. Tail-merge: successful case (small tail)

Conditions: after bin_pack, last chunk falls short of lower_bound_input_bytes_per_task, and sum with previous fits in upper_bound_input_bytes_per_task.

Input: after bin_pack got 2 chunks — chunk 1 = 80 MB (one cluster took its chunk), chunk 2 = 30 MB (remainder).

Output:

  • chunk 2 (30 MB) < lower_bound_input_bytes_per_task (64 MB) → qualifies for tail_merge.
  • chunk 1 + chunk 2 = 110 MB ≤ upper_bound_input_bytes_per_task (128 MB) → merge chunk 2 into chunk 1.

1 parquet file 110 MB.

Downsides: none. Fewer files, final file in target range [lower_bound_input_bytes_per_task, upper_bound_input_bytes_per_task].

Case 11. Tail-merge blocked by upper_bound_input_bytes_per_task

Conditions: after bin_pack, last chunk is small, but sum with previous exceeds upper_bound_input_bytes_per_task.

Input: after bin_pack got 2 chunks — chunk 1 = 120 MB (one large cluster took its chunk), chunk 2 = 10 MB (remainder).

Output:

  • chunk 2 (10 MB) < lower_bound_input_bytes_per_task (64 MB) → qualifies for tail_merge by precondition.
  • chunk 1 + chunk 2 = 130 MB > upper_bound_input_bytes_per_task (128 MB) → don't merge.
  • i.e., file 1 = 120 MB, file 2 = 10 MB (remainder).
  • must additionally check same partition tuple. If different — don't merge.

2 parquet files: 120 MB and 10 MB.

Downsides:

  • Small 10 MB file preserved as tail. Trade-off to respect upper_bound_input_bytes_per_task.
  • Typical flow (without atomic clusters close to upper_bound_input_bytes_per_task) encounters rarely.

Case 12. Single WAL row group exceeds lower_bound_input_bytes_per_task (theoretical)

Conditions: WAL row group sized > lower_bound_input_bytes_per_task. Impossible in practice (WAL request never that big — max_bytes_per_flush ~8 MB), but algorithm should handle it.

Input: 1 row group sized 80 MB.

Output: swept_line_cluster → 1 cluster 80 MB. bin_pack: 80 MB ≤ upper_bound_input_bytes_per_task (128 MB) → own chunk without split_oversized. 1 parquet file ~80 MB.

Downsides:

  • File larger than lower_bound_input_bytes_per_task, but within upper_bound_input_bytes_per_task. Writer doesn't do rollover (writer_max_parquet_bytes = 256 MB).
  • If row group sized > upper_bound_input_bytes_per_task (e.g., 200 MB) → split_oversized works like Case 2.

Case 13. Multiple WAL segments, mixed services by time

Conditions: 3 WAL segments, each contains row groups for multiple services at different times (typical multi-producer ingest).

Input: WAL seg 1: [service_name-A@t10, service_name-B@t10]. Seg 2: [service_name-A@t20, service_name-B@t20]. Seg 3: [service_name-A@t30]. Total 30 MB.

Output: sort by composite sort key: first all row groups of service_name-A (by ts DESC), then service_name-B (by ts DESC). swept_line_cluster: 5 disjoint clusters (different services + monotonic time within each). bin_pack: 30 MB ≤ upper_bound_input_bytes_per_task → 1 chunk. tail_merge not applicable (no previous chunk). 1 parquet file 30 MB with bounds [service_name-A..service_name-B], row groups within ordered by composite sort key.

Downsides: file (30 MB) smaller than lower_bound_input_bytes_per_task (64 MB), but no more data in commit. Identical to old behavior in file count, plus output sorted by key → better for row-group-level pruning within file.

Case 14. Cancellation during sort/swept_line_cluster/bin_pack/tail_merge

Conditions: cancel_token signals when planner already collected entries and started algorithm.

Input: list of ~1000 row groups.

Output: existing cancellation-checks in plan_runner.run remain. Between algorithm phases (sort → swept_line_cluster → bin_pack → tail_merge) add checks. On cancel — complete_task with TaskStatus::Cancelled.

Downsides:

  • All algorithm phases — pure compute, very fast (microseconds per 1000 row groups). Cancellation latency doesn't suffer.
  • No preemption inside sort (sort_by non-preemptable). Acceptable.

Case 15. Pairwise-complete overlap within oversized cluster (degenerate)

Conditions: same as Case 2 (one service, oversized cluster > upper_bound_input_bytes_per_task), but all row groups have identical or nearly identical time-ranges. Happens under very high load and large pod count for one service. Accept overlapping ranges without min_bound drift.

Input: 30 WAL row groups (~6.7 MB each), all service_name=A, time-ranges [10..80], [10..80], [10..80], ... (or with drift < 1 second). Total 200 MB.

Output: swept_line_cluster merges into one atomic cluster 200 MB. split_oversized cuts greedily:

  • sub-chunk 1: 19 row groups × 6.7 MB = ~127 MB, bounds [10..80].
  • sub-chunk 2: 11 row groups × 6.7 MB = ~74 MB, bounds [10..80].

tail_merge: sub-chunk 2 (74 MB) >= lower_bound_input_bytes_per_task (64 MB) → don't merge.

Output: 2 parquet files with identical bounds [10..80].

Downsides:

  • File-level pruning by timestamp completely non-functional for this window — both files contain same range. Unlike Case 2 where min_bound drift gives narrow overlap zone at sub-chunk junction, here disjointness unachievable at all.
  • Metric shift_planner_oversized_clusters_total increments → signal for investigation root cause (retry storm, producer bug, multi-replica without clock skew).
  • In typical prod (multi-pod with natural jitter) unlikely. Solved at compaction level.

Activity

  1. self-assigned this
    on Apr 27, 2026
  2. added 3 commits that reference this issue on Apr 30, 2026
  3. added 5 commits that reference this issue on May 4, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions