From 978c396f7e65a5d54477a716a7eff96219d4846a Mon Sep 17 00:00:00 2001 From: yufei Date: Mon, 10 Apr 2023 17:19:29 -0700 Subject: [PATCH 01/20] Output the net changes across snapshots in CDC --- .../TestCreateChangelogViewProcedure.java | 29 +++++ .../iceberg/spark/ChangelogIterator.java | 5 + .../iceberg/spark/NetChangelogIterator.java | 121 ++++++++++++++++++ .../spark/RemoveCarryoverIterator.java | 2 +- .../CreateChangelogViewProcedure.java | 44 +++++-- .../iceberg/spark/TestChangelogIterator.java | 95 ++++++++++---- 6 files changed, 259 insertions(+), 37 deletions(-) create mode 100644 spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/NetChangelogIterator.java diff --git a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java index dc12b0145d50..133e9835da3a 100644 --- a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java +++ b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java @@ -253,6 +253,35 @@ public void testUpdate() { sql("select * from %s order by _change_ordinal, id, data", viewName)); } + @Test + public void testNetChanges() { + createTableWith2Columns(); + sql("ALTER TABLE %s DROP PARTITION FIELD data", tableName); + sql("ALTER TABLE %s ADD PARTITION FIELD id", tableName); + + sql("INSERT INTO %s VALUES (1, 'a'), (2, 'b')", tableName); + Table table = validationCatalog.loadTable(tableIdent); + Snapshot snap1 = table.currentSnapshot(); + + sql("INSERT OVERWRITE %s VALUES (3, 'c'), (2, 'd')", tableName); + table.refresh(); + Snapshot snap2 = table.currentSnapshot(); + + List returns = + sql( + "CALL %s.system.create_changelog_view(table => '%s', identifier_columns => array('id'), net_changes => true)", + catalogName, tableName); + + String viewName = (String) returns.get(0)[0]; + assertEquals( + "Rows should match", + ImmutableList.of( + row(1, "a", INSERT, 0, snap1.snapshotId()), + row(2, "d", INSERT, 1, snap2.snapshotId()), + row(3, "c", INSERT, 1, snap2.snapshotId())), + sql("select * from %s order by _change_ordinal, id, data", viewName)); + } + @Test public void testUpdateWithIdentifierField() { createTableWithIdentifierField(); diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java index 5f30c5fd4e56..6932bade96f5 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java @@ -79,6 +79,11 @@ public static Iterator removeCarryovers(Iterator rowIterator, StructTy return Iterators.filter(changelogIterator, Objects::nonNull); } + public static Iterator netChanges(Iterator rowIterator, StructType rowType) { + ChangelogIterator changelogIterator = new NetChangelogIterator(rowIterator, rowType); + return Iterators.filter(changelogIterator, Objects::nonNull); + } + protected boolean isDifferentValue(Row currentRow, Row nextRow, int idx) { return !Objects.equals(nextRow.get(idx), currentRow.get(idx)); } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/NetChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/NetChangelogIterator.java new file mode 100644 index 000000000000..ab5fc47197d8 --- /dev/null +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/NetChangelogIterator.java @@ -0,0 +1,121 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark; + +import java.util.Iterator; +import java.util.List; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.types.StructType; + +/** + * This class is used to compute the net changes between multiple snapshots. It takes a row + * iterator, and returns the net changes between multiple snapshots. It assumes the following: + * + *
    + *
  • The row iterator is partitioned by the primary key. + *
  • The row iterator is sorted by the primary key, change order, and change type. The change + * order is 1-to-1 mapping to snapshot id. + *
+ */ +public class NetChangelogIterator extends RemoveCarryoverIterator { + + private Row cachedNextRow = null; + private final List cachedRows = Lists.newArrayList(); + + protected NetChangelogIterator(Iterator rowIterator, StructType rowType) { + super(rowIterator, rowType); + } + + @Override + public boolean hasNext() { + if(!cachedRows.isEmpty()) { + return true; + } + + if(cachedNextRow != null) { + return true; + } + + return rowIterator().hasNext(); + } + + @Override + public Row next() { + // if there are cached rows, return one of them from the beginning + if (!cachedRows.isEmpty()) { + return cachedRows.remove(0); + } + + Row currentRow = getCurrentRow(); + + // if there is no next row, return the current row + if (!rowIterator().hasNext()) { + return currentRow; + } + + Row nextRow = rowIterator().next(); + + cachedRows.add(currentRow); + // if they are the same row, remove the cached row or stack to the cached row + while (isSameRecord(currentRow, nextRow)) { + if (matched(currentRow, nextRow)) { + // if matched, remove both rows, remove the last row in the cache + cachedRows.remove(cachedRows.size() - 1); + nextRow = null; + } else { + // stack the row into the cache + cachedRows.add(nextRow); + nextRow = null; + } + + // break the loop if there is no next row or the cache is empty + if (cachedRows.isEmpty() || !rowIterator().hasNext()) { + break; + } + + // get the current row from the cache, and get the next row from the iterator for the next + // loop + currentRow = cachedRows.get(cachedRows.size() - 1); + nextRow = rowIterator().next(); + } + + // if they are different rows, hit the boundary, cache the next row + cachedNextRow = nextRow; + return null; + } + + private Row getCurrentRow() { + Row currentRow; + if (cachedNextRow != null) { + currentRow = cachedNextRow; + cachedNextRow = null; + } else { + currentRow = rowIterator().next(); + } + return currentRow; + } + + private boolean matched(Row currentRow, Row nextRow) { + return (nextRow.getString(changeTypeIndex()).equals(INSERT) + && currentRow.getString(changeTypeIndex()).equals(DELETE)) + || (nextRow.getString(changeTypeIndex()).equals(DELETE) + && currentRow.getString(changeTypeIndex()).equals(INSERT)); + } +} diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java index 70b160e13fee..2b33ed31281c 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java @@ -151,7 +151,7 @@ private int[] generateIndicesToIdentifySameRow(int columnSize) { return indices; } - private boolean isSameRecord(Row currentRow, Row nextRow) { + protected boolean isSameRecord(Row currentRow, Row nextRow) { for (int idx : indicesToIdentifySameRow) { if (isDifferentValue(currentRow, nextRow, idx)) { return false; diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index 85043d2df3d6..d4b30c36bcca 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -92,6 +92,8 @@ public class CreateChangelogViewProcedure extends BaseProcedure { ProcedureParameter.optional("remove_carryovers", DataTypes.BooleanType); private static final ProcedureParameter IDENTIFIER_COLUMNS_PARAM = ProcedureParameter.optional("identifier_columns", STRING_ARRAY); + private static final ProcedureParameter NET_CHANGES = + ProcedureParameter.optional("net_changes", DataTypes.BooleanType); private static final ProcedureParameter[] PARAMETERS = new ProcedureParameter[] { @@ -101,6 +103,7 @@ public class CreateChangelogViewProcedure extends BaseProcedure { COMPUTE_UPDATES_PARAM, REMOVE_CARRYOVERS_PARAM, IDENTIFIER_COLUMNS_PARAM, + NET_CHANGES, }; private static final StructType OUTPUT_TYPE = @@ -142,10 +145,12 @@ public InternalRow[] call(InternalRow args) { Identifier changelogTableIdent = changelogTableIdent(tableIdent); Dataset df = loadRows(changelogTableIdent, options(input)); + boolean netChanges = input.asBoolean(NET_CHANGES, false); + if (shouldComputeUpdateImages(input)) { - df = computeUpdateImages(identifierColumns(input, tableIdent), df); + df = computeUpdateImages(identifierColumns(input, tableIdent), df, netChanges); } else if (shouldRemoveCarryoverRows(input)) { - df = removeCarryoverRows(df); + df = removeCarryoverRows(df, netChanges); } String viewName = viewName(input, tableIdent.name()); @@ -155,18 +160,23 @@ public InternalRow[] call(InternalRow args) { return toOutputRows(viewName); } - private Dataset computeUpdateImages(String[] identifierColumns, Dataset df) { + private Dataset computeUpdateImages( + String[] identifierColumns, Dataset df, boolean netChanges) { Preconditions.checkArgument( identifierColumns.length > 0, "Cannot compute the update images because identifier columns are not set"); - Column[] repartitionSpec = new Column[identifierColumns.length + 1]; + int length = netChanges ? identifierColumns.length : identifierColumns.length + 1; + Column[] repartitionSpec = new Column[length]; for (int i = 0; i < identifierColumns.length; i++) { repartitionSpec[i] = df.col(identifierColumns[i]); } - repartitionSpec[repartitionSpec.length - 1] = df.col(MetadataColumns.CHANGE_ORDINAL.name()); - return applyChangelogIterator(df, repartitionSpec); + if (!netChanges) { + repartitionSpec[repartitionSpec.length - 1] = df.col(MetadataColumns.CHANGE_ORDINAL.name()); + } + + return applyChangelogIterator(df, repartitionSpec, netChanges); } private boolean shouldComputeUpdateImages(ProcedureInput input) { @@ -179,7 +189,7 @@ private boolean shouldRemoveCarryoverRows(ProcedureInput input) { return input.asBoolean(REMOVE_CARRYOVERS_PARAM, true); } - private Dataset removeCarryoverRows(Dataset df) { + private Dataset removeCarryoverRows(Dataset df, boolean netChanges) { Column[] repartitionSpec = Arrays.stream(df.columns()) .filter(c -> !c.equals(MetadataColumns.CHANGE_TYPE.name())) @@ -213,7 +223,8 @@ private String viewName(ProcedureInput input, String tableName) { return input.asString(CHANGELOG_VIEW_PARAM, defaultValue); } - private Dataset applyChangelogIterator(Dataset df, Column[] repartitionSpec) { + private Dataset applyChangelogIterator( + Dataset df, Column[] repartitionSpec, boolean netChanges) { Column[] sortSpec = sortSpec(df, repartitionSpec); StructType schema = df.schema(); String[] identifierFields = @@ -229,6 +240,7 @@ private Dataset applyChangelogIterator(Dataset df, Column[] repartitio } private Dataset applyCarryoverRemoveIterator(Dataset df, Column[] repartitionSpec) { + // todo double check we should sort by (change order, change type) for net changes Column[] sortSpec = sortSpec(df, repartitionSpec); StructType schema = df.schema(); @@ -241,8 +253,22 @@ private Dataset applyCarryoverRemoveIterator(Dataset df, Column[] repa } private static Column[] sortSpec(Dataset df, Column[] repartitionSpec) { - Column[] sortSpec = new Column[repartitionSpec.length + 1]; + Column changeOrdinal = df.col(MetadataColumns.CHANGE_ORDINAL.name()); + boolean noChangeOrdinal = + Arrays.stream(repartitionSpec).noneMatch(c -> c.equals(changeOrdinal)); + + Column[] sortSpec; + if (noChangeOrdinal) { + sortSpec = new Column[repartitionSpec.length + 2]; + } else { + sortSpec = new Column[repartitionSpec.length + 1]; + } + System.arraycopy(repartitionSpec, 0, sortSpec, 0, repartitionSpec.length); + + if (noChangeOrdinal) { + sortSpec[sortSpec.length - 2] = changeOrdinal; + } sortSpec[sortSpec.length - 1] = df.col(MetadataColumns.CHANGE_TYPE.name()); return sortSpec; } diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java index bf98bebb9d50..a1cebe672d55 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java @@ -228,20 +228,15 @@ public void testCarryRowsRemoveWithDuplicates() { new GenericRowWithSchema(new Object[] {2, "a", "data", INSERT}, null), new GenericRowWithSchema(new Object[] {3, "a", "new_data", INSERT}, null)); - Iterator iterator = - ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); - List result = Lists.newArrayList(iterator); - - assertEquals( - "Rows should match.", - Lists.newArrayList( + List expectedRows = Lists.newArrayList( new Object[] {0, "a", "data", DELETE}, new Object[] {0, "a", "data", DELETE}, new Object[] {0, "a", "data", DELETE}, new Object[] {1, "a", "old_data", DELETE}, new Object[] {1, "a", "old_data", DELETE}, - new Object[] {3, "a", "new_data", INSERT}), - rowsToJava(result)); + new Object[] {3, "a", "new_data", INSERT}); + + validateIterators(rowsWithDuplication, expectedRows); } @Test @@ -254,15 +249,9 @@ public void testCarryRowsRemoveLessInsertRows() { new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), new GenericRowWithSchema(new Object[] {2, "d", "data", INSERT}, null)); - Iterator iterator = - ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); - List result = Lists.newArrayList(iterator); + List expectedRows = Lists.newArrayList( new Object[] {1, "d", "data", DELETE}, new Object[] {2, "d", "data", INSERT}); - assertEquals( - "Rows should match.", - Lists.newArrayList( - new Object[] {1, "d", "data", DELETE}, new Object[] {2, "d", "data", INSERT}), - rowsToJava(result)); + validateIterators(rowsWithDuplication, expectedRows); } @Test @@ -279,15 +268,10 @@ public void testCarryRowsRemoveMoreInsertRows() { new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null)); - Iterator iterator = - ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); - List result = Lists.newArrayList(iterator); + List expectedRows = Lists.newArrayList( + new Object[] {0, "d", "data", DELETE}, new Object[] {1, "d", "data", INSERT}); - assertEquals( - "Rows should match.", - Lists.newArrayList( - new Object[] {0, "d", "data", DELETE}, new Object[] {1, "d", "data", INSERT}), - rowsToJava(result)); + validateIterators(rowsWithDuplication, expectedRows); } @Test @@ -299,14 +283,71 @@ public void testCarryRowsRemoveNoInsertRows() { new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null)); + List expectedRows = Lists.newArrayList( + new Object[] {1, "d", "data", DELETE}, new Object[] {1, "d", "data", DELETE}); + + validateIterators(rowsWithDuplication, expectedRows); + } + + private void validateIterators(List rowsWithDuplication, List expectedRows) { Iterator iterator = - ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); + ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); + List result = Lists.newArrayList(iterator); + + assertEquals("Rows should match.", expectedRows, rowsToJava(result)); + + iterator = ChangelogIterator.netChanges(rowsWithDuplication.iterator(), SCHEMA); + result = Lists.newArrayList(iterator); + + assertEquals("Rows should match.", expectedRows, rowsToJava(result)); + } + + @Test + public void testNetChanges() { + List rowsWithDuplication = + Lists.newArrayList( + // 1. the first version is delete and last version is insert, then output pre/post + // images + new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT}, null), + // 2. the first version is delete and last version is delete, output the first version + new GenericRowWithSchema(new Object[] {1, "b", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "b", "new_data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "b", "new_data", DELETE}, null), + // 3. the first version is insert and last version is delete, output nothing + new GenericRowWithSchema(new Object[] {1, "c", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "c", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "c", "new_data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "c", "new_data", DELETE}, null), + // 4. the first version is insert and last version is insert, output the last version + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "d", "new_data", INSERT}, null), + // 5. two rows are identical, remove one of them, and output pre/post images + new GenericRowWithSchema(new Object[] {1, "e", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "e", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "e", "new_data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "e", "new_data", INSERT}, null), + // 6. two rows are identical, and carry-over rows, output nothing + new GenericRowWithSchema(new Object[] {1, "f", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "f", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "f", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "f", "data", INSERT}, null)); + + Iterator iterator = ChangelogIterator.netChanges(rowsWithDuplication.iterator(), SCHEMA); List result = Lists.newArrayList(iterator); assertEquals( "Duplicate rows should not be removed", Lists.newArrayList( - new Object[] {1, "d", "data", DELETE}, new Object[] {1, "d", "data", DELETE}), + new Object[] {1, "a", "data", UPDATE_BEFORE}, + new Object[] {1, "a", "new_data", UPDATE_AFTER}, + new Object[] {1, "b", "data", DELETE}, + new Object[] {1, "d", "new_data", INSERT}, + new Object[] {1, "e", "data", UPDATE_BEFORE}, + new Object[] {1, "e", "new_data", UPDATE_AFTER}), rowsToJava(result)); } } From ccdcde95a2017898dc674d8f1ebabb9c64bc56d2 Mon Sep 17 00:00:00 2001 From: yufei Date: Wed, 17 May 2023 18:21:07 -0700 Subject: [PATCH 02/20] add tests --- .../iceberg/spark/ChangelogIterator.java | 2 +- ...r.java => RemoveNetCarryoverIterator.java} | 12 +-- .../iceberg/spark/TestChangelogIterator.java | 79 ++++++++----------- 3 files changed, 42 insertions(+), 51 deletions(-) rename spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/{NetChangelogIterator.java => RemoveNetCarryoverIterator.java} (89%) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java index 6932bade96f5..10271a708aeb 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java @@ -80,7 +80,7 @@ public static Iterator removeCarryovers(Iterator rowIterator, StructTy } public static Iterator netChanges(Iterator rowIterator, StructType rowType) { - ChangelogIterator changelogIterator = new NetChangelogIterator(rowIterator, rowType); + ChangelogIterator changelogIterator = new RemoveNetCarryoverIterator(rowIterator, rowType); return Iterators.filter(changelogIterator, Objects::nonNull); } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/NetChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java similarity index 89% rename from spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/NetChangelogIterator.java rename to spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index ab5fc47197d8..be99881b1105 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/NetChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -25,8 +25,8 @@ import org.apache.spark.sql.types.StructType; /** - * This class is used to compute the net changes between multiple snapshots. It takes a row - * iterator, and returns the net changes between multiple snapshots. It assumes the following: + * This class computes the net changes across multiple snapshots. It takes a row iterator, and + * assumes the following: * *
    *
  • The row iterator is partitioned by the primary key. @@ -34,22 +34,22 @@ * order is 1-to-1 mapping to snapshot id. *
*/ -public class NetChangelogIterator extends RemoveCarryoverIterator { +public class RemoveNetCarryoverIterator extends RemoveCarryoverIterator { private Row cachedNextRow = null; private final List cachedRows = Lists.newArrayList(); - protected NetChangelogIterator(Iterator rowIterator, StructType rowType) { + protected RemoveNetCarryoverIterator(Iterator rowIterator, StructType rowType) { super(rowIterator, rowType); } @Override public boolean hasNext() { - if(!cachedRows.isEmpty()) { + if (!cachedRows.isEmpty()) { return true; } - if(cachedNextRow != null) { + if (cachedNextRow != null) { return true; } diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java index a1cebe672d55..0d28f82399be 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java @@ -228,7 +228,8 @@ public void testCarryRowsRemoveWithDuplicates() { new GenericRowWithSchema(new Object[] {2, "a", "data", INSERT}, null), new GenericRowWithSchema(new Object[] {3, "a", "new_data", INSERT}, null)); - List expectedRows = Lists.newArrayList( + List expectedRows = + Lists.newArrayList( new Object[] {0, "a", "data", DELETE}, new Object[] {0, "a", "data", DELETE}, new Object[] {0, "a", "data", DELETE}, @@ -249,7 +250,9 @@ public void testCarryRowsRemoveLessInsertRows() { new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), new GenericRowWithSchema(new Object[] {2, "d", "data", INSERT}, null)); - List expectedRows = Lists.newArrayList( new Object[] {1, "d", "data", DELETE}, new Object[] {2, "d", "data", INSERT}); + List expectedRows = + Lists.newArrayList( + new Object[] {1, "d", "data", DELETE}, new Object[] {2, "d", "data", INSERT}); validateIterators(rowsWithDuplication, expectedRows); } @@ -268,7 +271,8 @@ public void testCarryRowsRemoveMoreInsertRows() { new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null)); - List expectedRows = Lists.newArrayList( + List expectedRows = + Lists.newArrayList( new Object[] {0, "d", "data", DELETE}, new Object[] {1, "d", "data", INSERT}); validateIterators(rowsWithDuplication, expectedRows); @@ -283,15 +287,16 @@ public void testCarryRowsRemoveNoInsertRows() { new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null)); - List expectedRows = Lists.newArrayList( + List expectedRows = + Lists.newArrayList( new Object[] {1, "d", "data", DELETE}, new Object[] {1, "d", "data", DELETE}); - validateIterators(rowsWithDuplication, expectedRows); + validateIterators(rowsWithDuplication, expectedRows); } private void validateIterators(List rowsWithDuplication, List expectedRows) { Iterator iterator = - ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); + ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); List result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); @@ -303,51 +308,37 @@ private void validateIterators(List rowsWithDuplication, List exp } @Test - public void testNetChanges() { + public void testRemoveNetCarryovers() { List rowsWithDuplication = Lists.newArrayList( - // 1. the first version is delete and last version is insert, then output pre/post - // images - new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT}, null), - // 2. the first version is delete and last version is delete, output the first version - new GenericRowWithSchema(new Object[] {1, "b", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "b", "new_data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "b", "new_data", DELETE}, null), - // 3. the first version is insert and last version is delete, output nothing - new GenericRowWithSchema(new Object[] {1, "c", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "c", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "c", "new_data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "c", "new_data", DELETE}, null), - // 4. the first version is insert and last version is insert, output the last version + new GenericRowWithSchema(new Object[] {0, "d", "data", DELETE}, null), + // a pair of delete and insert rows + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + // 2 delete rows and 2 insert rows + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "new_data", INSERT}, null), - // 5. two rows are identical, remove one of them, and output pre/post images - new GenericRowWithSchema(new Object[] {1, "e", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "e", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "e", "new_data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "e", "new_data", INSERT}, null), - // 6. two rows are identical, and carry-over rows, output nothing - new GenericRowWithSchema(new Object[] {1, "f", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "f", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "f", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "f", "data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + // a pair of insert and delete rows + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), + // extra insert rows + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + // different key + new GenericRowWithSchema(new Object[] {2, "d", "data", DELETE}, null)); + + List expectedRows = + Lists.newArrayList( + new Object[] {0, "d", "data", DELETE}, + new Object[] {1, "d", "data", INSERT}, + new Object[] {1, "d", "data", INSERT}, + new Object[] {2, "d", "data", DELETE}); Iterator iterator = ChangelogIterator.netChanges(rowsWithDuplication.iterator(), SCHEMA); List result = Lists.newArrayList(iterator); - assertEquals( - "Duplicate rows should not be removed", - Lists.newArrayList( - new Object[] {1, "a", "data", UPDATE_BEFORE}, - new Object[] {1, "a", "new_data", UPDATE_AFTER}, - new Object[] {1, "b", "data", DELETE}, - new Object[] {1, "d", "new_data", INSERT}, - new Object[] {1, "e", "data", UPDATE_BEFORE}, - new Object[] {1, "e", "new_data", UPDATE_AFTER}), - rowsToJava(result)); + assertEquals("Rows should match.", expectedRows, rowsToJava(result)); } } From dd95e5fbaf2f6d8e1abf61ca0334ef616c99b251 Mon Sep 17 00:00:00 2001 From: yufei Date: Thu, 18 May 2023 17:06:39 -0700 Subject: [PATCH 03/20] add tests --- .../TestCreateChangelogViewProcedure.java | 79 ++++++++++++------- .../iceberg/spark/ChangelogIterator.java | 2 +- .../spark/RemoveCarryoverIterator.java | 18 ++++- .../spark/RemoveNetCarryoverIterator.java | 32 +++++++- .../CreateChangelogViewProcedure.java | 57 +++++++------ .../iceberg/spark/TestChangelogIterator.java | 5 +- 6 files changed, 131 insertions(+), 62 deletions(-) diff --git a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java index 133e9835da3a..e105df183961 100644 --- a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java +++ b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java @@ -18,6 +18,8 @@ */ package org.apache.iceberg.spark.extensions; +import static org.junit.Assert.assertThrows; + import java.util.List; import java.util.Map; import org.apache.iceberg.ChangelogOperation; @@ -253,35 +255,6 @@ public void testUpdate() { sql("select * from %s order by _change_ordinal, id, data", viewName)); } - @Test - public void testNetChanges() { - createTableWith2Columns(); - sql("ALTER TABLE %s DROP PARTITION FIELD data", tableName); - sql("ALTER TABLE %s ADD PARTITION FIELD id", tableName); - - sql("INSERT INTO %s VALUES (1, 'a'), (2, 'b')", tableName); - Table table = validationCatalog.loadTable(tableIdent); - Snapshot snap1 = table.currentSnapshot(); - - sql("INSERT OVERWRITE %s VALUES (3, 'c'), (2, 'd')", tableName); - table.refresh(); - Snapshot snap2 = table.currentSnapshot(); - - List returns = - sql( - "CALL %s.system.create_changelog_view(table => '%s', identifier_columns => array('id'), net_changes => true)", - catalogName, tableName); - - String viewName = (String) returns.get(0)[0]; - assertEquals( - "Rows should match", - ImmutableList.of( - row(1, "a", INSERT, 0, snap1.snapshotId()), - row(2, "d", INSERT, 1, snap2.snapshotId()), - row(3, "c", INSERT, 1, snap2.snapshotId())), - sql("select * from %s order by _change_ordinal, id, data", viewName)); - } - @Test public void testUpdateWithIdentifierField() { createTableWithIdentifierField(); @@ -440,6 +413,54 @@ public void testRemoveCarryOversWithoutUpdatedRows() { sql("select * from %s order by _change_ordinal, id, data", viewName)); } + @Test + public void testNetChangesWithRemoveCarryOvers() { + // partitioned by id + createTableWith3Columns(); + + // insert rows: (1, 'a', 12) (2, 'b', 11) (2, 'e', 12) + sql("INSERT INTO %s VALUES (1, 'a', 12), (2, 'b', 11), (2, 'e', 12)", tableName); + Table table = validationCatalog.loadTable(tableIdent); + Snapshot snap1 = table.currentSnapshot(); + + // delete rows: (2, 'b', 11) (2, 'e', 12), insert rows: (3, 'c', 13) (2, 'd', 11) (2, 'e', 12) + sql("INSERT OVERWRITE %s VALUES (3, 'c', 13), (2, 'd', 11), (2, 'e', 12)", tableName); + table.refresh(); + Snapshot snap2 = table.currentSnapshot(); + + // delete rows: (2, 'd', 11) (2, 'e', 12) (3, 'c', 13), insert rows: (3, 'c', 15) (2, 'e', 12) + sql("INSERT OVERWRITE %s VALUES (3, 'c', 15), (2, 'e', 12)", tableName); + table.refresh(); + Snapshot snap3 = table.currentSnapshot(); + + List returns = + sql( + "CALL %s.system.create_changelog_view(table => '%s', net_changes => true)", + catalogName, tableName); + + String viewName = (String) returns.get(0)[0]; + + assertEquals( + "Rows should match", + ImmutableList.of( + row(1, "a", 12, INSERT, 0, snap1.snapshotId()), + row(3, "c", 15, INSERT, 2, snap3.snapshotId()), + row(2, "e", 12, INSERT, 2, snap3.snapshotId())), + sql("select * from %s order by _change_ordinal, data", viewName)); + } + + @Test + public void testNetChangesWithComputeUpdates() { + createTableWith2Columns(); + assertThrows( + "Should fail because net_changes is not supported with computing updates", + IllegalArgumentException.class, + () -> + sql( + "CALL %s.system.create_changelog_view(table => '%s', identifier_columns => array('id'), net_changes => true)", + catalogName, tableName)); + } + @Test public void testNotRemoveCarryOvers() { createTableWith3Columns(); diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java index 10271a708aeb..e848e61d6204 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java @@ -79,7 +79,7 @@ public static Iterator removeCarryovers(Iterator rowIterator, StructTy return Iterators.filter(changelogIterator, Objects::nonNull); } - public static Iterator netChanges(Iterator rowIterator, StructType rowType) { + public static Iterator removeNetCarryovers(Iterator rowIterator, StructType rowType) { ChangelogIterator changelogIterator = new RemoveNetCarryoverIterator(rowIterator, rowType); return Iterators.filter(changelogIterator, Objects::nonNull); } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java index 2b33ed31281c..b1431c2b879c 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java @@ -48,6 +48,7 @@ */ class RemoveCarryoverIterator extends ChangelogIterator { private final int[] indicesToIdentifySameRow; + private final StructType rowType; private Row cachedDeletedRow = null; private long deletedRowCount = 0; @@ -55,7 +56,8 @@ class RemoveCarryoverIterator extends ChangelogIterator { RemoveCarryoverIterator(Iterator rowIterator, StructType rowType) { super(rowIterator, rowType); - this.indicesToIdentifySameRow = generateIndicesToIdentifySameRow(rowType.size()); + this.rowType = rowType; + this.indicesToIdentifySameRow = generateIndicesToIdentifySameRow(); } @Override @@ -139,8 +141,8 @@ private boolean hasCachedDeleteRow() { return cachedDeletedRow != null; } - private int[] generateIndicesToIdentifySameRow(int columnSize) { - int[] indices = new int[columnSize - 1]; + private int[] generateIndicesToIdentifySameRow() { + int[] indices = new int[rowType.size() - 1]; for (int i = 0; i < indices.length; i++) { if (i < changeTypeIndex()) { indices[i] = i; @@ -152,7 +154,7 @@ private int[] generateIndicesToIdentifySameRow(int columnSize) { } protected boolean isSameRecord(Row currentRow, Row nextRow) { - for (int idx : indicesToIdentifySameRow) { + for (int idx : indicesToIdentifySameRow()) { if (isDifferentValue(currentRow, nextRow, idx)) { return false; } @@ -160,4 +162,12 @@ protected boolean isSameRecord(Row currentRow, Row nextRow) { return true; } + + protected StructType rowType() { + return rowType; + } + + protected int[] indicesToIdentifySameRow() { + return indicesToIdentifySameRow; + } } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index be99881b1105..84f48ee63735 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -20,6 +20,7 @@ import java.util.Iterator; import java.util.List; +import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.StructType; @@ -29,18 +30,21 @@ * assumes the following: * *
    - *
  • The row iterator is partitioned by the primary key. - *
  • The row iterator is sorted by the primary key, change order, and change type. The change - * order is 1-to-1 mapping to snapshot id. + *
  • The row iterator is partitioned by all columns. + *
  • The row iterator is sorted by all columns, change order, and change type. The change order + * is 1-to-1 mapping to snapshot id. *
*/ public class RemoveNetCarryoverIterator extends RemoveCarryoverIterator { - private Row cachedNextRow = null; private final List cachedRows = Lists.newArrayList(); + private final int[] indicesToIdentifySameRow; + + private Row cachedNextRow = null; protected RemoveNetCarryoverIterator(Iterator rowIterator, StructType rowType) { super(rowIterator, rowType); + this.indicesToIdentifySameRow = generateIndicesToIdentifySameRow(); } @Override @@ -118,4 +122,24 @@ private boolean matched(Row currentRow, Row nextRow) { || (nextRow.getString(changeTypeIndex()).equals(DELETE) && currentRow.getString(changeTypeIndex()).equals(INSERT)); } + + private int[] generateIndicesToIdentifySameRow() { + int changeOrdinalIndex = rowType().fieldIndex(MetadataColumns.CHANGE_ORDINAL.name()); + int snapshotIdIndex = rowType().fieldIndex(MetadataColumns.COMMIT_SNAPSHOT_ID.name()); + + int[] indices = new int[rowType().size() - 3]; + + for (int i = 0, j = 0; i < indices.length; i++) { + if (i != changeTypeIndex() && i != changeOrdinalIndex && i != snapshotIdIndex) { + indices[j] = i; + j++; + } + } + return indices; + } + + @Override + protected int[] indicesToIdentifySameRow() { + return indicesToIdentifySameRow; + } } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index d4b30c36bcca..c16a082c44de 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -21,6 +21,7 @@ import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.function.Predicate; import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.Table; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; @@ -148,7 +149,8 @@ public InternalRow[] call(InternalRow args) { boolean netChanges = input.asBoolean(NET_CHANGES, false); if (shouldComputeUpdateImages(input)) { - df = computeUpdateImages(identifierColumns(input, tableIdent), df, netChanges); + Preconditions.checkArgument(!netChanges, "Not support net changes with update images"); + df = computeUpdateImages(identifierColumns(input, tableIdent), df); } else if (shouldRemoveCarryoverRows(input)) { df = removeCarryoverRows(df, netChanges); } @@ -160,23 +162,19 @@ public InternalRow[] call(InternalRow args) { return toOutputRows(viewName); } - private Dataset computeUpdateImages( - String[] identifierColumns, Dataset df, boolean netChanges) { + private Dataset computeUpdateImages(String[] identifierColumns, Dataset df) { Preconditions.checkArgument( identifierColumns.length > 0, "Cannot compute the update images because identifier columns are not set"); - int length = netChanges ? identifierColumns.length : identifierColumns.length + 1; - Column[] repartitionSpec = new Column[length]; + Column[] repartitionSpec = new Column[identifierColumns.length + 1]; for (int i = 0; i < identifierColumns.length; i++) { repartitionSpec[i] = df.col(identifierColumns[i]); } - if (!netChanges) { - repartitionSpec[repartitionSpec.length - 1] = df.col(MetadataColumns.CHANGE_ORDINAL.name()); - } + repartitionSpec[repartitionSpec.length - 1] = df.col(MetadataColumns.CHANGE_ORDINAL.name()); - return applyChangelogIterator(df, repartitionSpec, netChanges); + return applyChangelogIterator(df, repartitionSpec); } private boolean shouldComputeUpdateImages(ProcedureInput input) { @@ -190,12 +188,23 @@ private boolean shouldRemoveCarryoverRows(ProcedureInput input) { } private Dataset removeCarryoverRows(Dataset df, boolean netChanges) { + Predicate columnsToRemove; + if (netChanges) { + columnsToRemove = + column -> + column.equals(MetadataColumns.CHANGE_TYPE.name()) + || column.equals(MetadataColumns.CHANGE_ORDINAL.name()) + || column.equals(MetadataColumns.COMMIT_SNAPSHOT_ID.name()); + } else { + columnsToRemove = column -> column.equals(MetadataColumns.CHANGE_TYPE.name()); + } + Column[] repartitionSpec = Arrays.stream(df.columns()) - .filter(c -> !c.equals(MetadataColumns.CHANGE_TYPE.name())) + .filter(columnsToRemove.negate()) .map(df::col) .toArray(Column[]::new); - return applyCarryoverRemoveIterator(df, repartitionSpec); + return applyCarryoverRemoveIterator(df, repartitionSpec, netChanges); } private String[] identifierColumns(ProcedureInput input, Identifier tableIdent) { @@ -223,9 +232,8 @@ private String viewName(ProcedureInput input, String tableName) { return input.asString(CHANGELOG_VIEW_PARAM, defaultValue); } - private Dataset applyChangelogIterator( - Dataset df, Column[] repartitionSpec, boolean netChanges) { - Column[] sortSpec = sortSpec(df, repartitionSpec); + private Dataset applyChangelogIterator(Dataset df, Column[] repartitionSpec) { + Column[] sortSpec = sortSpec(df, repartitionSpec, false); StructType schema = df.schema(); String[] identifierFields = Arrays.stream(repartitionSpec).map(Column::toString).toArray(String[]::new); @@ -239,26 +247,30 @@ private Dataset applyChangelogIterator( RowEncoder.apply(schema)); } - private Dataset applyCarryoverRemoveIterator(Dataset df, Column[] repartitionSpec) { - // todo double check we should sort by (change order, change type) for net changes - Column[] sortSpec = sortSpec(df, repartitionSpec); + private Dataset applyCarryoverRemoveIterator( + Dataset df, Column[] repartitionSpec, boolean netChanges) { + Column[] sortSpec = sortSpec(df, repartitionSpec, netChanges); StructType schema = df.schema(); return df.repartition(repartitionSpec) .sortWithinPartitions(sortSpec) .mapPartitions( (MapPartitionsFunction) - rowIterator -> ChangelogIterator.removeCarryovers(rowIterator, schema), + rowIterator -> + netChanges + ? ChangelogIterator.removeNetCarryovers(rowIterator, schema) + : ChangelogIterator.removeCarryovers(rowIterator, schema), RowEncoder.apply(schema)); } - private static Column[] sortSpec(Dataset df, Column[] repartitionSpec) { + private static Column[] sortSpec(Dataset df, Column[] repartitionSpec, boolean netChanges) { Column changeOrdinal = df.col(MetadataColumns.CHANGE_ORDINAL.name()); boolean noChangeOrdinal = Arrays.stream(repartitionSpec).noneMatch(c -> c.equals(changeOrdinal)); + Preconditions.checkState(noChangeOrdinal, "Change ordinal should not be in repartition spec"); Column[] sortSpec; - if (noChangeOrdinal) { + if (netChanges) { sortSpec = new Column[repartitionSpec.length + 2]; } else { sortSpec = new Column[repartitionSpec.length + 1]; @@ -266,10 +278,11 @@ private static Column[] sortSpec(Dataset df, Column[] repartitionSpec) { System.arraycopy(repartitionSpec, 0, sortSpec, 0, repartitionSpec.length); - if (noChangeOrdinal) { + sortSpec[sortSpec.length - 1] = df.col(MetadataColumns.CHANGE_TYPE.name()); + + if (netChanges) { sortSpec[sortSpec.length - 2] = changeOrdinal; } - sortSpec[sortSpec.length - 1] = df.col(MetadataColumns.CHANGE_TYPE.name()); return sortSpec; } diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java index 0d28f82399be..e4b7394599b1 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java @@ -301,7 +301,7 @@ private void validateIterators(List rowsWithDuplication, List exp assertEquals("Rows should match.", expectedRows, rowsToJava(result)); - iterator = ChangelogIterator.netChanges(rowsWithDuplication.iterator(), SCHEMA); + iterator = ChangelogIterator.removeNetCarryovers(rowsWithDuplication.iterator(), SCHEMA); result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); @@ -336,7 +336,8 @@ public void testRemoveNetCarryovers() { new Object[] {1, "d", "data", INSERT}, new Object[] {2, "d", "data", DELETE}); - Iterator iterator = ChangelogIterator.netChanges(rowsWithDuplication.iterator(), SCHEMA); + Iterator iterator = + ChangelogIterator.removeNetCarryovers(rowsWithDuplication.iterator(), SCHEMA); List result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); From 79efbaf18645302b2bc86303f1444df2b8c96400 Mon Sep 17 00:00:00 2001 From: yufei Date: Thu, 18 May 2023 18:09:05 -0700 Subject: [PATCH 04/20] Fix the test failures --- .../spark/procedures/CreateChangelogViewProcedure.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index c16a082c44de..d5c43a98d7ec 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -267,7 +267,10 @@ private static Column[] sortSpec(Dataset df, Column[] repartitionSpec, bool Column changeOrdinal = df.col(MetadataColumns.CHANGE_ORDINAL.name()); boolean noChangeOrdinal = Arrays.stream(repartitionSpec).noneMatch(c -> c.equals(changeOrdinal)); - Preconditions.checkState(noChangeOrdinal, "Change ordinal should not be in repartition spec"); + + if (netChanges) { + Preconditions.checkState(noChangeOrdinal, "Change ordinal should not be in repartition spec"); + } Column[] sortSpec; if (netChanges) { From cec8dbff18799c384405a633d3a1167d3a29a6c0 Mon Sep 17 00:00:00 2001 From: yufei Date: Thu, 18 May 2023 21:38:46 -0700 Subject: [PATCH 05/20] Fix the test failures --- .../iceberg/spark/TestChangelogIterator.java | 131 +++++++++++------- 1 file changed, 79 insertions(+), 52 deletions(-) diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java index e4b7394599b1..6b5de5e1a6f0 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java @@ -51,6 +51,26 @@ public class TestChangelogIterator extends SparkTestHelperBase { new StructField( MetadataColumns.CHANGE_TYPE.name(), DataTypes.StringType, false, Metadata.empty()) }); + + private static final StructType SCHEMA_WITH_METADATA_COLUMNS = + new StructType( + new StructField[] { + new StructField("id", DataTypes.IntegerType, false, Metadata.empty()), + new StructField("name", DataTypes.StringType, false, Metadata.empty()), + new StructField("data", DataTypes.StringType, true, Metadata.empty()), + new StructField( + MetadataColumns.CHANGE_TYPE.name(), DataTypes.StringType, false, Metadata.empty()), + new StructField( + MetadataColumns.CHANGE_ORDINAL.name(), + DataTypes.IntegerType, + false, + Metadata.empty()), + new StructField( + MetadataColumns.COMMIT_SNAPSHOT_ID.name(), + DataTypes.LongType, + false, + Metadata.empty()) + }); private static final String[] IDENTIFIER_FIELDS = new String[] {"id", "name"}; private enum RowType { @@ -216,26 +236,26 @@ public void testCarryRowsRemoveWithDuplicates() { List rowsWithDuplication = Lists.newArrayList( // keep all delete rows for id 0 and id 1 since there is no insert row for them - new GenericRowWithSchema(new Object[] {0, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {0, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {0, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "a", "old_data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "a", "old_data", DELETE}, null), + new GenericRowWithSchema(new Object[] {0, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {0, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {0, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "a", "old_data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "a", "old_data", DELETE, 0, 0}, null), // the same number of delete and insert rows for id 2 - new GenericRowWithSchema(new Object[] {2, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {2, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {2, "a", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {2, "a", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {3, "a", "new_data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {2, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {2, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {2, "a", "data", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {2, "a", "data", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {3, "a", "new_data", INSERT, 0, 0}, null)); List expectedRows = Lists.newArrayList( - new Object[] {0, "a", "data", DELETE}, - new Object[] {0, "a", "data", DELETE}, - new Object[] {0, "a", "data", DELETE}, - new Object[] {1, "a", "old_data", DELETE}, - new Object[] {1, "a", "old_data", DELETE}, - new Object[] {3, "a", "new_data", INSERT}); + new Object[] {0, "a", "data", DELETE, 0, 0}, + new Object[] {0, "a", "data", DELETE, 0, 0}, + new Object[] {0, "a", "data", DELETE, 0, 0}, + new Object[] {1, "a", "old_data", DELETE, 0, 0}, + new Object[] {1, "a", "old_data", DELETE, 0, 0}, + new Object[] {3, "a", "new_data", INSERT, 0, 0}); validateIterators(rowsWithDuplication, expectedRows); } @@ -245,14 +265,15 @@ public void testCarryRowsRemoveLessInsertRows() { // less insert rows than delete rows List rowsWithDuplication = Lists.newArrayList( - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {2, "d", "data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {2, "d", "data", INSERT, 0, 0}, null)); List expectedRows = Lists.newArrayList( - new Object[] {1, "d", "data", DELETE}, new Object[] {2, "d", "data", INSERT}); + new Object[] {1, "d", "data", DELETE, 0, 0}, + new Object[] {2, "d", "data", INSERT, 0, 0}); validateIterators(rowsWithDuplication, expectedRows); } @@ -261,19 +282,20 @@ public void testCarryRowsRemoveLessInsertRows() { public void testCarryRowsRemoveMoreInsertRows() { List rowsWithDuplication = Lists.newArrayList( - new GenericRowWithSchema(new Object[] {0, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {0, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), // more insert rows than delete rows, should keep extra insert rows - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 0, 0}, null)); List expectedRows = Lists.newArrayList( - new Object[] {0, "d", "data", DELETE}, new Object[] {1, "d", "data", INSERT}); + new Object[] {0, "d", "data", DELETE, 0, 0}, + new Object[] {1, "d", "data", INSERT, 0, 0}); validateIterators(rowsWithDuplication, expectedRows); } @@ -284,24 +306,28 @@ public void testCarryRowsRemoveNoInsertRows() { List rowsWithDuplication = Lists.newArrayList( // next two rows are identical - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null)); + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null)); List expectedRows = Lists.newArrayList( - new Object[] {1, "d", "data", DELETE}, new Object[] {1, "d", "data", DELETE}); + new Object[] {1, "d", "data", DELETE, 0, 0}, + new Object[] {1, "d", "data", DELETE, 0, 0}); validateIterators(rowsWithDuplication, expectedRows); } private void validateIterators(List rowsWithDuplication, List expectedRows) { Iterator iterator = - ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); + ChangelogIterator.removeCarryovers( + rowsWithDuplication.iterator(), SCHEMA_WITH_METADATA_COLUMNS); List result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); - iterator = ChangelogIterator.removeNetCarryovers(rowsWithDuplication.iterator(), SCHEMA); + iterator = + ChangelogIterator.removeNetCarryovers( + rowsWithDuplication.iterator(), SCHEMA_WITH_METADATA_COLUMNS); result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); @@ -311,33 +337,34 @@ private void validateIterators(List rowsWithDuplication, List exp public void testRemoveNetCarryovers() { List rowsWithDuplication = Lists.newArrayList( - new GenericRowWithSchema(new Object[] {0, "d", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {0, "d", "data", DELETE, 0, 0}, null), // a pair of delete and insert rows - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 0, 0}, null), // 2 delete rows and 2 insert rows - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 1, 1}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 1, 1}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 1, 1}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 1, 1}, null), // a pair of insert and delete rows - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 2, 2}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 2, 2}, null), // extra insert rows - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 2, 2}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 2, 2}, null), // different key - new GenericRowWithSchema(new Object[] {2, "d", "data", DELETE}, null)); + new GenericRowWithSchema(new Object[] {2, "d", "data", DELETE, 3, 3}, null)); List expectedRows = Lists.newArrayList( - new Object[] {0, "d", "data", DELETE}, - new Object[] {1, "d", "data", INSERT}, - new Object[] {1, "d", "data", INSERT}, - new Object[] {2, "d", "data", DELETE}); + new Object[] {0, "d", "data", DELETE, 0, 0}, + new Object[] {1, "d", "data", INSERT, 2, 2}, + new Object[] {1, "d", "data", INSERT, 2, 2}, + new Object[] {2, "d", "data", DELETE, 3, 3}); Iterator iterator = - ChangelogIterator.removeNetCarryovers(rowsWithDuplication.iterator(), SCHEMA); + ChangelogIterator.removeNetCarryovers( + rowsWithDuplication.iterator(), SCHEMA_WITH_METADATA_COLUMNS); List result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); From 10292fcd9938498f358a2aaf5473fb1ad0169f56 Mon Sep 17 00:00:00 2001 From: yufei Date: Fri, 19 May 2023 10:12:12 -0700 Subject: [PATCH 06/20] Add tests --- .../TestCreateChangelogViewProcedure.java | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java index e105df183961..4e7e15fed2e1 100644 --- a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java +++ b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java @@ -433,6 +433,7 @@ public void testNetChangesWithRemoveCarryOvers() { table.refresh(); Snapshot snap3 = table.currentSnapshot(); + // test with all snapshots List returns = sql( "CALL %s.system.create_changelog_view(table => '%s', net_changes => true)", @@ -447,6 +448,20 @@ public void testNetChangesWithRemoveCarryOvers() { row(3, "c", 15, INSERT, 2, snap3.snapshotId()), row(2, "e", 12, INSERT, 2, snap3.snapshotId())), sql("select * from %s order by _change_ordinal, data", viewName)); + + // test with snap2 and snap3 + sql( + "CALL %s.system.create_changelog_view(table => '%s', " + + "options => map('start-snapshot-id','%s'), " + + "net_changes => true)", + catalogName, tableName, snap1.snapshotId()); + + assertEquals( + "Rows should match", + ImmutableList.of( + row(2, "b", 11, DELETE, 0, snap2.snapshotId()), + row(3, "c", 15, INSERT, 1, snap3.snapshotId())), + sql("select * from %s order by _change_ordinal, data", viewName)); } @Test From 516040f5b6d054b6691316d5904719a8b7316b87 Mon Sep 17 00:00:00 2001 From: yufei Date: Wed, 31 May 2023 09:47:15 -0700 Subject: [PATCH 07/20] Use cache row count instead of list to reduce memory usage --- .../spark/RemoveNetCarryoverIterator.java | 28 +++++++++---------- .../iceberg/spark/TestChangelogIterator.java | 14 +++++----- 2 files changed, 20 insertions(+), 22 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index 84f48ee63735..d2bf5cfbd947 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -19,9 +19,7 @@ package org.apache.iceberg.spark; import java.util.Iterator; -import java.util.List; import org.apache.iceberg.MetadataColumns; -import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.StructType; @@ -37,10 +35,11 @@ */ public class RemoveNetCarryoverIterator extends RemoveCarryoverIterator { - private final List cachedRows = Lists.newArrayList(); private final int[] indicesToIdentifySameRow; private Row cachedNextRow = null; + private Row cachedRow = null; + private long cachedRowCount = 0; protected RemoveNetCarryoverIterator(Iterator rowIterator, StructType rowType) { super(rowIterator, rowType); @@ -49,7 +48,7 @@ protected RemoveNetCarryoverIterator(Iterator rowIterator, StructType rowTy @Override public boolean hasNext() { - if (!cachedRows.isEmpty()) { + if (cachedRowCount > 0) { return true; } @@ -63,8 +62,9 @@ public boolean hasNext() { @Override public Row next() { // if there are cached rows, return one of them from the beginning - if (!cachedRows.isEmpty()) { - return cachedRows.remove(0); + if (cachedRowCount > 0) { + cachedRowCount--; + return cachedRow; } Row currentRow = getCurrentRow(); @@ -76,27 +76,25 @@ public Row next() { Row nextRow = rowIterator().next(); - cachedRows.add(currentRow); + cachedRow = currentRow; + cachedRowCount = 1; // if they are the same row, remove the cached row or stack to the cached row while (isSameRecord(currentRow, nextRow)) { if (matched(currentRow, nextRow)) { - // if matched, remove both rows, remove the last row in the cache - cachedRows.remove(cachedRows.size() - 1); + // remove both rows + cachedRowCount--; nextRow = null; } else { - // stack the row into the cache - cachedRows.add(nextRow); + // stack the next row to the cached row nextRow = null; + cachedRowCount++; } // break the loop if there is no next row or the cache is empty - if (cachedRows.isEmpty() || !rowIterator().hasNext()) { + if (cachedRowCount <= 0 || !rowIterator().hasNext()) { break; } - // get the current row from the cache, and get the next row from the iterator for the next - // loop - currentRow = cachedRows.get(cachedRows.size() - 1); nextRow = rowIterator().next(); } diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java index 6b5de5e1a6f0..88ed981bd779 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java @@ -348,19 +348,19 @@ public void testRemoveNetCarryovers() { new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 1, 1}, null), // a pair of insert and delete rows new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 2, 2}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 2, 2}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 3, 3}, null), // extra insert rows - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 2, 2}, null), - new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 2, 2}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 4, 4}, null), + new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 4, 4}, null), // different key - new GenericRowWithSchema(new Object[] {2, "d", "data", DELETE, 3, 3}, null)); + new GenericRowWithSchema(new Object[] {2, "d", "data", DELETE, 4, 4}, null)); List expectedRows = Lists.newArrayList( new Object[] {0, "d", "data", DELETE, 0, 0}, - new Object[] {1, "d", "data", INSERT, 2, 2}, - new Object[] {1, "d", "data", INSERT, 2, 2}, - new Object[] {2, "d", "data", DELETE, 3, 3}); + new Object[] {1, "d", "data", INSERT, 4, 4}, + new Object[] {1, "d", "data", INSERT, 4, 4}, + new Object[] {2, "d", "data", DELETE, 4, 4}); Iterator iterator = ChangelogIterator.removeNetCarryovers( From 058eb7d16b10441f0893a08123a3305a10606eb4 Mon Sep 17 00:00:00 2001 From: yufei Date: Fri, 2 Jun 2023 13:36:31 -0700 Subject: [PATCH 08/20] Resolve comments --- .../spark/RemoveNetCarryoverIterator.java | 15 +++++++------- .../CreateChangelogViewProcedure.java | 20 ++++++++----------- 2 files changed, 16 insertions(+), 19 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index d2bf5cfbd947..626c313107ad 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -69,7 +69,7 @@ public Row next() { Row currentRow = getCurrentRow(); - // if there is no next row, return the current row + // return it directly if the current row is the last row if (!rowIterator().hasNext()) { return currentRow; } @@ -78,19 +78,20 @@ public Row next() { cachedRow = currentRow; cachedRowCount = 1; - // if they are the same row, remove the cached row or stack to the cached row + + // pull rows from the iterator until two consecutive rows are different while (isSameRecord(currentRow, nextRow)) { - if (matched(currentRow, nextRow)) { - // remove both rows + if (oppositeChangeType(currentRow, nextRow)) { + // two rows with opposite change types means no net changes cachedRowCount--; nextRow = null; } else { - // stack the next row to the cached row + // two rows with same change types means potential net changes nextRow = null; cachedRowCount++; } - // break the loop if there is no next row or the cache is empty + // stop pulling rows if there is no more rows or the next row is different if (cachedRowCount <= 0 || !rowIterator().hasNext()) { break; } @@ -114,7 +115,7 @@ private Row getCurrentRow() { return currentRow; } - private boolean matched(Row currentRow, Row nextRow) { + private boolean oppositeChangeType(Row currentRow, Row nextRow) { return (nextRow.getString(changeTypeIndex()).equals(INSERT) && currentRow.getString(changeTypeIndex()).equals(DELETE)) || (nextRow.getString(changeTypeIndex()).equals(DELETE) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index d5c43a98d7ec..493bf54ef746 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -264,28 +264,24 @@ private Dataset applyCarryoverRemoveIterator( } private static Column[] sortSpec(Dataset df, Column[] repartitionSpec, boolean netChanges) { - Column changeOrdinal = df.col(MetadataColumns.CHANGE_ORDINAL.name()); - boolean noChangeOrdinal = - Arrays.stream(repartitionSpec).noneMatch(c -> c.equals(changeOrdinal)); + Column[] sortSpec; if (netChanges) { - Preconditions.checkState(noChangeOrdinal, "Change ordinal should not be in repartition spec"); - } + Column changeOrdinal = df.col(MetadataColumns.CHANGE_ORDINAL.name()); + Preconditions.checkState( + Arrays.stream(repartitionSpec).noneMatch(c -> c.equals(changeOrdinal)), + "Change ordinal should not be in repartition spec"); - Column[] sortSpec; - if (netChanges) { sortSpec = new Column[repartitionSpec.length + 2]; + sortSpec[sortSpec.length - 2] = changeOrdinal; } else { sortSpec = new Column[repartitionSpec.length + 1]; } - System.arraycopy(repartitionSpec, 0, sortSpec, 0, repartitionSpec.length); - sortSpec[sortSpec.length - 1] = df.col(MetadataColumns.CHANGE_TYPE.name()); - if (netChanges) { - sortSpec[sortSpec.length - 2] = changeOrdinal; - } + System.arraycopy(repartitionSpec, 0, sortSpec, 0, repartitionSpec.length); + return sortSpec; } From 1bcff458f6ae7372c6ccd2eab4c22e8c7a3c9a60 Mon Sep 17 00:00:00 2001 From: yufei Date: Mon, 5 Jun 2023 16:00:59 -0700 Subject: [PATCH 09/20] Resolve comments --- .../iceberg/spark/TestChangelogIterator.java | 99 ++++++++----------- 1 file changed, 43 insertions(+), 56 deletions(-) diff --git a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java index 88ed981bd779..0539598f147e 100644 --- a/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java +++ b/spark/v3.3/spark/src/test/java/org/apache/iceberg/spark/TestChangelogIterator.java @@ -43,16 +43,6 @@ public class TestChangelogIterator extends SparkTestHelperBase { private static final String UPDATE_AFTER = ChangelogOperation.UPDATE_AFTER.name(); private static final StructType SCHEMA = - new StructType( - new StructField[] { - new StructField("id", DataTypes.IntegerType, false, Metadata.empty()), - new StructField("name", DataTypes.StringType, false, Metadata.empty()), - new StructField("data", DataTypes.StringType, true, Metadata.empty()), - new StructField( - MetadataColumns.CHANGE_TYPE.name(), DataTypes.StringType, false, Metadata.empty()) - }); - - private static final StructType SCHEMA_WITH_METADATA_COLUMNS = new StructType( new StructField[] { new StructField("id", DataTypes.IntegerType, false, Metadata.empty()), @@ -113,18 +103,18 @@ private List toOriginalRows(RowType rowType, int index) { switch (rowType) { case DELETED: return Lists.newArrayList( - new GenericRowWithSchema(new Object[] {index, "b", "data", DELETE}, null)); + new GenericRowWithSchema(new Object[] {index, "b", "data", DELETE, 0, 0}, null)); case INSERTED: return Lists.newArrayList( - new GenericRowWithSchema(new Object[] {index, "c", "data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {index, "c", "data", INSERT, 0, 0}, null)); case CARRY_OVER: return Lists.newArrayList( - new GenericRowWithSchema(new Object[] {index, "d", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {index, "d", "data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {index, "d", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {index, "d", "data", INSERT, 0, 0}, null)); case UPDATED: return Lists.newArrayList( - new GenericRowWithSchema(new Object[] {index, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {index, "a", "new_data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {index, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {index, "a", "new_data", INSERT, 0, 0}, null)); default: throw new IllegalArgumentException("Unknown row type: " + rowType); } @@ -134,18 +124,18 @@ private List toExpectedRows(RowType rowType, int order) { switch (rowType) { case DELETED: List rows = Lists.newArrayList(); - rows.add(new Object[] {order, "b", "data", DELETE}); + rows.add(new Object[] {order, "b", "data", DELETE, 0, 0}); return rows; case INSERTED: List insertedRows = Lists.newArrayList(); - insertedRows.add(new Object[] {order, "c", "data", INSERT}); + insertedRows.add(new Object[] {order, "c", "data", INSERT, 0, 0}); return insertedRows; case CARRY_OVER: return Lists.newArrayList(); case UPDATED: return Lists.newArrayList( - new Object[] {order, "a", "data", UPDATE_BEFORE}, - new Object[] {order, "a", "new_data", UPDATE_AFTER}); + new Object[] {order, "a", "data", UPDATE_BEFORE, 0, 0}, + new Object[] {order, "a", "new_data", UPDATE_AFTER, 0, 0}); default: throw new IllegalArgumentException("Unknown row type: " + rowType); } @@ -166,16 +156,16 @@ private void permute(List arr, int start, List pm) { public void testRowsWithNullValue() { final List rowsWithNull = Lists.newArrayList( - new GenericRowWithSchema(new Object[] {2, null, null, DELETE}, null), - new GenericRowWithSchema(new Object[] {3, null, null, INSERT}, null), - new GenericRowWithSchema(new Object[] {4, null, null, DELETE}, null), - new GenericRowWithSchema(new Object[] {4, null, null, INSERT}, null), + new GenericRowWithSchema(new Object[] {2, null, null, DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {3, null, null, INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {4, null, null, DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {4, null, null, INSERT, 0, 0}, null), // mixed null and non-null value in non-identifier columns - new GenericRowWithSchema(new Object[] {5, null, null, DELETE}, null), - new GenericRowWithSchema(new Object[] {5, null, "data", INSERT}, null), + new GenericRowWithSchema(new Object[] {5, null, null, DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {5, null, "data", INSERT, 0, 0}, null), // mixed null and non-null value in identifier columns - new GenericRowWithSchema(new Object[] {6, null, null, DELETE}, null), - new GenericRowWithSchema(new Object[] {6, "name", null, INSERT}, null)); + new GenericRowWithSchema(new Object[] {6, null, null, DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {6, "name", null, INSERT, 0, 0}, null)); Iterator iterator = ChangelogIterator.computeUpdates(rowsWithNull.iterator(), SCHEMA, IDENTIFIER_FIELDS); @@ -184,12 +174,12 @@ public void testRowsWithNullValue() { assertEquals( "Rows should match", Lists.newArrayList( - new Object[] {2, null, null, DELETE}, - new Object[] {3, null, null, INSERT}, - new Object[] {5, null, null, UPDATE_BEFORE}, - new Object[] {5, null, "data", UPDATE_AFTER}, - new Object[] {6, null, null, DELETE}, - new Object[] {6, "name", null, INSERT}), + new Object[] {2, null, null, DELETE, 0, 0}, + new Object[] {3, null, null, INSERT, 0, 0}, + new Object[] {5, null, null, UPDATE_BEFORE, 0, 0}, + new Object[] {5, null, "data", UPDATE_AFTER, 0, 0}, + new Object[] {6, null, null, DELETE, 0, 0}, + new Object[] {6, "name", null, INSERT, 0, 0}), rowsToJava(result)); } @@ -198,10 +188,10 @@ public void testUpdatedRowsWithDuplication() { List rowsWithDuplication = Lists.newArrayList( // two rows with same identifier fields(id, name) - new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT}, null)); + new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "a", "new_data", INSERT, 0, 0}, null)); Iterator iterator = ChangelogIterator.computeUpdates(rowsWithDuplication.iterator(), SCHEMA, IDENTIFIER_FIELDS); @@ -214,9 +204,9 @@ public void testUpdatedRowsWithDuplication() { // still allow extra insert rows rowsWithDuplication = Lists.newArrayList( - new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE}, null), - new GenericRowWithSchema(new Object[] {1, "a", "new_data1", INSERT}, null), - new GenericRowWithSchema(new Object[] {1, "a", "new_data2", INSERT}, null)); + new GenericRowWithSchema(new Object[] {1, "a", "data", DELETE, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "a", "new_data1", INSERT, 0, 0}, null), + new GenericRowWithSchema(new Object[] {1, "a", "new_data2", INSERT, 0, 0}, null)); Iterator iterator1 = ChangelogIterator.computeUpdates(rowsWithDuplication.iterator(), SCHEMA, IDENTIFIER_FIELDS); @@ -224,9 +214,9 @@ public void testUpdatedRowsWithDuplication() { assertEquals( "Rows should match.", Lists.newArrayList( - new Object[] {1, "a", "data", UPDATE_BEFORE}, - new Object[] {1, "a", "new_data1", UPDATE_AFTER}, - new Object[] {1, "a", "new_data2", INSERT}), + new Object[] {1, "a", "data", UPDATE_BEFORE, 0, 0}, + new Object[] {1, "a", "new_data1", UPDATE_AFTER, 0, 0}, + new Object[] {1, "a", "new_data2", INSERT, 0, 0}), rowsToJava(Lists.newArrayList(iterator1))); } @@ -319,15 +309,12 @@ public void testCarryRowsRemoveNoInsertRows() { private void validateIterators(List rowsWithDuplication, List expectedRows) { Iterator iterator = - ChangelogIterator.removeCarryovers( - rowsWithDuplication.iterator(), SCHEMA_WITH_METADATA_COLUMNS); + ChangelogIterator.removeCarryovers(rowsWithDuplication.iterator(), SCHEMA); List result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); - iterator = - ChangelogIterator.removeNetCarryovers( - rowsWithDuplication.iterator(), SCHEMA_WITH_METADATA_COLUMNS); + iterator = ChangelogIterator.removeNetCarryovers(rowsWithDuplication.iterator(), SCHEMA); result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); @@ -337,22 +324,23 @@ private void validateIterators(List rowsWithDuplication, List exp public void testRemoveNetCarryovers() { List rowsWithDuplication = Lists.newArrayList( + // this row are different from other rows, it is a net change, should be kept new GenericRowWithSchema(new Object[] {0, "d", "data", DELETE, 0, 0}, null), - // a pair of delete and insert rows + // a pair of delete and insert rows, should be removed new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 0, 0}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 0, 0}, null), - // 2 delete rows and 2 insert rows + // 2 delete rows and 2 insert rows, should be removed new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 1, 1}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 1, 1}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 1, 1}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 1, 1}, null), - // a pair of insert and delete rows + // a pair of insert and delete rows across snapshots, should be removed new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 2, 2}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", DELETE, 3, 3}, null), - // extra insert rows + // extra insert rows, they are net changes, should be kept new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 4, 4}, null), new GenericRowWithSchema(new Object[] {1, "d", "data", INSERT, 4, 4}, null), - // different key + // different key, net changes, should be kept new GenericRowWithSchema(new Object[] {2, "d", "data", DELETE, 4, 4}, null)); List expectedRows = @@ -363,8 +351,7 @@ public void testRemoveNetCarryovers() { new Object[] {2, "d", "data", DELETE, 4, 4}); Iterator iterator = - ChangelogIterator.removeNetCarryovers( - rowsWithDuplication.iterator(), SCHEMA_WITH_METADATA_COLUMNS); + ChangelogIterator.removeNetCarryovers(rowsWithDuplication.iterator(), SCHEMA); List result = Lists.newArrayList(iterator); assertEquals("Rows should match.", expectedRows, rowsToJava(result)); From 95fa7a6cc956e1e1890257630e48d608cf2499b2 Mon Sep 17 00:00:00 2001 From: yufei Date: Thu, 22 Jun 2023 13:48:28 -0700 Subject: [PATCH 10/20] Resolve comments --- gradle.properties | 2 +- .../TestCreateChangelogViewProcedure.java | 34 ++++++++++--------- .../iceberg/spark/ChangelogIterator.java | 4 +++ .../iceberg/spark/ComputeUpdateIterator.java | 8 ++--- .../spark/RemoveCarryoverIterator.java | 4 +-- .../spark/RemoveNetCarryoverIterator.java | 18 +++++----- .../CreateChangelogViewProcedure.java | 20 ++++------- 7 files changed, 43 insertions(+), 47 deletions(-) diff --git a/gradle.properties b/gradle.properties index eb0da0ac8547..ffb2454dcdde 100644 --- a/gradle.properties +++ b/gradle.properties @@ -20,7 +20,7 @@ systemProp.defaultFlinkVersions=1.17 systemProp.knownFlinkVersions=1.15,1.16,1.17 systemProp.defaultHiveVersions=2 systemProp.knownHiveVersions=2,3 -systemProp.defaultSparkVersions=3.4 +systemProp.defaultSparkVersions=3.4,3.3 systemProp.knownSparkVersions=3.1,3.2,3.3,3.4 systemProp.defaultScalaVersion=2.12 systemProp.knownScalaVersions=2.12,2.13 diff --git a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java index 4e7e15fed2e1..0a4a1073c3f9 100644 --- a/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java +++ b/spark/v3.3/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestCreateChangelogViewProcedure.java @@ -47,13 +47,13 @@ public void removeTable() { sql("DROP TABLE IF EXISTS %s", tableName); } - public void createTableWith2Columns() { + public void createTableWithTwoColumns() { sql("CREATE TABLE %s (id INT, data STRING) USING iceberg", tableName); sql("ALTER TABLE %s SET TBLPROPERTIES ('format-version'='%d')", tableName, 1); sql("ALTER TABLE %s ADD PARTITION FIELD data", tableName); } - private void createTableWith3Columns() { + private void createTableWithThreeColumns() { sql("CREATE TABLE %s (id INT, data STRING, age INT) USING iceberg", tableName); sql("ALTER TABLE %s SET TBLPROPERTIES ('format-version'='%d')", tableName, 1); sql("ALTER TABLE %s ADD PARTITION FIELD id", tableName); @@ -67,7 +67,7 @@ private void createTableWithIdentifierField() { @Test public void testCustomizedViewName() { - createTableWith2Columns(); + createTableWithTwoColumns(); sql("INSERT INTO %s VALUES (1, 'a')", tableName); sql("INSERT INTO %s VALUES (2, 'b')", tableName); @@ -100,7 +100,7 @@ public void testCustomizedViewName() { @Test public void testNoSnapshotIdInput() { - createTableWith2Columns(); + createTableWithTwoColumns(); sql("INSERT INTO %s VALUES (1, 'a')", tableName); Table table = validationCatalog.loadTable(tableIdent); Snapshot snap0 = table.currentSnapshot(); @@ -131,7 +131,7 @@ public void testNoSnapshotIdInput() { @Test public void testTimestampsBasedQuery() { - createTableWith2Columns(); + createTableWithTwoColumns(); long beginning = System.currentTimeMillis(); sql("INSERT INTO %s VALUES (1, 'a')", tableName); @@ -191,7 +191,7 @@ public void testTimestampsBasedQuery() { @Test public void testWithCarryovers() { - createTableWith2Columns(); + createTableWithTwoColumns(); sql("INSERT INTO %s VALUES (1, 'a')", tableName); Table table = validationCatalog.loadTable(tableIdent); Snapshot snap0 = table.currentSnapshot(); @@ -226,7 +226,7 @@ public void testWithCarryovers() { @Test public void testUpdate() { - createTableWith2Columns(); + createTableWithTwoColumns(); sql("ALTER TABLE %s DROP PARTITION FIELD data", tableName); sql("ALTER TABLE %s ADD PARTITION FIELD id", tableName); @@ -285,7 +285,7 @@ public void testUpdateWithIdentifierField() { @Test public void testUpdateWithFilter() { - createTableWith2Columns(); + createTableWithTwoColumns(); sql("ALTER TABLE %s DROP PARTITION FIELD data", tableName); sql("ALTER TABLE %s ADD PARTITION FIELD id", tableName); @@ -317,7 +317,7 @@ public void testUpdateWithFilter() { @Test public void testUpdateWithMultipleIdentifierColumns() { - createTableWith3Columns(); + createTableWithThreeColumns(); sql("INSERT INTO %s VALUES (1, 'a', 12), (2, 'b', 11)", tableName); Table table = validationCatalog.loadTable(tableIdent); @@ -349,7 +349,7 @@ public void testUpdateWithMultipleIdentifierColumns() { @Test public void testRemoveCarryOvers() { - createTableWith3Columns(); + createTableWithThreeColumns(); sql("INSERT INTO %s VALUES (1, 'a', 12), (2, 'b', 11), (2, 'e', 12)", tableName); Table table = validationCatalog.loadTable(tableIdent); @@ -383,7 +383,7 @@ public void testRemoveCarryOvers() { @Test public void testRemoveCarryOversWithoutUpdatedRows() { - createTableWith3Columns(); + createTableWithThreeColumns(); sql("INSERT INTO %s VALUES (1, 'a', 12), (2, 'b', 11), (2, 'e', 12)", tableName); Table table = validationCatalog.loadTable(tableIdent); @@ -416,19 +416,21 @@ public void testRemoveCarryOversWithoutUpdatedRows() { @Test public void testNetChangesWithRemoveCarryOvers() { // partitioned by id - createTableWith3Columns(); + createTableWithThreeColumns(); // insert rows: (1, 'a', 12) (2, 'b', 11) (2, 'e', 12) sql("INSERT INTO %s VALUES (1, 'a', 12), (2, 'b', 11), (2, 'e', 12)", tableName); Table table = validationCatalog.loadTable(tableIdent); Snapshot snap1 = table.currentSnapshot(); - // delete rows: (2, 'b', 11) (2, 'e', 12), insert rows: (3, 'c', 13) (2, 'd', 11) (2, 'e', 12) + // delete rows: (2, 'b', 11) (2, 'e', 12) + // insert rows: (3, 'c', 13) (2, 'd', 11) (2, 'e', 12) sql("INSERT OVERWRITE %s VALUES (3, 'c', 13), (2, 'd', 11), (2, 'e', 12)", tableName); table.refresh(); Snapshot snap2 = table.currentSnapshot(); - // delete rows: (2, 'd', 11) (2, 'e', 12) (3, 'c', 13), insert rows: (3, 'c', 15) (2, 'e', 12) + // delete rows: (2, 'd', 11) (2, 'e', 12) (3, 'c', 13) + // insert rows: (3, 'c', 15) (2, 'e', 12) sql("INSERT OVERWRITE %s VALUES (3, 'c', 15), (2, 'e', 12)", tableName); table.refresh(); Snapshot snap3 = table.currentSnapshot(); @@ -466,7 +468,7 @@ public void testNetChangesWithRemoveCarryOvers() { @Test public void testNetChangesWithComputeUpdates() { - createTableWith2Columns(); + createTableWithTwoColumns(); assertThrows( "Should fail because net_changes is not supported with computing updates", IllegalArgumentException.class, @@ -478,7 +480,7 @@ public void testNetChangesWithComputeUpdates() { @Test public void testNotRemoveCarryOvers() { - createTableWith3Columns(); + createTableWithThreeColumns(); sql("INSERT INTO %s VALUES (1, 'a', 12), (2, 'b', 11), (2, 'e', 12)", tableName); Table table = validationCatalog.loadTable(tableIdent); diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java index e848e61d6204..e3f3f14a30a9 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java @@ -45,6 +45,10 @@ protected int changeTypeIndex() { return changeTypeIndex; } + protected String changeType(Row row) { + return row.getString(changeTypeIndex()); + } + protected Iterator rowIterator() { return rowIterator; } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ComputeUpdateIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ComputeUpdateIterator.java index 23e6a19a17e7..6951c33e51aa 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ComputeUpdateIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ComputeUpdateIterator.java @@ -81,15 +81,13 @@ public Row next() { // either a cached record which is not an UPDATE or the next record in the iterator. Row currentRow = currentRow(); - if (currentRow.getString(changeTypeIndex()).equals(DELETE) && rowIterator().hasNext()) { + if (changeType(currentRow).equals(DELETE) && rowIterator().hasNext()) { Row nextRow = rowIterator().next(); cachedRow = nextRow; if (sameLogicalRow(currentRow, nextRow)) { - String nextRowChangeType = nextRow.getString(changeTypeIndex()); - Preconditions.checkState( - nextRowChangeType.equals(INSERT), + changeType(nextRow).equals(INSERT), "Cannot compute updates because there are multiple rows with the same identifier" + " fields([%s]). Please make sure the rows are unique.", String.join(",", identifierFields)); @@ -118,7 +116,7 @@ private Row modify(Row row, int valueIndex, Object value) { } private boolean cachedUpdateRecord() { - return cachedRow != null && cachedRow.getString(changeTypeIndex()).equals(UPDATE_AFTER); + return cachedRow != null && changeType(cachedRow).equals(UPDATE_AFTER); } private Row currentRow() { diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java index b1431c2b879c..bd2bb291db5b 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java @@ -90,7 +90,7 @@ public Row next() { } // If the current row is a delete row, drain all identical delete rows - if (currentRow.getString(changeTypeIndex()).equals(DELETE) && rowIterator().hasNext()) { + if (changeType(currentRow).equals(DELETE) && rowIterator().hasNext()) { cachedDeletedRow = currentRow; deletedRowCount = 1; @@ -101,7 +101,7 @@ public Row next() { while (nextRow != null && cachedDeletedRow != null && isSameRecord(cachedDeletedRow, nextRow)) { - if (nextRow.getString(changeTypeIndex()).equals(INSERT)) { + if (changeType(nextRow).equals(INSERT)) { deletedRowCount--; if (deletedRowCount == 0) { cachedDeletedRow = null; diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index 626c313107ad..d2f47a117264 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -24,8 +24,9 @@ import org.apache.spark.sql.types.StructType; /** - * This class computes the net changes across multiple snapshots. It takes a row iterator, and - * assumes the following: + * This class computes the net changes across multiple snapshots. It is different from {@link + * RemoveCarryoverIterator}, which only removes carry-over rows within a single snapshot. It takes a + * row iterator, and assumes the following: * *
    *
  • The row iterator is partitioned by all columns. @@ -82,13 +83,14 @@ public Row next() { // pull rows from the iterator until two consecutive rows are different while (isSameRecord(currentRow, nextRow)) { if (oppositeChangeType(currentRow, nextRow)) { - // two rows with opposite change types means no net changes + // two rows with opposite change types means no net changes, remove both cachedRowCount--; nextRow = null; } else { - // two rows with same change types means potential net changes - nextRow = null; + // two rows with same change types means potential net changes, cache the next row, reset it + // to null cachedRowCount++; + nextRow = null; } // stop pulling rows if there is no more rows or the next row is different @@ -116,10 +118,8 @@ private Row getCurrentRow() { } private boolean oppositeChangeType(Row currentRow, Row nextRow) { - return (nextRow.getString(changeTypeIndex()).equals(INSERT) - && currentRow.getString(changeTypeIndex()).equals(DELETE)) - || (nextRow.getString(changeTypeIndex()).equals(DELETE) - && currentRow.getString(changeTypeIndex()).equals(INSERT)); + return (changeType(nextRow).equals(INSERT) && changeType(currentRow).equals(DELETE)) + || (changeType(nextRow).equals(DELETE) && changeType(currentRow).equals(INSERT)); } private int[] generateIndicesToIdentifySameRow() { diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index 493bf54ef746..12b13a5a8a83 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -264,23 +264,15 @@ private Dataset applyCarryoverRemoveIterator( } private static Column[] sortSpec(Dataset df, Column[] repartitionSpec, boolean netChanges) { - Column[] sortSpec; + Column changeType = df.col(MetadataColumns.CHANGE_TYPE.name()); + Column changeOrdinal = df.col(MetadataColumns.CHANGE_ORDINAL.name()); + Column[] extraColumns = + netChanges ? new Column[] {changeOrdinal, changeType} : new Column[] {changeType}; - if (netChanges) { - Column changeOrdinal = df.col(MetadataColumns.CHANGE_ORDINAL.name()); - Preconditions.checkState( - Arrays.stream(repartitionSpec).noneMatch(c -> c.equals(changeOrdinal)), - "Change ordinal should not be in repartition spec"); - - sortSpec = new Column[repartitionSpec.length + 2]; - sortSpec[sortSpec.length - 2] = changeOrdinal; - } else { - sortSpec = new Column[repartitionSpec.length + 1]; - } - - sortSpec[sortSpec.length - 1] = df.col(MetadataColumns.CHANGE_TYPE.name()); + Column[] sortSpec = new Column[repartitionSpec.length + extraColumns.length]; System.arraycopy(repartitionSpec, 0, sortSpec, 0, repartitionSpec.length); + System.arraycopy(extraColumns, 0, sortSpec, repartitionSpec.length, extraColumns.length); return sortSpec; } From 2b3c5de999ca338d1f083a07a524b791e0f5b673 Mon Sep 17 00:00:00 2001 From: yufei Date: Thu, 22 Jun 2023 14:36:05 -0700 Subject: [PATCH 11/20] Resolve comments --- .../CreateChangelogViewProcedure.java | 23 ++++++++++--------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index 12b13a5a8a83..c0951367fd8f 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -21,7 +21,9 @@ import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.function.Predicate; +import org.apache.commons.compress.utils.Sets; import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.Table; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; @@ -188,22 +190,21 @@ private boolean shouldRemoveCarryoverRows(ProcedureInput input) { } private Dataset removeCarryoverRows(Dataset df, boolean netChanges) { - Predicate columnsToRemove; + Predicate columnsToKeep; if (netChanges) { - columnsToRemove = - column -> - column.equals(MetadataColumns.CHANGE_TYPE.name()) - || column.equals(MetadataColumns.CHANGE_ORDINAL.name()) - || column.equals(MetadataColumns.COMMIT_SNAPSHOT_ID.name()); + Set metadataColumn = + Sets.newHashSet( + MetadataColumns.CHANGE_TYPE.name(), + MetadataColumns.CHANGE_ORDINAL.name(), + MetadataColumns.COMMIT_SNAPSHOT_ID.name()); + + columnsToKeep = column -> !metadataColumn.contains(column); } else { - columnsToRemove = column -> column.equals(MetadataColumns.CHANGE_TYPE.name()); + columnsToKeep = column -> !column.equals(MetadataColumns.CHANGE_TYPE.name()); } Column[] repartitionSpec = - Arrays.stream(df.columns()) - .filter(columnsToRemove.negate()) - .map(df::col) - .toArray(Column[]::new); + Arrays.stream(df.columns()).filter(columnsToKeep).map(df::col).toArray(Column[]::new); return applyCarryoverRemoveIterator(df, repartitionSpec, netChanges); } From 5bab4e100aa42cff4c7652046a3ff767aaab1bfe Mon Sep 17 00:00:00 2001 From: yufei Date: Thu, 22 Jun 2023 14:55:43 -0700 Subject: [PATCH 12/20] Use the Sets in guava --- .../iceberg/spark/procedures/CreateChangelogViewProcedure.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index c0951367fd8f..b2cecb8cf15f 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -23,12 +23,12 @@ import java.util.Map; import java.util.Set; import java.util.function.Predicate; -import org.apache.commons.compress.utils.Sets; import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.Table; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; import org.apache.iceberg.spark.ChangelogIterator; import org.apache.iceberg.spark.source.SparkChangelogTable; import org.apache.spark.api.java.function.MapPartitionsFunction; From 414285c3e5151f463ae1c2ca8ea8b69a18ddb835 Mon Sep 17 00:00:00 2001 From: yufei Date: Tue, 27 Jun 2023 11:01:46 -0700 Subject: [PATCH 13/20] Resolve comments --- .../iceberg/spark/ChangelogIterator.java | 10 ++++++++++ .../spark/RemoveCarryoverIterator.java | 19 +------------------ .../spark/RemoveNetCarryoverIterator.java | 16 +++++++--------- 3 files changed, 18 insertions(+), 27 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java index e3f3f14a30a9..9e6cec44c4a5 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java @@ -88,6 +88,16 @@ public static Iterator removeNetCarryovers(Iterator rowIterator, Struc return Iterators.filter(changelogIterator, Objects::nonNull); } + protected boolean isSameRecord(Row currentRow, Row nextRow, int[] indicesToIdentifySameRow) { + for (int idx : indicesToIdentifySameRow) { + if (isDifferentValue(currentRow, nextRow, idx)) { + return false; + } + } + + return true; + } + protected boolean isDifferentValue(Row currentRow, Row nextRow, int idx) { return !Objects.equals(nextRow.get(idx), currentRow.get(idx)); } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java index bd2bb291db5b..93272a2af8eb 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java @@ -100,7 +100,7 @@ public Row next() { // row is the same record while (nextRow != null && cachedDeletedRow != null - && isSameRecord(cachedDeletedRow, nextRow)) { + && isSameRecord(cachedDeletedRow, nextRow, indicesToIdentifySameRow)) { if (changeType(nextRow).equals(INSERT)) { deletedRowCount--; if (deletedRowCount == 0) { @@ -153,21 +153,4 @@ private int[] generateIndicesToIdentifySameRow() { return indices; } - protected boolean isSameRecord(Row currentRow, Row nextRow) { - for (int idx : indicesToIdentifySameRow()) { - if (isDifferentValue(currentRow, nextRow, idx)) { - return false; - } - } - - return true; - } - - protected StructType rowType() { - return rowType; - } - - protected int[] indicesToIdentifySameRow() { - return indicesToIdentifySameRow; - } } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index d2f47a117264..5ba7697c5f28 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -34,9 +34,10 @@ * is 1-to-1 mapping to snapshot id. *
*/ -public class RemoveNetCarryoverIterator extends RemoveCarryoverIterator { +public class RemoveNetCarryoverIterator extends ChangelogIterator { private final int[] indicesToIdentifySameRow; + private final StructType rowType; private Row cachedNextRow = null; private Row cachedRow = null; @@ -44,6 +45,7 @@ public class RemoveNetCarryoverIterator extends RemoveCarryoverIterator { protected RemoveNetCarryoverIterator(Iterator rowIterator, StructType rowType) { super(rowIterator, rowType); + this.rowType = rowType; this.indicesToIdentifySameRow = generateIndicesToIdentifySameRow(); } @@ -81,7 +83,7 @@ public Row next() { cachedRowCount = 1; // pull rows from the iterator until two consecutive rows are different - while (isSameRecord(currentRow, nextRow)) { + while (isSameRecord(currentRow, nextRow, indicesToIdentifySameRow)) { if (oppositeChangeType(currentRow, nextRow)) { // two rows with opposite change types means no net changes, remove both cachedRowCount--; @@ -123,10 +125,10 @@ private boolean oppositeChangeType(Row currentRow, Row nextRow) { } private int[] generateIndicesToIdentifySameRow() { - int changeOrdinalIndex = rowType().fieldIndex(MetadataColumns.CHANGE_ORDINAL.name()); - int snapshotIdIndex = rowType().fieldIndex(MetadataColumns.COMMIT_SNAPSHOT_ID.name()); + int changeOrdinalIndex = rowType.fieldIndex(MetadataColumns.CHANGE_ORDINAL.name()); + int snapshotIdIndex = rowType.fieldIndex(MetadataColumns.COMMIT_SNAPSHOT_ID.name()); - int[] indices = new int[rowType().size() - 3]; + int[] indices = new int[rowType.size() - 3]; for (int i = 0, j = 0; i < indices.length; i++) { if (i != changeTypeIndex() && i != changeOrdinalIndex && i != snapshotIdIndex) { @@ -137,8 +139,4 @@ private int[] generateIndicesToIdentifySameRow() { return indices; } - @Override - protected int[] indicesToIdentifySameRow() { - return indicesToIdentifySameRow; - } } From 72617a6da3c90d5d4c34e393a548f139d7d4c4fb Mon Sep 17 00:00:00 2001 From: yufei Date: Tue, 27 Jun 2023 11:26:03 -0700 Subject: [PATCH 14/20] Resolve comments --- .../spark/RemoveCarryoverIterator.java | 1 - .../spark/RemoveNetCarryoverIterator.java | 26 +++++++------------ 2 files changed, 10 insertions(+), 17 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java index 93272a2af8eb..f562aa2ee54f 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java @@ -152,5 +152,4 @@ private int[] generateIndicesToIdentifySameRow() { } return indices; } - } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index 5ba7697c5f28..f72301078c53 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -70,29 +70,26 @@ public Row next() { return cachedRow; } - Row currentRow = getCurrentRow(); - - // return it directly if the current row is the last row + cachedRow = getCurrentRow(); + // return it directly if there is no more rows if (!rowIterator().hasNext()) { - return currentRow; + return cachedRow; } - - Row nextRow = rowIterator().next(); - - cachedRow = currentRow; cachedRowCount = 1; + cachedNextRow = rowIterator().next(); + // pull rows from the iterator until two consecutive rows are different - while (isSameRecord(currentRow, nextRow, indicesToIdentifySameRow)) { - if (oppositeChangeType(currentRow, nextRow)) { + while (isSameRecord(cachedRow, cachedNextRow, indicesToIdentifySameRow)) { + if (oppositeChangeType(cachedRow, cachedNextRow)) { // two rows with opposite change types means no net changes, remove both cachedRowCount--; - nextRow = null; + cachedNextRow = null; } else { // two rows with same change types means potential net changes, cache the next row, reset it // to null cachedRowCount++; - nextRow = null; + cachedNextRow = null; } // stop pulling rows if there is no more rows or the next row is different @@ -100,11 +97,9 @@ public Row next() { break; } - nextRow = rowIterator().next(); + cachedNextRow = rowIterator().next(); } - // if they are different rows, hit the boundary, cache the next row - cachedNextRow = nextRow; return null; } @@ -138,5 +133,4 @@ private int[] generateIndicesToIdentifySameRow() { } return indices; } - } From 7b20c7c3126f0ab596167664a4b98cfa4bdefb53 Mon Sep 17 00:00:00 2001 From: yufei Date: Tue, 27 Jun 2023 11:49:08 -0700 Subject: [PATCH 15/20] Resolve comments --- .../iceberg/spark/ChangelogIterator.java | 14 +++++++++++++ .../spark/RemoveCarryoverIterator.java | 13 ++++-------- .../spark/RemoveNetCarryoverIterator.java | 20 ++++++++----------- 3 files changed, 26 insertions(+), 21 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java index 9e6cec44c4a5..c4dd203fa14c 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java @@ -20,6 +20,7 @@ import java.util.Iterator; import java.util.Objects; +import java.util.Set; import org.apache.iceberg.ChangelogOperation; import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.relocated.com.google.common.collect.Iterators; @@ -101,4 +102,17 @@ protected boolean isSameRecord(Row currentRow, Row nextRow, int[] indicesToIdent protected boolean isDifferentValue(Row currentRow, Row nextRow, int idx) { return !Objects.equals(nextRow.get(idx), currentRow.get(idx)); } + + protected static int[] generateIndicesToIdentifySameRow( + int totalColumnCount, Set metadataColumnIndices) { + int[] indices = new int[totalColumnCount - metadataColumnIndices.size()]; + + for (int i = 0, j = 0; i < indices.length; i++) { + if (!metadataColumnIndices.contains(i)) { + indices[j] = i; + j++; + } + } + return indices; + } } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java index f562aa2ee54f..935fb299d26c 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java @@ -19,6 +19,8 @@ package org.apache.iceberg.spark; import java.util.Iterator; +import java.util.Set; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.StructType; @@ -142,14 +144,7 @@ private boolean hasCachedDeleteRow() { } private int[] generateIndicesToIdentifySameRow() { - int[] indices = new int[rowType.size() - 1]; - for (int i = 0; i < indices.length; i++) { - if (i < changeTypeIndex()) { - indices[i] = i; - } else { - indices[i] = i + 1; - } - } - return indices; + Set metadataColumnIndices = Sets.newHashSet(changeTypeIndex()); + return generateIndicesToIdentifySameRow(rowType.size(), metadataColumnIndices); } } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index f72301078c53..02173913b7e5 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -19,7 +19,9 @@ package org.apache.iceberg.spark; import java.util.Iterator; +import java.util.Set; import org.apache.iceberg.MetadataColumns; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.StructType; @@ -120,17 +122,11 @@ private boolean oppositeChangeType(Row currentRow, Row nextRow) { } private int[] generateIndicesToIdentifySameRow() { - int changeOrdinalIndex = rowType.fieldIndex(MetadataColumns.CHANGE_ORDINAL.name()); - int snapshotIdIndex = rowType.fieldIndex(MetadataColumns.COMMIT_SNAPSHOT_ID.name()); - - int[] indices = new int[rowType.size() - 3]; - - for (int i = 0, j = 0; i < indices.length; i++) { - if (i != changeTypeIndex() && i != changeOrdinalIndex && i != snapshotIdIndex) { - indices[j] = i; - j++; - } - } - return indices; + Set metadataColumnIndices = + Sets.newHashSet( + rowType.fieldIndex(MetadataColumns.CHANGE_ORDINAL.name()), + rowType.fieldIndex(MetadataColumns.COMMIT_SNAPSHOT_ID.name()), + changeTypeIndex()); + return generateIndicesToIdentifySameRow(rowType.size(), metadataColumnIndices); } } From 15d89445e429af034f8790f542290d303f569887 Mon Sep 17 00:00:00 2001 From: yufei Date: Tue, 27 Jun 2023 11:53:02 -0700 Subject: [PATCH 16/20] Resolve comments --- gradle.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/gradle.properties b/gradle.properties index ffb2454dcdde..eb0da0ac8547 100644 --- a/gradle.properties +++ b/gradle.properties @@ -20,7 +20,7 @@ systemProp.defaultFlinkVersions=1.17 systemProp.knownFlinkVersions=1.15,1.16,1.17 systemProp.defaultHiveVersions=2 systemProp.knownHiveVersions=2,3 -systemProp.defaultSparkVersions=3.4,3.3 +systemProp.defaultSparkVersions=3.4 systemProp.knownSparkVersions=3.1,3.2,3.3,3.4 systemProp.defaultScalaVersion=2.12 systemProp.knownScalaVersions=2.12,2.13 From bca5bba949ee5941743392c752196cb222308ed1 Mon Sep 17 00:00:00 2001 From: yufei Date: Tue, 27 Jun 2023 15:07:41 -0700 Subject: [PATCH 17/20] Resolve comments --- .../apache/iceberg/spark/RemoveNetCarryoverIterator.java | 7 +++---- .../spark/procedures/CreateChangelogViewProcedure.java | 7 +++++++ 2 files changed, 10 insertions(+), 4 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index 02173913b7e5..1c6362703ce4 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -86,16 +86,15 @@ public Row next() { if (oppositeChangeType(cachedRow, cachedNextRow)) { // two rows with opposite change types means no net changes, remove both cachedRowCount--; - cachedNextRow = null; } else { - // two rows with same change types means potential net changes, cache the next row, reset it - // to null + // two rows with same change types means potential net changes, cache the next row cachedRowCount++; - cachedNextRow = null; } // stop pulling rows if there is no more rows or the next row is different if (cachedRowCount <= 0 || !rowIterator().hasNext()) { + // reset the cached next row if there is no more rows + cachedNextRow = null; break; } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index b2cecb8cf15f..aea6586cbd1d 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -91,6 +91,13 @@ public class CreateChangelogViewProcedure extends BaseProcedure { ProcedureParameter.optional("options", STRING_MAP); private static final ProcedureParameter COMPUTE_UPDATES_PARAM = ProcedureParameter.optional("compute_updates", DataTypes.BooleanType); + + /** + * Enable or disable the remove carry-over rows. + * + * @deprecated since 1.4.0, will be removed in 1.5.0; The procedure will always remove carry-over rows. + * Please query {@link SparkChangelogTable} instead for the use cases doesn't remove carry-over rows. + */ private static final ProcedureParameter REMOVE_CARRYOVERS_PARAM = ProcedureParameter.optional("remove_carryovers", DataTypes.BooleanType); private static final ProcedureParameter IDENTIFIER_COLUMNS_PARAM = From ed90942a8fc9827ef483ba20a7ce545678023b16 Mon Sep 17 00:00:00 2001 From: yufei Date: Tue, 27 Jun 2023 15:13:27 -0700 Subject: [PATCH 18/20] Style fix --- .../spark/procedures/CreateChangelogViewProcedure.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index aea6586cbd1d..5716ed5b3ae3 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -95,11 +95,13 @@ public class CreateChangelogViewProcedure extends BaseProcedure { /** * Enable or disable the remove carry-over rows. * - * @deprecated since 1.4.0, will be removed in 1.5.0; The procedure will always remove carry-over rows. - * Please query {@link SparkChangelogTable} instead for the use cases doesn't remove carry-over rows. + * @deprecated since 1.4.0, will be removed in 1.5.0; The procedure will always remove carry-over + * rows. Please query {@link SparkChangelogTable} instead for the use cases doesn't remove + * carry-over rows. */ private static final ProcedureParameter REMOVE_CARRYOVERS_PARAM = ProcedureParameter.optional("remove_carryovers", DataTypes.BooleanType); + private static final ProcedureParameter IDENTIFIER_COLUMNS_PARAM = ProcedureParameter.optional("identifier_columns", STRING_ARRAY); private static final ProcedureParameter NET_CHANGES = From ac94c5d3d7e1451c90f57f6d697129aad12463e4 Mon Sep 17 00:00:00 2001 From: yufei Date: Tue, 27 Jun 2023 15:27:15 -0700 Subject: [PATCH 19/20] Style fix --- .../iceberg/spark/procedures/CreateChangelogViewProcedure.java | 1 + 1 file changed, 1 insertion(+) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java index 5716ed5b3ae3..259254aa2d51 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/procedures/CreateChangelogViewProcedure.java @@ -99,6 +99,7 @@ public class CreateChangelogViewProcedure extends BaseProcedure { * rows. Please query {@link SparkChangelogTable} instead for the use cases doesn't remove * carry-over rows. */ + @Deprecated private static final ProcedureParameter REMOVE_CARRYOVERS_PARAM = ProcedureParameter.optional("remove_carryovers", DataTypes.BooleanType); From 7a5607a5fadd2a8d20b65f354e4536d1811aa4a5 Mon Sep 17 00:00:00 2001 From: yufei Date: Thu, 29 Jun 2023 13:49:00 -0700 Subject: [PATCH 20/20] Resolve comments --- .../apache/iceberg/spark/ChangelogIterator.java | 11 ++++++++++- .../iceberg/spark/RemoveCarryoverIterator.java | 4 +--- .../iceberg/spark/RemoveNetCarryoverIterator.java | 14 ++++++-------- 3 files changed, 17 insertions(+), 12 deletions(-) diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java index c4dd203fa14c..cc44b1f3992c 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/ChangelogIterator.java @@ -23,6 +23,7 @@ import java.util.Set; import org.apache.iceberg.ChangelogOperation; import org.apache.iceberg.MetadataColumns; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Iterators; import org.apache.spark.sql.Row; import org.apache.spark.sql.types.StructType; @@ -36,9 +37,11 @@ public abstract class ChangelogIterator implements Iterator { private final Iterator rowIterator; private final int changeTypeIndex; + private final StructType rowType; protected ChangelogIterator(Iterator rowIterator, StructType rowType) { this.rowIterator = rowIterator; + this.rowType = rowType; this.changeTypeIndex = rowType.fieldIndex(MetadataColumns.CHANGE_TYPE.name()); } @@ -46,8 +49,14 @@ protected int changeTypeIndex() { return changeTypeIndex; } + protected StructType rowType() { + return rowType; + } + protected String changeType(Row row) { - return row.getString(changeTypeIndex()); + String changeType = row.getString(changeTypeIndex()); + Preconditions.checkNotNull(changeType, "Change type should not be null"); + return changeType; } protected Iterator rowIterator() { diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java index 935fb299d26c..2e90dc7749d1 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveCarryoverIterator.java @@ -50,7 +50,6 @@ */ class RemoveCarryoverIterator extends ChangelogIterator { private final int[] indicesToIdentifySameRow; - private final StructType rowType; private Row cachedDeletedRow = null; private long deletedRowCount = 0; @@ -58,7 +57,6 @@ class RemoveCarryoverIterator extends ChangelogIterator { RemoveCarryoverIterator(Iterator rowIterator, StructType rowType) { super(rowIterator, rowType); - this.rowType = rowType; this.indicesToIdentifySameRow = generateIndicesToIdentifySameRow(); } @@ -145,6 +143,6 @@ private boolean hasCachedDeleteRow() { private int[] generateIndicesToIdentifySameRow() { Set metadataColumnIndices = Sets.newHashSet(changeTypeIndex()); - return generateIndicesToIdentifySameRow(rowType.size(), metadataColumnIndices); + return generateIndicesToIdentifySameRow(rowType().size(), metadataColumnIndices); } } diff --git a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java index 1c6362703ce4..e4234755cdcf 100644 --- a/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java +++ b/spark/v3.3/spark/src/main/java/org/apache/iceberg/spark/RemoveNetCarryoverIterator.java @@ -39,15 +39,13 @@ public class RemoveNetCarryoverIterator extends ChangelogIterator { private final int[] indicesToIdentifySameRow; - private final StructType rowType; - private Row cachedNextRow = null; - private Row cachedRow = null; - private long cachedRowCount = 0; + private Row cachedNextRow; + private Row cachedRow; + private long cachedRowCount; protected RemoveNetCarryoverIterator(Iterator rowIterator, StructType rowType) { super(rowIterator, rowType); - this.rowType = rowType; this.indicesToIdentifySameRow = generateIndicesToIdentifySameRow(); } @@ -123,9 +121,9 @@ private boolean oppositeChangeType(Row currentRow, Row nextRow) { private int[] generateIndicesToIdentifySameRow() { Set metadataColumnIndices = Sets.newHashSet( - rowType.fieldIndex(MetadataColumns.CHANGE_ORDINAL.name()), - rowType.fieldIndex(MetadataColumns.COMMIT_SNAPSHOT_ID.name()), + rowType().fieldIndex(MetadataColumns.CHANGE_ORDINAL.name()), + rowType().fieldIndex(MetadataColumns.COMMIT_SNAPSHOT_ID.name()), changeTypeIndex()); - return generateIndicesToIdentifySameRow(rowType.size(), metadataColumnIndices); + return generateIndicesToIdentifySameRow(rowType().size(), metadataColumnIndices); } }