From b6e3e7e7dc7f1eab0fc0b3d5b49c7e3c42210176 Mon Sep 17 00:00:00 2001 From: Sreeram Garlapati Date: Fri, 25 Jun 2021 19:23:51 -0700 Subject: [PATCH 1/8] DSV2 micro_batch read - implement skipReplace and skipDelete (#1) * idea checkpoint * implement skipDelete and skipReplace options * revert changes in SnapshotUtil --- .../iceberg/spark/SparkReadOptions.java | 6 +++ .../org/apache/iceberg/spark/Spark3Util.java | 13 +++++ .../spark/source/SparkMicroBatchStream.java | 49 +++++++++++++------ .../source/TestStructuredStreamingRead3.java | 27 ++++++++++ 4 files changed, 79 insertions(+), 16 deletions(-) diff --git a/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java b/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java index bce0bf4e8bb5..f21d4fd0a344 100644 --- a/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java +++ b/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java @@ -50,4 +50,10 @@ private SparkReadOptions() { // Set ID that is used to fetch file scan tasks public static final String FILE_SCAN_TASK_SET_ID = "file-scan-task-set-id"; + + // skip snapshots of type delete while reading stream out of iceberg table + public static final String READ_STREAM_SKIP_DELETE = "read-stream-skip-delete"; + + // skip snapshots of type replace while reading stream out of iceberg table + public static final String READ_STREAM_SKIP_REPLACE = "read-stream-skip-replace"; } diff --git a/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java b/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java index 9a24b3e5ffb6..7b52ba118890 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java @@ -562,6 +562,19 @@ public static Integer propertyAsInt(CaseInsensitiveStringMap options, String pro return null; } + public static Boolean propertyAsBoolean(CaseInsensitiveStringMap options, String property, Boolean defaultValue) { + if (defaultValue != null) { + return options.getBoolean(property, defaultValue); + } + + String value = options.get(property); + if (value != null) { + return Boolean.parseBoolean(value); + } + + return null; + } + public static class DescribeSchemaVisitor extends TypeUtil.SchemaVisitor { private static final Joiner COMMA = Joiner.on(','); private static final DescribeSchemaVisitor INSTANCE = new DescribeSchemaVisitor(); diff --git a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java index 1a7f217af7f9..90e7dfd896c0 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java @@ -27,6 +27,7 @@ import java.io.UncheckedIOException; import java.nio.charset.StandardCharsets; import java.util.List; +import java.util.Set; import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.DataOperations; import org.apache.iceberg.FileScanTask; @@ -45,6 +46,7 @@ import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; 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.Spark3Util; import org.apache.iceberg.spark.SparkReadOptions; import org.apache.iceberg.spark.source.SparkBatchScan.ReadTask; @@ -79,6 +81,7 @@ public class SparkMicroBatchStream implements MicroBatchStream { private final Long splitOpenFileCost; private final boolean localityPreferred; private final StreamingOffset initialOffset; + private final Set skippableDataOperations; SparkMicroBatchStream(JavaSparkContext sparkContext, Table table, boolean caseSensitive, Schema expectedSchema, CaseInsensitiveStringMap options, String checkpointLocation) { @@ -100,6 +103,15 @@ public class SparkMicroBatchStream implements MicroBatchStream { InitialOffsetStore initialOffsetStore = new InitialOffsetStore(table, checkpointLocation); this.initialOffset = initialOffsetStore.initialOffset(); + + this.skippableDataOperations = Sets.newHashSet(); + if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_DELETE, false)) { + this.skippableDataOperations.add(DataOperations.DELETE); + } + + if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_REPLACE, false)) { + this.skippableDataOperations.add(DataOperations.REPLACE); + } } @Override @@ -169,37 +181,42 @@ public void stop() { private List planFiles(StreamingOffset startOffset, StreamingOffset endOffset) { List fileScanTasks = Lists.newArrayList(); - MicroBatch latestMicroBatch = null; StreamingOffset batchStartOffset = StreamingOffset.START_OFFSET.equals(startOffset) ? new StreamingOffset(SnapshotUtil.oldestSnapshot(table).snapshotId(), 0, false) : startOffset; + StreamingOffset currentOffset = null; + do { - StreamingOffset currentOffset = - latestMicroBatch != null && latestMicroBatch.lastIndexOfSnapshot() ? - new StreamingOffset(snapshotAfter(latestMicroBatch.snapshotId()), 0L, false) : - batchStartOffset; + if (currentOffset == null) { + currentOffset = batchStartOffset; + } else { + Snapshot snapshotAfter = SnapshotUtil.snapshotAfter(table, currentOffset.snapshotId()); + boolean shouldSkip = skippableDataOperations.contains(snapshotAfter.operation()); + + // TODO: fix error message + Preconditions.checkState( + snapshotAfter.operation().equals(DataOperations.APPEND) || shouldSkip, + "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshotAfter.operation()); + + currentOffset = new StreamingOffset(snapshotAfter.snapshotId(), 0L, false); - latestMicroBatch = MicroBatches.from(table.snapshot(currentOffset.snapshotId()), table.io()) + if (shouldSkip) { + continue; + } + } + + MicroBatch latestMicroBatch = MicroBatches.from(table.snapshot(currentOffset.snapshotId()), table.io()) .caseSensitive(caseSensitive) .specsById(table.specs()) .generate(currentOffset.position(), Long.MAX_VALUE, currentOffset.shouldScanAllFiles()); fileScanTasks.addAll(latestMicroBatch.tasks()); - } while (latestMicroBatch.snapshotId() != endOffset.snapshotId()); + } while (currentOffset.snapshotId() != endOffset.snapshotId()); return fileScanTasks; } - private long snapshotAfter(long snapshotId) { - Snapshot snapshotAfter = SnapshotUtil.snapshotAfter(table, snapshotId); - - Preconditions.checkState(snapshotAfter.operation().equals(DataOperations.APPEND), - "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshotAfter.operation()); - - return snapshotAfter.snapshotId(); - } - private static class InitialOffsetStore { private final Table table; private final FileIO io; diff --git a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java index 07f3df4ea4aa..bb2169bdb132 100644 --- a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java +++ b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java @@ -42,6 +42,7 @@ import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.spark.SparkCatalogTestBase; +import org.apache.iceberg.spark.SparkReadOptions; import org.apache.iceberg.types.Types; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Encoders; @@ -321,6 +322,32 @@ public void testReadStreamWithSnapshotTypeReplaceErrorsOut() throws Exception { ); } + @SuppressWarnings("unchecked") + @Test + public void testReadStreamWithSnapshotTypeReplaceAndSkipReplaceOption() throws Exception { + // fill table with some data + List> dataAcrossSnapshots = TEST_DATA_MULTIPLE_SNAPSHOTS; + appendDataAsMultipleSnapshots(dataAcrossSnapshots, tableIdentifier); + + table.refresh(); + + // this should create a snapshot with type Replace. + table.rewriteManifests() + .clusterBy(f -> 1) + .commit(); + + // check pre-condition + Assert.assertEquals(DataOperations.REPLACE, table.currentSnapshot().operation()); + + Dataset df = spark.readStream() + .format("iceberg") + .option(SparkReadOptions.READ_STREAM_SKIP_REPLACE, "true") + .load(tableIdentifier); + + Assertions.assertThat(processAvailable(df)) + .containsExactlyInAnyOrderElementsOf(Iterables.concat(dataAcrossSnapshots)); + } + @SuppressWarnings("unchecked") @Test public void testReadStreamWithSnapshotTypeDeleteErrorsOut() throws Exception { From 51df3d3f9b5e4f98c733bacedd5d7a64f243b01c Mon Sep 17 00:00:00 2001 From: Sreeram Garlapati Date: Fri, 25 Jun 2021 19:25:45 -0700 Subject: [PATCH 2/8] implement skipDelete and skipReplace (#2) * implement skipDelete and skipReplace options * revert changes in SnapshotUtil --- .../iceberg/spark/SparkReadOptions.java | 6 +++ .../org/apache/iceberg/spark/Spark3Util.java | 13 +++++ .../spark/source/SparkMicroBatchStream.java | 49 +++++++++++++------ .../source/TestStructuredStreamingRead3.java | 27 ++++++++++ 4 files changed, 79 insertions(+), 16 deletions(-) diff --git a/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java b/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java index bce0bf4e8bb5..f21d4fd0a344 100644 --- a/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java +++ b/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java @@ -50,4 +50,10 @@ private SparkReadOptions() { // Set ID that is used to fetch file scan tasks public static final String FILE_SCAN_TASK_SET_ID = "file-scan-task-set-id"; + + // skip snapshots of type delete while reading stream out of iceberg table + public static final String READ_STREAM_SKIP_DELETE = "read-stream-skip-delete"; + + // skip snapshots of type replace while reading stream out of iceberg table + public static final String READ_STREAM_SKIP_REPLACE = "read-stream-skip-replace"; } diff --git a/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java b/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java index 9a24b3e5ffb6..7b52ba118890 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java @@ -562,6 +562,19 @@ public static Integer propertyAsInt(CaseInsensitiveStringMap options, String pro return null; } + public static Boolean propertyAsBoolean(CaseInsensitiveStringMap options, String property, Boolean defaultValue) { + if (defaultValue != null) { + return options.getBoolean(property, defaultValue); + } + + String value = options.get(property); + if (value != null) { + return Boolean.parseBoolean(value); + } + + return null; + } + public static class DescribeSchemaVisitor extends TypeUtil.SchemaVisitor { private static final Joiner COMMA = Joiner.on(','); private static final DescribeSchemaVisitor INSTANCE = new DescribeSchemaVisitor(); diff --git a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java index 1a7f217af7f9..90e7dfd896c0 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java @@ -27,6 +27,7 @@ import java.io.UncheckedIOException; import java.nio.charset.StandardCharsets; import java.util.List; +import java.util.Set; import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.DataOperations; import org.apache.iceberg.FileScanTask; @@ -45,6 +46,7 @@ import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; 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.Spark3Util; import org.apache.iceberg.spark.SparkReadOptions; import org.apache.iceberg.spark.source.SparkBatchScan.ReadTask; @@ -79,6 +81,7 @@ public class SparkMicroBatchStream implements MicroBatchStream { private final Long splitOpenFileCost; private final boolean localityPreferred; private final StreamingOffset initialOffset; + private final Set skippableDataOperations; SparkMicroBatchStream(JavaSparkContext sparkContext, Table table, boolean caseSensitive, Schema expectedSchema, CaseInsensitiveStringMap options, String checkpointLocation) { @@ -100,6 +103,15 @@ public class SparkMicroBatchStream implements MicroBatchStream { InitialOffsetStore initialOffsetStore = new InitialOffsetStore(table, checkpointLocation); this.initialOffset = initialOffsetStore.initialOffset(); + + this.skippableDataOperations = Sets.newHashSet(); + if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_DELETE, false)) { + this.skippableDataOperations.add(DataOperations.DELETE); + } + + if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_REPLACE, false)) { + this.skippableDataOperations.add(DataOperations.REPLACE); + } } @Override @@ -169,37 +181,42 @@ public void stop() { private List planFiles(StreamingOffset startOffset, StreamingOffset endOffset) { List fileScanTasks = Lists.newArrayList(); - MicroBatch latestMicroBatch = null; StreamingOffset batchStartOffset = StreamingOffset.START_OFFSET.equals(startOffset) ? new StreamingOffset(SnapshotUtil.oldestSnapshot(table).snapshotId(), 0, false) : startOffset; + StreamingOffset currentOffset = null; + do { - StreamingOffset currentOffset = - latestMicroBatch != null && latestMicroBatch.lastIndexOfSnapshot() ? - new StreamingOffset(snapshotAfter(latestMicroBatch.snapshotId()), 0L, false) : - batchStartOffset; + if (currentOffset == null) { + currentOffset = batchStartOffset; + } else { + Snapshot snapshotAfter = SnapshotUtil.snapshotAfter(table, currentOffset.snapshotId()); + boolean shouldSkip = skippableDataOperations.contains(snapshotAfter.operation()); + + // TODO: fix error message + Preconditions.checkState( + snapshotAfter.operation().equals(DataOperations.APPEND) || shouldSkip, + "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshotAfter.operation()); + + currentOffset = new StreamingOffset(snapshotAfter.snapshotId(), 0L, false); - latestMicroBatch = MicroBatches.from(table.snapshot(currentOffset.snapshotId()), table.io()) + if (shouldSkip) { + continue; + } + } + + MicroBatch latestMicroBatch = MicroBatches.from(table.snapshot(currentOffset.snapshotId()), table.io()) .caseSensitive(caseSensitive) .specsById(table.specs()) .generate(currentOffset.position(), Long.MAX_VALUE, currentOffset.shouldScanAllFiles()); fileScanTasks.addAll(latestMicroBatch.tasks()); - } while (latestMicroBatch.snapshotId() != endOffset.snapshotId()); + } while (currentOffset.snapshotId() != endOffset.snapshotId()); return fileScanTasks; } - private long snapshotAfter(long snapshotId) { - Snapshot snapshotAfter = SnapshotUtil.snapshotAfter(table, snapshotId); - - Preconditions.checkState(snapshotAfter.operation().equals(DataOperations.APPEND), - "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshotAfter.operation()); - - return snapshotAfter.snapshotId(); - } - private static class InitialOffsetStore { private final Table table; private final FileIO io; diff --git a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java index 07f3df4ea4aa..bb2169bdb132 100644 --- a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java +++ b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java @@ -42,6 +42,7 @@ import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.spark.SparkCatalogTestBase; +import org.apache.iceberg.spark.SparkReadOptions; import org.apache.iceberg.types.Types; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Encoders; @@ -321,6 +322,32 @@ public void testReadStreamWithSnapshotTypeReplaceErrorsOut() throws Exception { ); } + @SuppressWarnings("unchecked") + @Test + public void testReadStreamWithSnapshotTypeReplaceAndSkipReplaceOption() throws Exception { + // fill table with some data + List> dataAcrossSnapshots = TEST_DATA_MULTIPLE_SNAPSHOTS; + appendDataAsMultipleSnapshots(dataAcrossSnapshots, tableIdentifier); + + table.refresh(); + + // this should create a snapshot with type Replace. + table.rewriteManifests() + .clusterBy(f -> 1) + .commit(); + + // check pre-condition + Assert.assertEquals(DataOperations.REPLACE, table.currentSnapshot().operation()); + + Dataset df = spark.readStream() + .format("iceberg") + .option(SparkReadOptions.READ_STREAM_SKIP_REPLACE, "true") + .load(tableIdentifier); + + Assertions.assertThat(processAvailable(df)) + .containsExactlyInAnyOrderElementsOf(Iterables.concat(dataAcrossSnapshots)); + } + @SuppressWarnings("unchecked") @Test public void testReadStreamWithSnapshotTypeDeleteErrorsOut() throws Exception { From 40eeb255abeb55fce9d3ed2548514475f1e85612 Mon Sep 17 00:00:00 2001 From: Daksha Asrani Date: Mon, 28 Jun 2021 18:05:26 -0700 Subject: [PATCH 3/8] Revert "DSV2 micro_batch read - implement skipReplace and skipDelete (#1)" This reverts commit b6e3e7e7dc7f1eab0fc0b3d5b49c7e3c42210176. --- .../iceberg/spark/SparkReadOptions.java | 6 --- .../org/apache/iceberg/spark/Spark3Util.java | 13 ----- .../spark/source/SparkMicroBatchStream.java | 49 ++++++------------- .../source/TestStructuredStreamingRead3.java | 27 ---------- 4 files changed, 16 insertions(+), 79 deletions(-) diff --git a/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java b/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java index f21d4fd0a344..bce0bf4e8bb5 100644 --- a/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java +++ b/spark/src/main/java/org/apache/iceberg/spark/SparkReadOptions.java @@ -50,10 +50,4 @@ private SparkReadOptions() { // Set ID that is used to fetch file scan tasks public static final String FILE_SCAN_TASK_SET_ID = "file-scan-task-set-id"; - - // skip snapshots of type delete while reading stream out of iceberg table - public static final String READ_STREAM_SKIP_DELETE = "read-stream-skip-delete"; - - // skip snapshots of type replace while reading stream out of iceberg table - public static final String READ_STREAM_SKIP_REPLACE = "read-stream-skip-replace"; } diff --git a/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java b/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java index 7b52ba118890..9a24b3e5ffb6 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/Spark3Util.java @@ -562,19 +562,6 @@ public static Integer propertyAsInt(CaseInsensitiveStringMap options, String pro return null; } - public static Boolean propertyAsBoolean(CaseInsensitiveStringMap options, String property, Boolean defaultValue) { - if (defaultValue != null) { - return options.getBoolean(property, defaultValue); - } - - String value = options.get(property); - if (value != null) { - return Boolean.parseBoolean(value); - } - - return null; - } - public static class DescribeSchemaVisitor extends TypeUtil.SchemaVisitor { private static final Joiner COMMA = Joiner.on(','); private static final DescribeSchemaVisitor INSTANCE = new DescribeSchemaVisitor(); diff --git a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java index bacd15067245..84fdbd288de7 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java @@ -27,7 +27,6 @@ import java.io.UncheckedIOException; import java.nio.charset.StandardCharsets; import java.util.List; -import java.util.Set; import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.DataOperations; import org.apache.iceberg.FileScanTask; @@ -46,7 +45,6 @@ import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; 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.Spark3Util; import org.apache.iceberg.spark.SparkReadOptions; import org.apache.iceberg.spark.source.SparkBatchScan.ReadTask; @@ -81,7 +79,6 @@ public class SparkMicroBatchStream implements MicroBatchStream { private final Long splitOpenFileCost; private final boolean localityPreferred; private final StreamingOffset initialOffset; - private final Set skippableDataOperations; SparkMicroBatchStream(JavaSparkContext sparkContext, Table table, boolean caseSensitive, Schema expectedSchema, CaseInsensitiveStringMap options, String checkpointLocation) { @@ -104,15 +101,6 @@ public class SparkMicroBatchStream implements MicroBatchStream { InitialOffsetStore initialOffsetStore = new InitialOffsetStore(table, checkpointLocation); this.initialOffset = initialOffsetStore.initialOffset(); - - this.skippableDataOperations = Sets.newHashSet(); - if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_DELETE, false)) { - this.skippableDataOperations.add(DataOperations.DELETE); - } - - if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_REPLACE, false)) { - this.skippableDataOperations.add(DataOperations.REPLACE); - } } @Override @@ -182,42 +170,37 @@ public void stop() { private List planFiles(StreamingOffset startOffset, StreamingOffset endOffset) { List fileScanTasks = Lists.newArrayList(); + MicroBatch latestMicroBatch = null; StreamingOffset batchStartOffset = StreamingOffset.START_OFFSET.equals(startOffset) ? new StreamingOffset(SnapshotUtil.oldestSnapshot(table).snapshotId(), 0, false) : startOffset; - StreamingOffset currentOffset = null; - do { - if (currentOffset == null) { - currentOffset = batchStartOffset; - } else { - Snapshot snapshotAfter = SnapshotUtil.snapshotAfter(table, currentOffset.snapshotId()); - boolean shouldSkip = skippableDataOperations.contains(snapshotAfter.operation()); - - // TODO: fix error message - Preconditions.checkState( - snapshotAfter.operation().equals(DataOperations.APPEND) || shouldSkip, - "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshotAfter.operation()); - - currentOffset = new StreamingOffset(snapshotAfter.snapshotId(), 0L, false); + StreamingOffset currentOffset = + latestMicroBatch != null && latestMicroBatch.lastIndexOfSnapshot() ? + new StreamingOffset(snapshotAfter(latestMicroBatch.snapshotId()), 0L, false) : + batchStartOffset; - if (shouldSkip) { - continue; - } - } - - MicroBatch latestMicroBatch = MicroBatches.from(table.snapshot(currentOffset.snapshotId()), table.io()) + latestMicroBatch = MicroBatches.from(table.snapshot(currentOffset.snapshotId()), table.io()) .caseSensitive(caseSensitive) .specsById(table.specs()) .generate(currentOffset.position(), Long.MAX_VALUE, currentOffset.shouldScanAllFiles()); fileScanTasks.addAll(latestMicroBatch.tasks()); - } while (currentOffset.snapshotId() != endOffset.snapshotId()); + } while (latestMicroBatch.snapshotId() != endOffset.snapshotId()); return fileScanTasks; } + private long snapshotAfter(long snapshotId) { + Snapshot snapshotAfter = SnapshotUtil.snapshotAfter(table, snapshotId); + + Preconditions.checkState(snapshotAfter.operation().equals(DataOperations.APPEND), + "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshotAfter.operation()); + + return snapshotAfter.snapshotId(); + } + private static class InitialOffsetStore { private final Table table; private final FileIO io; diff --git a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java index bb2169bdb132..07f3df4ea4aa 100644 --- a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java +++ b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java @@ -42,7 +42,6 @@ import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.spark.SparkCatalogTestBase; -import org.apache.iceberg.spark.SparkReadOptions; import org.apache.iceberg.types.Types; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Encoders; @@ -322,32 +321,6 @@ public void testReadStreamWithSnapshotTypeReplaceErrorsOut() throws Exception { ); } - @SuppressWarnings("unchecked") - @Test - public void testReadStreamWithSnapshotTypeReplaceAndSkipReplaceOption() throws Exception { - // fill table with some data - List> dataAcrossSnapshots = TEST_DATA_MULTIPLE_SNAPSHOTS; - appendDataAsMultipleSnapshots(dataAcrossSnapshots, tableIdentifier); - - table.refresh(); - - // this should create a snapshot with type Replace. - table.rewriteManifests() - .clusterBy(f -> 1) - .commit(); - - // check pre-condition - Assert.assertEquals(DataOperations.REPLACE, table.currentSnapshot().operation()); - - Dataset df = spark.readStream() - .format("iceberg") - .option(SparkReadOptions.READ_STREAM_SKIP_REPLACE, "true") - .load(tableIdentifier); - - Assertions.assertThat(processAvailable(df)) - .containsExactlyInAnyOrderElementsOf(Iterables.concat(dataAcrossSnapshots)); - } - @SuppressWarnings("unchecked") @Test public void testReadStreamWithSnapshotTypeDeleteErrorsOut() throws Exception { From 7bb4f90cba9c5320e85ea5e44b178320c9b9912b Mon Sep 17 00:00:00 2001 From: Daksha Asrani Date: Mon, 28 Jun 2021 18:25:52 -0700 Subject: [PATCH 4/8] Spark3 Streaming Read: Adding options to skip delete and replace snapshots --- .../spark/source/SparkMicroBatchStream.java | 51 ++++++++++++------- .../source/TestStructuredStreamingRead3.java | 39 ++++++++++++++ 2 files changed, 73 insertions(+), 17 deletions(-) diff --git a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java index 90e7dfd896c0..91c36b47ec2b 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java @@ -27,7 +27,6 @@ import java.io.UncheckedIOException; import java.nio.charset.StandardCharsets; import java.util.List; -import java.util.Set; import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.DataOperations; import org.apache.iceberg.FileScanTask; @@ -46,7 +45,6 @@ import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; 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.Spark3Util; import org.apache.iceberg.spark.SparkReadOptions; import org.apache.iceberg.spark.source.SparkBatchScan.ReadTask; @@ -81,7 +79,8 @@ public class SparkMicroBatchStream implements MicroBatchStream { private final Long splitOpenFileCost; private final boolean localityPreferred; private final StreamingOffset initialOffset; - private final Set skippableDataOperations; + private final boolean skipDelete; + private final boolean skipReplace; SparkMicroBatchStream(JavaSparkContext sparkContext, Table table, boolean caseSensitive, Schema expectedSchema, CaseInsensitiveStringMap options, String checkpointLocation) { @@ -104,14 +103,8 @@ public class SparkMicroBatchStream implements MicroBatchStream { InitialOffsetStore initialOffsetStore = new InitialOffsetStore(table, checkpointLocation); this.initialOffset = initialOffsetStore.initialOffset(); - this.skippableDataOperations = Sets.newHashSet(); - if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_DELETE, false)) { - this.skippableDataOperations.add(DataOperations.DELETE); - } - - if (Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_REPLACE, false)) { - this.skippableDataOperations.add(DataOperations.REPLACE); - } + this.skipDelete = Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_DELETE, false); + this.skipReplace = Spark3Util.propertyAsBoolean(options, SparkReadOptions.READ_STREAM_SKIP_REPLACE, false); } @Override @@ -192,12 +185,7 @@ private List planFiles(StreamingOffset startOffset, StreamingOffse currentOffset = batchStartOffset; } else { Snapshot snapshotAfter = SnapshotUtil.snapshotAfter(table, currentOffset.snapshotId()); - boolean shouldSkip = skippableDataOperations.contains(snapshotAfter.operation()); - - // TODO: fix error message - Preconditions.checkState( - snapshotAfter.operation().equals(DataOperations.APPEND) || shouldSkip, - "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshotAfter.operation()); + boolean shouldSkip = shouldSkip(snapshotAfter); currentOffset = new StreamingOffset(snapshotAfter.snapshotId(), 0L, false); @@ -217,6 +205,35 @@ private List planFiles(StreamingOffset startOffset, StreamingOffse return fileScanTasks; } + private boolean shouldSkip(Snapshot snapshot) { + if (snapshot.operation().equals(DataOperations.DELETE)) { + return shouldSkipDelete(snapshot); + } else if (snapshot.operation().equals(DataOperations.REPLACE)) { + return shouldSkipReplace(snapshot); + } + + Preconditions.checkState( + snapshot.operation().equals(DataOperations.APPEND), + "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshot.operation()); + return false; + } + + private boolean shouldSkipDelete(Snapshot snapshot) { + Preconditions.checkState(skipDelete, + "Invalid Snapshot operation: %s, only APPEND is allowed. To skip delete, set Spark Option %s", + snapshot.operation(), + SparkReadOptions.READ_STREAM_SKIP_DELETE); + return true; + } + + private boolean shouldSkipReplace(Snapshot snapshot) { + Preconditions.checkState(skipReplace, + "Invalid Snapshot operation: %s, only APPEND is allowed. To skip replace, set Spark Option %s", + snapshot.operation(), + SparkReadOptions.READ_STREAM_SKIP_REPLACE); + return true; + } + private static class InitialOffsetStore { private final Table table; private final FileIO io; diff --git a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java index bb2169bdb132..17b7b0cb0755 100644 --- a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java +++ b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java @@ -390,6 +390,45 @@ public void testReadStreamWithSnapshotTypeDeleteErrorsOut() throws Exception { ); } + @SuppressWarnings("unchecked") + @Test + public void testReadStreamWithSnapshotTypeDeleteAndSkipDeleteOption() throws Exception { + table.updateSpec() + .removeField("id_bucket") + .addField(ref("id")) + .commit(); + + table.refresh(); + + // fill table with some data + List> dataAcrossSnapshots = TEST_DATA_MULTIPLE_SNAPSHOTS; + appendDataAsMultipleSnapshots(dataAcrossSnapshots, tableIdentifier); + + table.refresh(); + + // this should create a snapshot with type delete. + table.newDelete() + .deleteFromRowFilter(Expressions.equal("id", 4)) + .commit(); + + // check pre-condition - that the above delete operation on table resulted in Snapshot of Type DELETE. + table.refresh(); + Assert.assertEquals(DataOperations.DELETE, table.currentSnapshot().operation()); + + Dataset df = spark.readStream() + .format("iceberg") + .option(SparkReadOptions.READ_STREAM_SKIP_DELETE, "true") + .load(tableIdentifier); + StreamingQuery streamingQuery = df.writeStream() + .format("memory") + .queryName("testtablewithdelete") + .outputMode(OutputMode.Append()) + .start(); + + Assertions.assertThat(processAvailable(df)) + .containsExactlyInAnyOrderElementsOf(Iterables.concat(dataAcrossSnapshots)); + } + private static List processMicroBatch(DataStreamWriter singleBatchWriter, String viewName) throws TimeoutException, StreamingQueryException { StreamingQuery streamingQuery = singleBatchWriter.start(); From 8a6b56eb16f1fa62d4aed289855870a17eceb6c0 Mon Sep 17 00:00:00 2001 From: Daksha Asrani Date: Mon, 28 Jun 2021 19:32:01 -0700 Subject: [PATCH 5/8] Fixed spacing --- .../spark/source/SparkMicroBatchStream.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java index 91c36b47ec2b..3224f96db0c4 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java @@ -213,24 +213,24 @@ private boolean shouldSkip(Snapshot snapshot) { } Preconditions.checkState( - snapshot.operation().equals(DataOperations.APPEND), - "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshot.operation()); + snapshot.operation().equals(DataOperations.APPEND), + "Invalid Snapshot operation: %s, only APPEND is allowed.", snapshot.operation()); return false; } private boolean shouldSkipDelete(Snapshot snapshot) { Preconditions.checkState(skipDelete, - "Invalid Snapshot operation: %s, only APPEND is allowed. To skip delete, set Spark Option %s", - snapshot.operation(), - SparkReadOptions.READ_STREAM_SKIP_DELETE); + "Invalid Snapshot operation: %s, only APPEND is allowed. To skip delete, set Spark Option %s", + snapshot.operation(), + SparkReadOptions.READ_STREAM_SKIP_DELETE); return true; } private boolean shouldSkipReplace(Snapshot snapshot) { Preconditions.checkState(skipReplace, - "Invalid Snapshot operation: %s, only APPEND is allowed. To skip replace, set Spark Option %s", - snapshot.operation(), - SparkReadOptions.READ_STREAM_SKIP_REPLACE); + "Invalid Snapshot operation: %s, only APPEND is allowed. To skip replace, set Spark Option %s", + snapshot.operation(), + SparkReadOptions.READ_STREAM_SKIP_REPLACE); return true; } From e99870ded26647f5202848b71897352801b5d7d7 Mon Sep 17 00:00:00 2001 From: daksha121 Date: Fri, 16 Jul 2021 15:40:17 -0700 Subject: [PATCH 6/8] Make SparkWrite classes public --- .../main/java/org/apache/iceberg/spark/source/SparkWrite.java | 2 +- .../java/org/apache/iceberg/spark/source/SparkWriteBuilder.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWrite.java b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWrite.java index 213fee4a6dd5..ed3358ec422b 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWrite.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWrite.java @@ -89,7 +89,7 @@ import static org.apache.iceberg.TableProperties.WRITE_TARGET_FILE_SIZE_BYTES; import static org.apache.iceberg.TableProperties.WRITE_TARGET_FILE_SIZE_BYTES_DEFAULT; -class SparkWrite { +public class SparkWrite { private static final Logger LOG = LoggerFactory.getLogger(SparkWrite.class); private final JavaSparkContext sparkContext; diff --git a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWriteBuilder.java b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWriteBuilder.java index b23e0a7935cf..65a1a46ef9a8 100644 --- a/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWriteBuilder.java +++ b/spark3/src/main/java/org/apache/iceberg/spark/source/SparkWriteBuilder.java @@ -43,7 +43,7 @@ import org.apache.spark.sql.types.StructType; import org.apache.spark.sql.util.CaseInsensitiveStringMap; -class SparkWriteBuilder implements WriteBuilder, SupportsDynamicOverwrite, SupportsOverwrite { +public class SparkWriteBuilder implements WriteBuilder, SupportsDynamicOverwrite, SupportsOverwrite { private final SparkSession spark; private final Table table; From 28b70647b014751807e0ba98f37ef61cac605f3a Mon Sep 17 00:00:00 2001 From: daksha121 Date: Fri, 16 Jul 2021 15:46:22 -0700 Subject: [PATCH 7/8] Revert "Merge branch 'master' of https://github.com/daksha121/iceberg" This reverts commit b58f3743e8a9a91fff39233f64cd74df898414cb, reversing changes made to 762f5bb9a9ff2196733dc092196a28c96a65007f. --- .../source/TestStructuredStreamingRead3.java | 26 ------------------- 1 file changed, 26 deletions(-) diff --git a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java index 9b172a3ebf4d..4e50cf2a3e21 100644 --- a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java +++ b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java @@ -313,32 +313,6 @@ public void testReadStreamWithSnapshotTypeReplaceIgnoresReplace() throws Excepti Assertions.assertThat(actual).containsExactlyInAnyOrderElementsOf(Iterables.concat(expected)); } - @SuppressWarnings("unchecked") - @Test - public void testReadStreamWithSnapshotTypeReplaceAndSkipReplaceOption() throws Exception { - // fill table with some data - List> dataAcrossSnapshots = TEST_DATA_MULTIPLE_SNAPSHOTS; - appendDataAsMultipleSnapshots(dataAcrossSnapshots, tableIdentifier); - - table.refresh(); - - // this should create a snapshot with type Replace. - table.rewriteManifests() - .clusterBy(f -> 1) - .commit(); - - // check pre-condition - Assert.assertEquals(DataOperations.REPLACE, table.currentSnapshot().operation()); - - Dataset df = spark.readStream() - .format("iceberg") - .option(SparkReadOptions.READ_STREAM_SKIP_REPLACE, "true") - .load(tableIdentifier); - - Assertions.assertThat(processAvailable(df)) - .containsExactlyInAnyOrderElementsOf(Iterables.concat(dataAcrossSnapshots)); - } - @SuppressWarnings("unchecked") @Test public void testReadStreamWithSnapshotTypeDeleteErrorsOut() throws Exception { From c28a99f16d7f534193096b3ca9beb9ae03a258ef Mon Sep 17 00:00:00 2001 From: daksha121 Date: Wed, 21 Jul 2021 17:01:43 -0700 Subject: [PATCH 8/8] Merge with master --- .../source/TestStructuredStreamingRead3.java | 26 +++++++++++++++++++ 1 file changed, 26 insertions(+) diff --git a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java index 4e50cf2a3e21..2d17fa52189d 100644 --- a/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java +++ b/spark3/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreamingRead3.java @@ -313,6 +313,32 @@ public void testReadStreamWithSnapshotTypeReplaceIgnoresReplace() throws Excepti Assertions.assertThat(actual).containsExactlyInAnyOrderElementsOf(Iterables.concat(expected)); } + @SuppressWarnings("unchecked") + @Test + public void testReadStreamWithSnapshotTypeReplaceAndSkipReplaceOption() throws Exception { + // fill table with some data + List> dataAcrossSnapshots = TEST_DATA_MULTIPLE_SNAPSHOTS; + appendDataAsMultipleSnapshots(dataAcrossSnapshots, tableIdentifier); + + table.refresh(); + + // this should create a snapshot with type Replace. + table.rewriteManifests() + .clusterBy(f -> 1) + .commit(); + + // check pre-condition + Assert.assertEquals(DataOperations.REPLACE, table.currentSnapshot().operation()); + + Dataset df = spark.readStream() + .format("iceberg") + .option(SparkReadOptions.READ_STREAM_SKIP_REPLACE, "true") + .load(tableIdentifier); + + Assertions.assertThat(processAvailable(df)) + .containsExactlyInAnyOrderElementsOf(Iterables.concat(dataAcrossSnapshots)); + } + @SuppressWarnings("unchecked") @Test public void testReadStreamWithSnapshotTypeDeleteErrorsOut() throws Exception {