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; 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 {