HIVE-30064: Report accurate affected-row counts for copy-on-write UPDATE and MERGE - #6794
ryukobayashi wants to merge 7 commits into
Conversation
a48b4cf to
5cb4bac
Compare
|
I don't think the marker column is the right layer for this. It adds a column to every row, survivors included, and carries it through the shuffle and the SerDe. It relies on the SerDe ignoring a trailing field. It uses a query-wide conf flag plus a positional "last field" lookup that applies to every FileSinkOperator in the DAG. It also changes what The information is already in the plan. t AS (
select <acid cols, ROW__POSITION := -cnt>, ...
from (
select ..., row_number() over (partition by file_path) rn,
count(*) over (partition by file_path) cnt
from target where <cond>
) q
where rn = 1
)The inner query adds the In if (positionDelete.pos() < 0) {
matchedRows += -positionDelete.pos(); // rows this operation replaced in that file
...
replacedDataFiles.add(dataFile);
} else {
writer.write(rowData, specs.get(currentSpecId), partition(rowData, currentSpecId));
}The count still has to reach // FileSinkOperator.closeOp
long affectedRows = numRows;
if (allWritersProvideAffectedRows()) {
affectedRows = sumAffectedRowsFromWriters(); // across all fpaths.outWriters
}
row_count.set(conf.isDeleteOfSplitUpdate() ? 0 : affectedRows);Summed over all markers, this gives the updated rows for UPDATE and the deleted rows for DELETE. CoW DELETE has the same miscount today (it reports survivors plus one row per file), and this PR doesn't address it. For MERGE it gives updated + deleted rows, because |
|
@deniskuzZ Thanks for the detailed review. I agree that adding a trailing marker column is not the right layer. I changed the implementation to use the existing negative ROW__POSITION marker instead. COWWithClauseBuilder now calculates count(*) over the file_path partition and emits one replacement marker per file with ROW__POSITION = -matched_count. HiveIcebergCopyOnWriteRecordWriter sums the absolute values of these negative positions and exposes the result through a dedicated affected-row writer capability. FileSinkOperator uses the affected-row count only when all output writers support the capability. Otherwise, it falls back to the existing physical output-row count, so regular INSERTs and non-CoW writers retain their previous behavior. This also avoids adding a column to the query projection or passing an extra field through the shuffle and SerDe. For MERGE, the affected-row count follows the semantics described in the review: it counts matched target rows, i.e. UPDATE + DELETE rows, and does not include INSERT rows. Therefore, the test case with two updates, one delete, and two inserts now expects 3 affected rows. The COW UPDATE and MERGE regression tests pass, and the affected Iceberg qtests were regenerated and pass with the new execution plans. |
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
The implementation, tests, interface contract, and PR description disagree on which MERGE and DELETE rows contribute to the reported count.
Review effort: Balanced
Findings: 2
Open (2)
What changed in this PR
Adds affected-row tracking for Iceberg copy-on-write DML by encoding matched-row counts in replacement markers and exposing writer-provided counts through FileSinkOperator.
Changes:
- Adds an affected-row reporting writer interface and counter aggregation.
- Encodes per-file matched counts in CoW replacement rows.
- Adds Iceberg operation-status tests and updates query-plan expectations.
| File | Description |
|---|---|
ql/.../NonNativeAcidMultiInsertSqlGenerator.java |
Supports custom deleted-row positions. |
ql/.../MultiInsertSqlGenerator.java |
Adds the corresponding generator API. |
ql/.../COWWithClauseBuilder.java |
Computes per-file matched-row counts. |
ql/.../CopyOnWriteMergeRewriter.java |
Contains formatting-only changes. |
ql/.../AffectedRowsProvidingRecordWriter.java |
Defines the affected-row provider interface. |
ql/.../FileSinkOperator.java |
Aggregates writer-provided affected-row counts. |
.../update_iceberg_cow_null_identity_partition.q.out |
Updates expected UPDATE plans. |
.../update_iceberg_copy_on_write_unpartitioned.q.out |
Updates unpartitioned UPDATE plans. |
.../update_iceberg_copy_on_write_partitioned.q.out |
Updates partitioned UPDATE plans. |
.../merge_with_null_check_on_joining_col.q.out |
Updates MERGE CBO plans. |
.../merge_iceberg_copy_on_write_unpartitioned.q.out |
Updates unpartitioned MERGE plans. |
.../merge_iceberg_copy_on_write_partitioned.q.out |
Updates partitioned MERGE plans. |
.../delete_iceberg_copy_on_write_unpartitioned.q.out |
Updates unpartitioned DELETE plans. |
.../delete_iceberg_copy_on_write_partitioned.q.out |
Updates partitioned DELETE plans. |
.../TestHiveIcebergCRUD.java |
Adds UPDATE and MERGE row-count tests. |
.../TestHiveShell.java |
Adds an operation-status helper. |
.../HiveIcebergCopyOnWriteRecordWriter.java |
Accumulates counts from replacement markers. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| "WHEN MATCHED THEN UPDATE SET b = 'merged' " + | ||
| "WHEN NOT MATCHED THEN INSERT VALUES (src.a, src.b)"); | ||
|
|
||
| Assert.assertEquals(3, numModifiedRows); |
| sqlGenerator.append(", row_number() OVER (partition by ").append(filePathCol).append(") rn"); | ||
| sqlGenerator.append(", count(*) OVER (partition by ").append(filePathCol).append(") ") | ||
| .append(matchedRowCountCol); |
|




What changes were proposed in this pull request?
This patch adds matched-target-row tracking for copy-on-write
UPDATE,DELETE, andMERGEoperations.The CoW rewriters reuse the existing negative
ROW__POSITIONreplacement marker. For each rewritten data file,COWWithClauseBuildercomputes the number of matching target rows and emits one marker withROW__POSITION = -matched_count. No extra projection column is added, so the marker does not pass through the shuffle or SerDe as an additional field.HiveIcebergCopyOnWriteRecordWritersums the absolute values of the negative positions and exposes the result throughAffectedRowsProvidingRecordWriter.FileSinkOperatoruses the writer-provided count only when all output writers support the capability; otherwise it falls back to the existing physical output-row count.For
MERGE, the reported count is the number of matched target rows: matchedUPDATErows plus matchedDELETErows.INSERTrows are intentionally excluded from this count.Why are the changes needed?
Copy-on-write operations rewrite complete data files. Therefore, the rows sent to the
FileSinkOperatorinclude both modified rows and unchanged rows copied from the rewritten files.For
MERGE, a single output also combines multiple logical branches, including matched updates, not-matched inserts, deletes, and unchanged survivor rows.Previously,
FileSinkOperatorcounted all output rows, so Hive could report an incorrectnumModifiedRowsvalue. The reported count could include unchanged rows or rows from otherMERGEbranches.The marker makes the logical operation explicit and avoids relying on the physical Tez plan shape or operator layout to determine the count.
Does this PR introduce any user-facing change?
Yes.
The reported affected-row count for copy-on-write
UPDATE,DELETE, andMERGEoperations is corrected to reflect matched target rows rather than all rows rewritten by the CoW plan.The underlying table data and write behavior are unchanged.
For example, a CoW
MERGEthat updates two rows, inserts two rows, and deletes one row reports3matched target rows (2updates +1delete). The two inserted rows are excluded.How was this patch tested?
Added Iceberg CoW regression tests for:
UPDATEwith two matching rowsDELETEwith two matching rowsMERGEwith two matched updates, two inserts, and one deleteThe tests retrieve
numModifiedRowsthrough Hive operation status and verify the expected count. The affected Iceberg qtests were also updated for the new execution plans.