Hello Peter Rozsa, Noemi Pap-Takacs, Impala Public Jenkins,
I'd like you to reexamine a change. Please visit
http://gerrit.cloudera.org:8080/24614
to look at the new patch set (#15).
Change subject: IMPALA-12587: Respect MAX_FS_WRITERS for Iceberg
DELETE/UPDATE/MERGE
......................................................................
IMPALA-12587: Respect MAX_FS_WRITERS for Iceberg DELETE/UPDATE/MERGE
Previously, MAX_FS_WRITERS only capped parallelism for HdfsTableSink
(INSERT). Iceberg DELETE (IcebergBufferedDeleteSink), UPDATE and MERGE
(MultiDataSink) ignored the limit, potentially creating too many small
files. The scheduler's IsExceedMaxFsWriters check also only recognized
HdfsTableSink, so multi-sink fragments were never capped at schedule
time.
This patch:
- Introduces the HasQuantityLimit interface with getNumInstances() and
getNumNodes(), replacing the instanceof cascade in
PlanFragment.getNumInstances()/getNumNodes(). HdfsTableSink,
JoinBuildSink, IcebergBufferedDeleteSink, and MultiDataSink all
implement it.
- Adds maxTableSinks_ to IcebergBufferedDeleteSink so
getNumInstances()/getNumNodes() cap writer parallelism via
MAX_FS_WRITERS.
- Makes MultiDataSink delegate getNumInstances()/getNumNodes() to the
minimum across child sinks that implement HasQuantityLimit.
- Simplifies MultiDataSink construction with a varargs constructor.
- Extends the scheduler's IsExceedMaxFsWriters to inspect
child_data_sinks, so multi-sink fragments (UPDATE/MERGE) are capped
at schedule time.
- Adds planner tests (PlannerTest/iceberg-writer-limit) and e2e tests
(TestIcebergWriterLimit) that validate instance counts and actual
file counts for INSERT, DELETE, UPDATE, and MERGE with
MAX_FS_WRITERS=1 vs. no limit.
- Minor cleanup in MaxRowsProcessedVisitor (unnecessary casts).
Change-Id: Ice7362cecb8b43fbff43b8ad82827e31e2b4eea3
Assisted-by: Claude Opus 4.6 (Claude Code)
---
M be/src/scheduling/scheduler.h
M fe/src/main/java/org/apache/impala/analysis/IcebergDeleteImpl.java
M fe/src/main/java/org/apache/impala/analysis/IcebergMergeImpl.java
M fe/src/main/java/org/apache/impala/analysis/IcebergUpdateImpl.java
A fe/src/main/java/org/apache/impala/planner/HasQuantityLimit.java
M fe/src/main/java/org/apache/impala/planner/HdfsTableSink.java
M fe/src/main/java/org/apache/impala/planner/IcebergBufferedDeleteSink.java
M fe/src/main/java/org/apache/impala/planner/IcebergMergeSink.java
M fe/src/main/java/org/apache/impala/planner/JoinBuildSink.java
M fe/src/main/java/org/apache/impala/planner/MultiDataSink.java
M fe/src/main/java/org/apache/impala/planner/PlanFragment.java
M fe/src/main/java/org/apache/impala/util/MaxRowsProcessedVisitor.java
M fe/src/test/java/org/apache/impala/catalog/IcebergContentFileStoreTest.java
M fe/src/test/java/org/apache/impala/planner/PlannerTest.java
A
testdata/workloads/functional-planner/queries/PlannerTest/iceberg-writer-limit.test
A
testdata/workloads/functional-query/queries/QueryTest/iceberg-writer-limit.test
M tests/query_test/test_iceberg.py
17 files changed, 598 insertions(+), 82 deletions(-)
git pull ssh://gerrit.cloudera.org:29418/Impala-ASF refs/changes/14/24614/15
--
To view, visit http://gerrit.cloudera.org:8080/24614
To unsubscribe, visit http://gerrit.cloudera.org:8080/settings
Gerrit-Project: Impala-ASF
Gerrit-Branch: master
Gerrit-MessageType: newpatchset
Gerrit-Change-Id: Ice7362cecb8b43fbff43b8ad82827e31e2b4eea3
Gerrit-Change-Number: 24614
Gerrit-PatchSet: 15
Gerrit-Owner: Nandor Kollar <[email protected]>
Gerrit-Reviewer: Impala Public Jenkins <[email protected]>
Gerrit-Reviewer: Nandor Kollar <[email protected]>
Gerrit-Reviewer: Noemi Pap-Takacs <[email protected]>
Gerrit-Reviewer: Peter Rozsa <[email protected]>