opensearch-project / opensearch-project/data-prepper

[FEATURE] Group SHUFFLE_WRITE tasks by file size in iceberg-source

Open
#6,724 0 comments 0 reactions 1 assignee View on GitHub

@lawofcycles is already working on this.

Since Apr 7, 2026.

performance
Dominant language
Java
Stars
374
Forks
354
Avg merge
3d 18h
Merged PRs (30d)
8

Description

Is your feature request related to a problem? Please describe.

The current iceberg-source creates one SHUFFLE_WRITE task per data file. When a large UPDATE/DELETE operation affects many data files (e.g. 500K rows scattered across thousands of files in a Copy on Write table), the task count can reach thousands. Source coordination stores all SHUFFLE_WRITE tasks under the same DynamoDB partition key, which has a per partition write throughput limit of 1,000 WCU/s regardless of provisioned capacity. This causes write throttling that degrades performance or stalls processing entirely.

Observed during performance testing with NYC Yellow Taxi (41 million rows, partitioned by day):

SHUFFLE_WRITE task count Result
~200 Completed normally
~2,000 Completed with throttling (3 to 7 minutes)
~5,000 Stalled

Additionally, when individual data files are small (e.g. median 20KB from Spark parallel writes), the coordination overhead per task dominates the actual file processing time, making the per file task granularity inefficient regardless of the coordination backend.

Describe the solution you'd like

Group multiple data files into a single SHUFFLE_WRITE task based on total file size, so that each task represents a meaningful amount of work. fileSizeInBytes is a required field in Iceberg file metadata, so it is available at planning time without reading the actual files.

Additional context

Performance test results: https://github.com/opensearch-project/data-prepper/pull/6682#issuecomment-4185591667
Related PR: #6682 (source-layer shuffle implementation)

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.