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]

Reply via email to