anoopj commented on code in PR #17347:
URL: https://github.com/apache/iceberg/pull/17347#discussion_r3658451788
##########
spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreaming.java:
##########
@@ -146,6 +148,71 @@ public void testStreamingWriteAppendMode() throws
Exception {
}
}
+ @Test
+ public void testStreamingWriteAppendModeWithMergeAppend() throws Exception {
+ File parent = temp.resolve("parquet").toFile();
+ File location = new File(parent, "test-table");
+ File checkpoint = new File(parent, "checkpoint");
+
+ HadoopTables tables = new HadoopTables(CONF);
+ PartitionSpec spec =
PartitionSpec.builderFor(SCHEMA).identity("data").build();
+ // set a low min-count-to-merge so the merge append visibly consolidates
manifests
+ Table table =
+ tables.create(
+ SCHEMA,
+ spec,
+ ImmutableMap.of(TableProperties.MANIFEST_MIN_MERGE_COUNT, "2"),
+ location.toString());
+
+ List<SimpleRecord> expected =
+ Lists.newArrayList(
+ new SimpleRecord(1, "1"),
+ new SimpleRecord(2, "2"),
+ new SimpleRecord(3, "3"),
+ new SimpleRecord(4, "4"));
+
+ MemoryStream<Integer> inputStream = newMemoryStream(1, spark,
Encoders.INT());
+ DataStreamWriter<Row> streamWriter =
+ inputStream
+ .toDF()
+ .selectExpr("value AS id", "CAST (value AS STRING) AS data")
+ .writeStream()
+ .outputMode("append")
+ .format("iceberg")
+ .option("checkpointLocation", checkpoint.toString())
+ .option("path", location.toString())
+ .option(SparkWriteOptions.STREAMING_MERGE_APPEND_ENABLED, "true");
+
+ try {
+ StreamingQuery query = streamWriter.start();
+ List<Integer> batch1 = Lists.newArrayList(1, 2);
+ send(batch1, inputStream);
+ query.processAllAvailable();
+ List<Integer> batch2 = Lists.newArrayList(3, 4);
+ send(batch2, inputStream);
+ query.processAllAvailable();
+ query.stop();
+
+ Dataset<Row> result =
spark.read().format("iceberg").load(location.toString());
+ List<SimpleRecord> actual =
+
result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+
+ assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
+ assertThat(table.snapshots()).as("Number of snapshots should
match").hasSize(2);
+
+ // the second streaming append uses the regular append, which merges the
two data manifests
+ // into a single manifest; a fast append would have left two separate
manifests
+ table.refresh();
Review Comment:
Should the snapshot count assertion on line 201 be run _after_
`table.refresh()`? (similar to the assertion on line 206)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]