lxy-9602 commented on code in PR #224:
URL: https://github.com/apache/paimon-cpp/pull/224#discussion_r3879577749


##########
test/inte/realtime_write_inte_test.cpp:
##########
@@ -723,6 +1086,866 @@ TEST_F(RealtimeWriteInteTest, TestAppendCommitAndRead) {
     FinalizeCommitAndCheck(writer.get(), /*realtime_commits=*/{}, 
/*prepare_identifier=*/0, rows);
 }
 
+TEST_F(RealtimeWriteInteTest, TestPkRead) {
+    CreatePkTable();
+    auto saw_query_predicate = std::make_shared<std::atomic<bool>>(false);
+    auto query_view = std::make_shared<std::weak_ptr<RealtimeReadView>>();
+    auto factory =
+        MakeDecoratingFactory<QueryTrackingRealtimeStore>(saw_query_predicate, 
query_view);
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create(factory));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    std::vector<Row> first_rows = {{1, "old", "p0"}, {2, "two", "p0"}, {1, 
"new-in-run", "p0"}};
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> first_batch,
+                         MakeBatch(first_rows, /*partitioned=*/false, 
/*bucket=*/0,
+                                   {RecordBatch::RowKind::INSERT, 
RecordBatch::RowKind::INSERT,
+                                    RecordBatch::RowKind::UPDATE_AFTER}));
+    ASSERT_OK(writer->Write(std::move(first_batch)));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> update_batch,
+                         MakeBatch({Row{1, "new", "p0"}}, 
/*partitioned=*/false, /*bucket=*/0,
+                                   {RecordBatch::RowKind::UPDATE_AFTER}));
+    ASSERT_OK(writer->Write(std::move(update_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> memory_rows, 
ReadRows(realtime_context));
+    ASSERT_EQ((std::vector<Row>{{1, "new", "p0"}, {2, "two", "p0"}}), 
memory_rows);
+
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_EQ(1, progress.size());
+    ASSERT_OK(Commit(progress, /*commit_identifier=*/0));
+
+    std::vector<Row> second_rows = {{1, "latest", "p0"}, {2, "gone", "p0"}, 
{3, "three", "p0"}};
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> second_batch,
+                         MakeBatch(second_rows, /*partitioned=*/false, 
/*bucket=*/0,
+                                   {RecordBatch::RowKind::UPDATE_AFTER,
+                                    RecordBatch::RowKind::DELETE, 
RecordBatch::RowKind::INSERT}));
+    ASSERT_OK(writer->Write(std::move(second_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> union_rows, 
ReadRows(realtime_context));
+    ASSERT_EQ((std::vector<Row>{{1, "latest", "p0"}, {3, "three", "p0"}}), 
union_rows);
+
+    const std::string expected_payload = "new";
+    std::shared_ptr<Predicate> predicate = PredicateBuilder::Equal(
+        /*field_index=*/1, /*field_name=*/"payload", FieldType::STRING,
+        Literal(FieldType::STRING, expected_payload.data(), 
expected_payload.size()));
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> filtered_plan,
+                         CreatePlan(realtime_context, predicate));
+    ASSERT_OK_AND_ASSIGN(
+        CollectedReadResult filtered_result,
+        ReadPlan(filtered_plan, realtime_context, {"id", "payload", "pt"}, 
predicate,
+                 /*enable_predicate_filter=*/true));
+    ASSERT_EQ(nullptr, filtered_result.data);
+    ASSERT_FALSE(saw_query_predicate->load(std::memory_order_acquire));
+    filtered_result.reader->Close();
+    filtered_result.reader.reset();
+    ASSERT_OK(writer->Close());
+    writer.reset();
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> lifetime_plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ReadContextBuilder read_builder(table_path_);
+    read_builder.SetOptions(options_)
+        .SetReadFieldNames({"id", "payload", "pt"})
+        .WithRealtimeContext(realtime_context)
+        .WithMemoryPool(pool_);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<ReadContext> read_context, 
read_builder.Finish());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableRead> table_read,
+                         TableRead::Create(std::move(read_context)));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<BatchReader> reader,
+                         table_read->CreateReader(lifetime_plan->Splits()));
+    ASSERT_FALSE(query_view->expired());
+
+    std::weak_ptr<RealtimeContext> weak_context = realtime_context;
+    table_read.reset();
+    lifetime_plan.reset();
+    realtime_context.reset();
+    ASSERT_TRUE(weak_context.expired());
+    ASSERT_FALSE(query_view->expired());
+    ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatch read_batch, 
reader->NextBatch());
+    ASSERT_FALSE(BatchReader::IsEofBatch(read_batch));
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::Array> read_array,
+                         ReadResultCollector::GetArray(std::move(read_batch)));
+    ASSERT_NE(nullptr, read_array);
+    read_array.reset();
+    reader->Close();
+    reader.reset();
+    ASSERT_TRUE(query_view->expired());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkRealtimeReadOptimizedScanUnsupported) {
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+    const std::vector<Row> rows = {{1, "one", "p0"}, {2, "two", "p0"}};
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> batch,
+                         MakeBatch(rows, /*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(progress, 
/*commit_identifier=*/0));
+    ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> disk_rows, ReadRows());
+    ASSERT_EQ(rows, disk_rows);
+
+    ScanContextBuilder scan_builder(table_path_ + "$ro");
+    
scan_builder.SetOptions(options_).WithRealtimeContext(realtime_context).WithMemoryPool(pool_);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<ScanContext> scan_context, 
scan_builder.Finish());
+    Result<std::unique_ptr<TableScan>> scan = 
TableScan::Create(std::move(scan_context));
+    ASSERT_TRUE(scan.status().IsNotImplemented()) << scan.status().ToString();
+    ASSERT_NE(std::string::npos, scan.status().ToString().find(
+                                     "PK real-time union read does not support 
read-optimized"));
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkDeleteInsertAndPinnedReadsAcrossRefresh) {
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> delete_batch,
+                         MakeBatch({Row{1, "deleted", "p0"}}, 
/*partitioned=*/false, /*bucket=*/0,
+                                   {RecordBatch::RowKind::DELETE}));
+    ASSERT_OK(writer->Write(std::move(delete_batch)));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> insert_batch,
+                         MakeBatch({Row{1, "inserted", "p0"}}, 
/*partitioned=*/false, /*bucket=*/0,
+                                   {RecordBatch::RowKind::INSERT}));
+    ASSERT_OK(writer->Write(std::move(insert_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_EQ(1, progress.size());
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> pinned_plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> reader_plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ReadContextBuilder read_builder(table_path_);
+    read_builder.SetOptions(options_)
+        .SetReadFieldNames({"id", "payload", "pt"})
+        .WithRealtimeContext(realtime_context)
+        .WithMemoryPool(pool_);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<ReadContext> read_context, 
read_builder.Finish());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableRead> table_read,
+                         TableRead::Create(std::move(read_context)));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<BatchReader> pinned_reader,
+                         table_read->CreateReader(reader_plan->Splits()));
+
+    ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(progress, 
/*commit_identifier=*/0));
+    ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> plan_rows, ReadRows(pinned_plan, 
realtime_context));
+    ASSERT_EQ((std::vector<Row>{{1, "inserted", "p0"}}), plan_rows);
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::ChunkedArray> reader_rows,
+                         
ReadResultCollector::CollectResult(pinned_reader.get()));
+    ASSERT_EQ(1, reader_rows->length());
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkMergeDiskSealedAndActive) {
+    options_[Options::READ_BATCH_SIZE] = "2";
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    const std::vector<std::vector<Row>> disk_batches = {
+        {{1, "disk-1", "p0"}, {2, "disk-2", "p0"}, {3, "disk-3", "p0"}},
+        {{10, "disk-10", "p0"}, {11, "disk-11", "p0"}},
+    };
+    int64_t commit_identifier = 0;
+    for (const std::vector<Row>& disk_rows : disk_batches) {
+        ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> batch,
+                             MakeBatch(disk_rows, /*partitioned=*/false));
+        ASSERT_OK(writer->Write(std::move(batch)));
+        ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                             
writer->PrepareCommitWithProgress(commit_identifier));
+        ASSERT_EQ(1, progress.size());
+        ASSERT_EQ(1, NewFiles(progress).size());
+        ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(progress, 
commit_identifier));
+        ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+        ++commit_identifier;
+    }
+
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> sealed_batch,
+        MakeBatch({Row{1, "sealed-1", "p0"}, Row{2, "deleted-2", "p0"}, Row{4, 
"sealed-4", "p0"}},
+                  /*partitioned=*/false, /*bucket=*/0,
+                  {RecordBatch::RowKind::UPDATE_AFTER, 
RecordBatch::RowKind::DELETE,
+                   RecordBatch::RowKind::INSERT}));
+    ASSERT_OK(writer->Write(std::move(sealed_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> sealed_progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/2));
+    ASSERT_EQ(1, sealed_progress.size());
+
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> active_batch,
+        MakeBatch({Row{1, "active-1", "p0"}, Row{4, "deleted-4", "p0"}, Row{5, 
"active-5", "p0"}},
+                  /*partitioned=*/false, /*bucket=*/0,
+                  {RecordBatch::RowKind::UPDATE_AFTER, 
RecordBatch::RowKind::DELETE,
+                   RecordBatch::RowKind::INSERT}));
+    ASSERT_OK(writer->Write(std::move(active_batch)));
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ASSERT_OK_AND_ASSIGN(CollectedReadResult result,
+                         ReadPlan(plan, realtime_context, {"payload", "id"}, 
/*predicate=*/nullptr,
+                                  /*enable_predicate_filter=*/false));
+    ASSERT_NE(nullptr, result.data);
+    ASSERT_GT(result.data->num_chunks(), 1);
+    for (const std::shared_ptr<arrow::Array>& chunk : result.data->chunks()) {
+        ASSERT_LE(chunk->length(), 2);
+    }
+    std::shared_ptr<arrow::DataType> result_type = arrow::struct_(
+        {arrow::field("_VALUE_KIND", arrow::int8()), arrow::field("payload", 
arrow::utf8()),
+         arrow::field("id", arrow::int64())});
+    std::shared_ptr<arrow::Array> expected =
+        arrow::ipc::internal::json::ArrayFromJSON(result_type, R"([
+            [0, "active-1", 1],
+            [0, "disk-3", 3],
+            [0, "active-5", 5],
+            [0, "disk-10", 10],
+            [0, "disk-11", 11]
+        ])")
+            .ValueOrDie();
+    
ASSERT_TRUE(std::make_shared<arrow::ChunkedArray>(expected)->Equals(*result.data))
+        << result.data->ToString();
+    result.reader->Close();
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkMergeAllDiskSplitsWithMemory) {
+    options_[Options::SOURCE_SPLIT_OPEN_FILE_COST] = "1";
+    options_[Options::SOURCE_SPLIT_TARGET_SIZE] = "1";
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    const std::vector<std::vector<Row>> disk_batches = {
+        {{1, "disk-1", "p0"}, {2, "disk-2", "p0"}},
+        {{10, "disk-10", "p0"}, {11, "disk-11", "p0"}},
+        {{20, "disk-20", "p0"}, {21, "disk-21", "p0"}},
+    };
+    for (int64_t commit_identifier = 0;
+         commit_identifier < static_cast<int64_t>(disk_batches.size()); 
++commit_identifier) {
+        ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> batch,
+                             MakeBatch(disk_batches[commit_identifier], 
/*partitioned=*/false));
+        ASSERT_OK(writer->Write(std::move(batch)));
+        ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                             
writer->PrepareCommitWithProgress(commit_identifier));
+        ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(progress, 
commit_identifier));
+        ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+    }
+
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> memory_batch,
+        MakeBatch({Row{1, "memory-1", "p0"}, Row{10, "deleted-10", "p0"}},
+                  /*partitioned=*/false, /*bucket=*/0,
+                  {RecordBatch::RowKind::UPDATE_AFTER, 
RecordBatch::RowKind::DELETE}));
+    ASSERT_OK(writer->Write(std::move(memory_batch)));
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ASSERT_EQ(1, plan->Splits().size());
+    std::shared_ptr<RealtimeSplit> realtime_split =
+        std::dynamic_pointer_cast<RealtimeSplit>(plan->Splits()[0]);
+    ASSERT_NE(nullptr, realtime_split);
+    ASSERT_EQ(3, realtime_split->DiskSplits().size());
+
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> actual_rows, ReadRows(plan, 
realtime_context));
+    ASSERT_EQ((std::vector<Row>{{1, "memory-1", "p0"},
+                                {2, "disk-2", "p0"},
+                                {11, "disk-11", "p0"},
+                                {20, "disk-20", "p0"},
+                                {21, "disk-21", "p0"}}),
+              actual_rows);
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkNestedProjectionAcrossDiskAndMemory) {
+    const std::shared_ptr<arrow::Field> projected_b = arrow::field("b", 
arrow::int64());
+    fields_ = {
+        arrow::field("id", arrow::int64()),
+        arrow::field("payload", arrow::struct_({arrow::field("a", 
arrow::int64()), projected_b})),
+        arrow::field("pt", arrow::utf8()),
+    };
+    schema_ = arrow::schema(fields_);
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+    auto make_batch = [&](const std::string& json) -> 
Result<std::unique_ptr<RecordBatch>> {
+        PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(
+            std::shared_ptr<arrow::Array> array,
+            arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields_), 
json));
+        ArrowArray c_array;
+        PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportArray(*array, &c_array));
+        RecordBatchBuilder builder(&c_array);
+        return builder.SetBucket(0).Finish();
+    };
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> disk_batch,
+                         make_batch(R"([[1, [101, 1001], "p0"], [2, [102, 
1002], "p0"]])"));
+    ASSERT_OK(writer->Write(std::move(disk_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> disk_progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(disk_progress, 
/*commit_identifier=*/0));
+    ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> sealed_batch,
+                         make_batch(R"([[1, [201, 2001], "p0"], [3, [203, 
2003], "p0"]])"));
+    ASSERT_OK(writer->Write(std::move(sealed_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> sealed_progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/1));
+    ASSERT_EQ(1, sealed_progress.size());
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> active_batch,
+                         make_batch(R"([[1, [301, 3001], "p0"], [4, [304, 
null], "p0"]])"));
+    ASSERT_OK(writer->Write(std::move(active_batch)));
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+
+    auto projected_schema = arrow::schema({
+        arrow::field("payload", arrow::struct_({projected_b})),
+        arrow::field("id", arrow::int64()),
+    });
+    auto c_schema = std::make_unique<ArrowSchema>();
+    ASSERT_TRUE(arrow::ExportSchema(*projected_schema, c_schema.get()).ok());
+    ReadContextBuilder read_builder(table_path_);
+    read_builder.SetOptions(options_)
+        .SetReadSchema(std::move(c_schema))
+        .WithRealtimeContext(realtime_context)
+        .WithMemoryPool(pool_);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<ReadContext> read_context, 
read_builder.Finish());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<TableRead> table_read,
+                         TableRead::Create(std::move(read_context)));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<BatchReader> reader,
+                         table_read->CreateReader(plan->Splits()));
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::ChunkedArray> actual,
+                         ReadResultCollector::CollectResult(reader.get()));
+    const std::shared_ptr<arrow::DataType> result_type = arrow::struct_({
+        arrow::field("_VALUE_KIND", arrow::int8()),
+        arrow::field("payload", arrow::struct_({projected_b})),
+        arrow::field("id", arrow::int64()),
+    });
+    const std::shared_ptr<arrow::Array> expected =
+        arrow::ipc::internal::json::ArrayFromJSON(result_type, R"([
+            [0, [3001], 1],
+            [0, [1002], 2],
+            [0, [2003], 3],
+            [0, [null], 4]
+        ])")
+            .ValueOrDie();
+    
ASSERT_TRUE(std::make_shared<arrow::ChunkedArray>(expected)->Equals(*actual))
+        << actual->ToString();
+    reader->Close();
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkKeylessProjection) {
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> disk_batch,
+                         MakeBatch({Row{1, "disk", "p0"}}, 
/*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(disk_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> disk_progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(disk_progress, 
/*commit_identifier=*/0));
+    ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> memory_batch,
+                         MakeBatch({Row{1, "memory", "p0"}}, 
/*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(memory_batch)));
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ReadPlanWithSchemaAndCheck(plan, realtime_context,
+                               arrow::schema({arrow::field("payload", 
arrow::utf8())}), R"([
+        [0, "memory"]
+    ])");
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestCompositePkKeylessProjection) {
+    CreatePkTable(/*partition_keys=*/{}, /*primary_keys=*/{"id", "payload"});
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> disk_batch,
+                         MakeBatch({Row{1, "key", "disk"}}, 
/*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(disk_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> disk_progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(disk_progress, 
/*commit_identifier=*/0));
+    ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> memory_batch,
+                         MakeBatch({Row{1, "key", "memory"}}, 
/*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(memory_batch)));
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ReadPlanWithSchemaAndCheck(plan, realtime_context,
+                               arrow::schema({arrow::field("pt", 
arrow::utf8())}), R"([
+        [0, "memory"]
+    ])");
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkCompositeMerge) {
+    CreatePkTable(/*partition_keys=*/{}, /*primary_keys=*/{"id", "payload"});
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> disk_batch,
+                         MakeBatch({Row{1, "a", "disk-1a"}, Row{1, "b", 
"disk-1b"},
+                                    Row{2, "a", "disk-2a"}, Row{3, "c", 
"disk-3c"}},
+                                   /*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(disk_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> disk_progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_EQ(1, disk_progress.size());
+    ASSERT_EQ(OffsetRange(0, 4), disk_progress[0].offset_range);
+    ASSERT_EQ(1, NewFiles(disk_progress).size());
+    ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(disk_progress, 
/*commit_identifier=*/0));
+    ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> sealed_batch,
+        MakeBatch({Row{1, "a", "sealed-1a"}, Row{1, "b", "deleted-1b"}, Row{2, 
"b", "sealed-2b"}},
+                  /*partitioned=*/false, /*bucket=*/0,
+                  {RecordBatch::RowKind::UPDATE_AFTER, 
RecordBatch::RowKind::DELETE,
+                   RecordBatch::RowKind::INSERT}));
+    ASSERT_OK(writer->Write(std::move(sealed_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> sealed_progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/1));
+    ASSERT_EQ(1, sealed_progress.size());
+    ASSERT_EQ(OffsetRange(4, 7), sealed_progress[0].offset_range);
+    ASSERT_EQ(1, NewFiles(sealed_progress).size());
+
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> active_batch,
+        MakeBatch({Row{1, "a", "active-1a"}, Row{1, "c", "active-1c"}, Row{2, 
"a", "active-2a"}},
+                  /*partitioned=*/false, /*bucket=*/0,
+                  {RecordBatch::RowKind::UPDATE_AFTER, 
RecordBatch::RowKind::INSERT,
+                   RecordBatch::RowKind::UPDATE_AFTER}));
+    ASSERT_OK(writer->Write(std::move(active_batch)));
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> plan,
+                         CreatePlan(realtime_context, /*predicate=*/nullptr));
+    ASSERT_EQ(1, plan->Splits().size());
+    std::shared_ptr<RealtimeSplit> split =
+        std::dynamic_pointer_cast<RealtimeSplit>(plan->Splits()[0]);
+    ASSERT_NE(nullptr, split);
+    ASSERT_FALSE(split->DiskSplits().empty());
+    ASSERT_EQ(4, split->CommittedEndOffset());
+    ASSERT_EQ(10, split->MemoryEndOffset());
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> actual_rows, ReadRows(plan, 
realtime_context));
+    ASSERT_EQ((std::vector<Row>{{1, "a", "active-1a"},
+                                {1, "c", "active-1c"},
+                                {2, "a", "active-2a"},
+                                {2, "b", "sealed-2b"},
+                                {3, "c", "disk-3c"}}),
+              actual_rows);
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkWriterHandoff) {
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> first_writer,
+                         CreateRealtimeWriter(realtime_context));
+    const std::vector<Row> first_rows = {
+        {0, "value-0", "p0"}, {1, "value-1", "p0"}, {2, "value-2", "p0"}};
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> first_batch,
+                         MakeBatch(first_rows, /*partitioned=*/false));
+    ASSERT_OK(first_writer->Write(std::move(first_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> first_progress,
+                         
first_writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_EQ(1, first_progress.size());
+    ASSERT_EQ(OffsetRange(0, 3), first_progress[0].offset_range);
+    ASSERT_EQ(1, NewFiles(first_progress).size());
+    ASSERT_EQ(0, NewFiles(first_progress)[0]->min_sequence_number);
+    ASSERT_EQ(2, NewFiles(first_progress)[0]->max_sequence_number);
+    ASSERT_OK(first_writer->Close());
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> second_writer,
+                         CreateRealtimeWriter(realtime_context));
+    const std::vector<Row> second_rows = {{0, "updated-0", "p0"}, {3, 
"value-3", "p0"}};
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> second_batch,
+                         MakeBatch(second_rows, /*partitioned=*/false));
+    ASSERT_OK(second_writer->Write(std::move(second_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> second_progress,
+                         
second_writer->PrepareCommitWithProgress(/*commit_identifier=*/1));
+    ASSERT_EQ(1, second_progress.size());
+    ASSERT_EQ(OffsetRange(3, 5), second_progress[0].offset_range);
+    ASSERT_EQ(1, NewFiles(second_progress).size());
+    ASSERT_EQ(3, NewFiles(second_progress)[0]->min_sequence_number);
+    ASSERT_EQ(4, NewFiles(second_progress)[0]->max_sequence_number);
+
+    first_progress.push_back(std::move(second_progress[0]));
+    ASSERT_OK(Commit(first_progress, /*commit_identifier=*/1));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> actual_rows, 
ReadRows(realtime_context));
+    ASSERT_EQ((std::vector<Row>{{0, "updated-0", "p0"},
+                                {1, "value-1", "p0"},
+                                {2, "value-2", "p0"},
+                                {3, "value-3", "p0"}}),
+              actual_rows);
+    ASSERT_OK(second_writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkPartitionBucketRecovery) {
+    options_[Options::BUCKET] = "2";
+    CreatePkTable(/*partition_keys=*/{"pt"});
+    const RealtimePartitionBucket p0b0({{"pt", "p0"}}, /*bucket=*/0);
+    const RealtimePartitionBucket p1b1({{"pt", "p1"}}, /*bucket=*/1);
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> first_context, 
RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> first_writer,
+                         CreateRealtimeWriter(first_context));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> p0_first_batch,
+                         MakeBatch({Row{0, "p0-zero", "p0"}, Row{1, "p0-one", 
"p0"}},
+                                   /*partitioned=*/true, /*bucket=*/0));
+    ASSERT_OK(first_writer->Write(std::move(p0_first_batch)));
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> p1_first_batch,
+        MakeBatch({Row{10, "p1-ten", "p1"}, Row{11, "p1-eleven", "p1"}, 
Row{12, "p1-twelve", "p1"}},
+                  /*partitioned=*/true, /*bucket=*/1));
+    ASSERT_OK(first_writer->Write(std::move(p1_first_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> first_progress,
+                         
first_writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_EQ(2, first_progress.size());
+    std::map<RealtimePartitionBucket, OffsetRange> first_ranges;
+    std::map<RealtimePartitionBucket, std::pair<int64_t, int64_t>> 
first_sequences;
+    for (const RealtimeCommitProgress& progress : first_progress) {
+        first_ranges.emplace(progress.partition_bucket, progress.offset_range);
+        std::shared_ptr<CommitMessageImpl> message =
+            
std::dynamic_pointer_cast<CommitMessageImpl>(progress.commit_message);
+        ASSERT_NE(nullptr, message);
+        const std::vector<std::shared_ptr<DataFileMeta>>& files =
+            message->GetNewFilesIncrement().NewFiles();
+        ASSERT_EQ(1, files.size());
+        first_sequences.emplace(
+            progress.partition_bucket,
+            std::make_pair(files[0]->min_sequence_number, 
files[0]->max_sequence_number));
+    }
+    ASSERT_EQ(OffsetRange(0, 2), first_ranges.at(p0b0));
+    ASSERT_EQ(OffsetRange(0, 3), first_ranges.at(p1b1));
+    ASSERT_EQ((std::pair<int64_t, int64_t>(0, 1)), first_sequences.at(p0b0));
+    ASSERT_EQ((std::pair<int64_t, int64_t>(0, 2)), first_sequences.at(p1b1));
+    ASSERT_OK_AND_ASSIGN(int64_t first_snapshot_id,
+                         Commit(first_progress, /*commit_identifier=*/0));
+    ASSERT_OK(first_writer->RefreshCommittedSnapshot(first_snapshot_id));
+    ASSERT_OK(first_writer->Close());
+    first_writer.reset();
+    first_context.reset();
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> second_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> second_writer,
+                         CreateRealtimeWriter(second_context));
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> p0_second_batch,
+        MakeBatch({Row{0, "p0-zero-new", "p0"}, Row{2, "p0-two", "p0"}},
+                  /*partitioned=*/true, /*bucket=*/0,
+                  {RecordBatch::RowKind::UPDATE_AFTER, 
RecordBatch::RowKind::INSERT}));
+    ASSERT_OK(second_writer->Write(std::move(p0_second_batch)));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> p1_second_batch,
+                         MakeBatch({Row{10, "p1-ten-deleted", "p1"}, Row{13, 
"p1-thirteen", "p1"}},
+                                   /*partitioned=*/true, /*bucket=*/1,
+                                   {RecordBatch::RowKind::DELETE, 
RecordBatch::RowKind::INSERT}));
+    ASSERT_OK(second_writer->Write(std::move(p1_second_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> second_progress,
+                         
second_writer->PrepareCommitWithProgress(/*commit_identifier=*/1));
+    ASSERT_EQ(2, second_progress.size());
+    std::map<RealtimePartitionBucket, OffsetRange> second_ranges;
+    std::map<RealtimePartitionBucket, std::pair<int64_t, int64_t>> 
second_sequences;
+    for (const RealtimeCommitProgress& progress : second_progress) {
+        second_ranges.emplace(progress.partition_bucket, 
progress.offset_range);
+        std::shared_ptr<CommitMessageImpl> message =
+            
std::dynamic_pointer_cast<CommitMessageImpl>(progress.commit_message);
+        ASSERT_NE(nullptr, message);
+        const std::vector<std::shared_ptr<DataFileMeta>>& files =
+            message->GetNewFilesIncrement().NewFiles();
+        ASSERT_EQ(1, files.size());
+        second_sequences.emplace(
+            progress.partition_bucket,
+            std::make_pair(files[0]->min_sequence_number, 
files[0]->max_sequence_number));
+    }
+    ASSERT_EQ(OffsetRange(2, 4), second_ranges.at(p0b0));
+    ASSERT_EQ(OffsetRange(3, 5), second_ranges.at(p1b1));
+    ASSERT_EQ((std::pair<int64_t, int64_t>(2, 3)), second_sequences.at(p0b0));
+    ASSERT_EQ((std::pair<int64_t, int64_t>(3, 4)), second_sequences.at(p1b1));
+    ASSERT_OK_AND_ASSIGN(int64_t second_snapshot_id,
+                         Commit(second_progress, /*commit_identifier=*/1));
+    ASSERT_OK(second_writer->RefreshCommittedSnapshot(second_snapshot_id));
+
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> actual_rows, 
ReadRows(second_context));
+    std::sort(actual_rows.begin(), actual_rows.end());
+    ASSERT_EQ((std::vector<Row>{{0, "p0-zero-new", "p0"},
+                                {1, "p0-one", "p0"},
+                                {2, "p0-two", "p0"},
+                                {11, "p1-eleven", "p1"},
+                                {12, "p1-twelve", "p1"},
+                                {13, "p1-thirteen", "p1"}}),
+              actual_rows);
+    ASSERT_OK(second_writer->Close());
+
+    ASSERT_OK_AND_ASSIGN(RealtimeOffsetMap offsets, ReadCommittedOffsets());
+    ASSERT_EQ(2, offsets.size());
+    ASSERT_EQ(4, offsets.at(p0b0));
+    ASSERT_EQ(5, offsets.at(p1b1));
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkRecovery) {
+    CreatePkTable();
+
+    WriteContextBuilder seed_builder(table_path_, commit_user_);
+    seed_builder.SetOptions(options_).WithStreamingMode(true);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<WriteContext> seed_context, 
seed_builder.Finish());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> seed_writer,
+                         FileStoreWrite::Create(std::move(seed_context)));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> seed_batch,
+                         MakeBatch({Row{99, "seed", "p0"}}, 
/*partitioned=*/false));
+    ASSERT_OK(seed_writer->Write(std::move(seed_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<std::shared_ptr<CommitMessage>> 
seed_messages,
+                         seed_writer->PrepareCommit(/*wait_compaction=*/false,
+                                                    /*commit_identifier=*/0));
+    CommitContextBuilder seed_commit_builder(table_path_, commit_user_);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<CommitContext> seed_commit_context,
+                         seed_commit_builder.SetOptions(options_).Finish());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> seed_commit,
+                         
FileStoreCommit::Create(std::move(seed_commit_context)));
+    ASSERT_OK(seed_commit->Commit(seed_messages, /*commit_identifier=*/0));
+    ASSERT_OK(seed_writer->Close());
+    const std::vector<Row> mutations = {
+        {1, "one", "p0"}, {1, "one-new", "p0"}, {2, "deleted", "p0"}, {3, 
"three", "p0"}};
+    const std::vector<RecordBatch::RowKind> mutation_kinds = {
+        RecordBatch::RowKind::INSERT, RecordBatch::RowKind::UPDATE_AFTER,
+        RecordBatch::RowKind::DELETE, RecordBatch::RowKind::INSERT};
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> first_context, 
RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> first_writer,
+                         CreateRealtimeWriter(first_context));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> batch,
+                         MakeBatch(mutations, /*partitioned=*/false, 
/*bucket=*/0, mutation_kinds));
+    ASSERT_OK(first_writer->Write(std::move(batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<int64_t> memory_sequences, 
ReadPkSequences(first_context));
+    ASSERT_EQ((std::vector<int64_t>{1, 2, 3, 4}), memory_sequences);
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                         
first_writer->PrepareCommitWithProgress(/*commit_identifier=*/1));
+    ASSERT_EQ(1, progress.size());
+    ASSERT_EQ(OffsetRange(0, 4), progress[0].offset_range);
+    ASSERT_EQ(1, NewFiles(progress).size());
+    ASSERT_EQ(2, NewFiles(progress)[0]->min_sequence_number);
+    ASSERT_EQ(memory_sequences.back(), 
NewFiles(progress)[0]->max_sequence_number);
+    ASSERT_EQ(1, NewFiles(progress)[0]->delete_row_count);
+    ASSERT_OK(Commit(progress, /*commit_identifier=*/1));
+    ASSERT_OK(first_writer->Close());
+    first_context.reset();
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> rows_after_replay, ReadRows());
+    ASSERT_EQ((std::vector<Row>{{1, "one-new", "p0"}, {3, "three", "p0"}, {99, 
"seed", "p0"}}),
+              rows_after_replay);
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> second_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> second_writer,
+                         CreateRealtimeWriter(second_context));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> restart_batch,
+                         MakeBatch({Row{4, "four", "p0"}}, 
/*partitioned=*/false));
+    ASSERT_OK(second_writer->Write(std::move(restart_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<int64_t> restart_sequences, 
ReadPkSequences(second_context));
+    ASSERT_EQ((std::vector<int64_t>{5}), restart_sequences);
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> restart_progress,
+                         
second_writer->PrepareCommitWithProgress(/*commit_identifier=*/2));
+    ASSERT_EQ(1, restart_progress.size());
+    ASSERT_EQ(OffsetRange(4, 5), restart_progress[0].offset_range);
+    ASSERT_EQ(5, NewFiles(restart_progress)[0]->min_sequence_number);
+    ASSERT_EQ(5, NewFiles(restart_progress)[0]->max_sequence_number);
+    ASSERT_OK(second_writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkCompaction) {
+    options_[Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER] = "1";
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+
+    int64_t latest_snapshot_id = -1;
+    constexpr int64_t kCommitRoundsBeforeCompaction = 4;
+    std::set<std::string> committed_file_names;
+    for (int64_t round = 0; round < kCommitRoundsBeforeCompaction; ++round) {
+        const bool delete_latest_live_row = round == 
kCommitRoundsBeforeCompaction - 1;
+        ASSERT_OK_AND_ASSIGN(
+            std::unique_ptr<RecordBatch> batch,
+            MakeBatch(
+                {Row{delete_latest_live_row ? round - 1 : round,
+                     delete_latest_live_row ? "deleted" : "value-" + 
std::to_string(round), "p0"}},
+                /*partitioned=*/false, /*bucket=*/0,
+                delete_latest_live_row
+                    ? 
std::vector<RecordBatch::RowKind>{RecordBatch::RowKind::DELETE}
+                    : std::vector<RecordBatch::RowKind>{}));
+        ASSERT_OK(writer->Write(std::move(batch)));
+        ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                             writer->PrepareCommitWithProgress(round));
+        ASSERT_EQ(1, progress.size());
+        std::shared_ptr<CommitMessageImpl> message =
+            
std::dynamic_pointer_cast<CommitMessageImpl>(progress[0].commit_message);
+        ASSERT_NE(nullptr, message);
+        ASSERT_TRUE(message->GetCompactIncrement().IsEmpty());
+        ASSERT_EQ(1, NewFiles(progress).size());
+        committed_file_names.insert(NewFiles(progress)[0]->file_name);
+        ASSERT_OK_AND_ASSIGN(latest_snapshot_id, Commit(progress, round));
+        ASSERT_OK(writer->RefreshCommittedSnapshot(latest_snapshot_id));
+        ASSERT_OK_AND_ASSIGN(uint64_t memory_usage, 
GetRealtimeMemoryUsage(realtime_context));
+        ASSERT_EQ(0, memory_usage);
+    }
+    WriteContextBuilder compact_builder(table_path_, commit_user_);
+    compact_builder.SetOptions(options_).WithStreamingMode(true);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<WriteContext> compact_context, 
compact_builder.Finish());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> compact_writer,
+                         FileStoreWrite::Create(std::move(compact_context)));
+    ASSERT_OK(compact_writer->Compact(/*partition=*/{}, /*bucket=*/0,
+                                      /*full_compaction=*/true));
+    ASSERT_OK_AND_ASSIGN(
+        std::vector<std::shared_ptr<CommitMessage>> compact_messages,
+        compact_writer->PrepareCommit(/*wait_compaction=*/true, 
/*commit_identifier=*/4));
+    ASSERT_EQ(1, compact_messages.size());
+    std::shared_ptr<CommitMessageImpl> compact_message =
+        std::dynamic_pointer_cast<CommitMessageImpl>(compact_messages[0]);
+    ASSERT_NE(nullptr, compact_message);
+    ASSERT_TRUE(compact_message->GetNewFilesIncrement().IsEmpty());
+    ASSERT_EQ(kCommitRoundsBeforeCompaction,
+              compact_message->GetCompactIncrement().CompactBefore().size());
+    std::set<std::string> compacted_file_names;
+    for (const std::shared_ptr<DataFileMeta>& file :
+         compact_message->GetCompactIncrement().CompactBefore()) {
+        compacted_file_names.insert(file->file_name);
+    }
+    ASSERT_EQ(committed_file_names, compacted_file_names);
+    
ASSERT_FALSE(compact_message->GetCompactIncrement().CompactAfter().empty());
+    constexpr int64_t kHistoricalMaxSequenceNumber = 
kCommitRoundsBeforeCompaction - 1;
+    int64_t compacted_live_max_sequence_number = -1;
+    for (const std::shared_ptr<DataFileMeta>& file :
+         compact_message->GetCompactIncrement().CompactAfter()) {
+        compacted_live_max_sequence_number =
+            std::max(compacted_live_max_sequence_number, 
file->max_sequence_number);
+    }
+    ASSERT_LT(compacted_live_max_sequence_number, 
kHistoricalMaxSequenceNumber);
+    CommitContextBuilder commit_builder(table_path_, commit_user_);
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<CommitContext> commit_context,
+                         commit_builder.SetOptions(options_).Finish());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> commit,
+                         FileStoreCommit::Create(std::move(commit_context)));
+    ASSERT_OK(commit->Commit(compact_messages, /*commit_identifier=*/4));
+    ASSERT_OK(compact_writer->Close());
+
+    ASSERT_OK_AND_ASSIGN(CoreOptions options, CoreOptions::FromMap(options_));
+    SnapshotManager snapshot_manager(options.GetFileSystem(), table_path_);
+    ASSERT_OK_AND_ASSIGN(std::optional<Snapshot> compact_snapshot,
+                         snapshot_manager.LatestSnapshot());
+    ASSERT_TRUE(compact_snapshot);
+    ASSERT_EQ(Snapshot::CommitKind::Compact(), 
compact_snapshot->GetCommitKind());
+    ASSERT_OK_AND_ASSIGN(RealtimeOffsetMap offsets, ReadCommittedOffsets());
+    ASSERT_EQ(4, offsets.at(RealtimePartitionBucket(/*partition=*/{}, 
/*bucket=*/0)));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> compacted_rows, ReadRows());
+    ASSERT_EQ((std::vector<Row>{{0, "value-0", "p0"}, {1, "value-1", "p0"}}), 
compacted_rows);
+    ASSERT_OK(writer->Close());
+    writer.reset();
+    realtime_context.reset();
+
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> fresh_context, 
RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> fresh_writer,
+                         CreateRealtimeWriter(fresh_context));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> fresh_batch,
+                         MakeBatch({Row{4, "value-4", "p0"}},
+                                   /*partitioned=*/false));
+    ASSERT_OK(fresh_writer->Write(std::move(fresh_batch)));
+    ASSERT_OK_AND_ASSIGN(std::vector<int64_t> fresh_sequences, 
ReadPkSequences(fresh_context));
+    ASSERT_EQ((std::vector<int64_t>{compacted_live_max_sequence_number + 1}), 
fresh_sequences);
+    ASSERT_LT(fresh_sequences.front(), kHistoricalMaxSequenceNumber);
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> fresh_progress,
+                         
fresh_writer->PrepareCommitWithProgress(/*commit_identifier=*/5));
+    ASSERT_EQ(1, fresh_progress.size());
+    ASSERT_EQ(OffsetRange(4, 5), fresh_progress[0].offset_range);
+    ASSERT_EQ(compacted_live_max_sequence_number + 1,
+              NewFiles(fresh_progress)[0]->min_sequence_number);
+    ASSERT_EQ(compacted_live_max_sequence_number + 1,
+              NewFiles(fresh_progress)[0]->max_sequence_number);
+    ASSERT_OK_AND_ASSIGN(latest_snapshot_id, Commit(fresh_progress, 
/*commit_identifier=*/5));
+    ASSERT_OK(fresh_writer->Close());
+
+    ASSERT_OK_AND_ASSIGN(offsets, ReadCommittedOffsets());
+    ASSERT_EQ(5, offsets.at(RealtimePartitionBucket(/*partition=*/{}, 
/*bucket=*/0)));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> final_rows, ReadRows());
+    ASSERT_EQ((std::vector<Row>{{0, "value-0", "p0"}, {1, "value-1", "p0"}, 
{4, "value-4", "p0"}}),
+              final_rows);
+}
+
+TEST_F(RealtimeWriteInteTest, 
TestPkMultipleStoredBatchesMergeForQueryAndCommit) {
+    CreatePkTable();
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create());
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> first_batch,
+                         MakeBatch({Row{4, "four", "p0"}, Row{2, "two", "p0"}, 
Row{1, "one", "p0"}},
+                                   /*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(first_batch)));
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<RecordBatch> second_batch,
+        MakeBatch({Row{3, "three", "p0"}, Row{2, "deleted", "p0"}, Row{1, 
"one-new", "p0"}},
+                  /*partitioned=*/false, /*bucket=*/0,
+                  {RecordBatch::RowKind::INSERT, RecordBatch::RowKind::DELETE,
+                   RecordBatch::RowKind::UPDATE_AFTER}));
+    ASSERT_OK(writer->Write(std::move(second_batch)));
+
+    const std::vector<Row> expected = {{1, "one-new", "p0"}, {3, "three", 
"p0"}, {4, "four", "p0"}};
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> query_rows, 
ReadRows(realtime_context));
+    ASSERT_EQ(expected, query_rows);
+
+    ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+                         
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+    ASSERT_EQ(1, progress.size());
+    ASSERT_EQ(OffsetRange(0, 6), progress[0].offset_range);
+    ASSERT_OK(Commit(progress, /*commit_identifier=*/0));
+    ASSERT_OK_AND_ASSIGN(std::vector<Row> rows, ReadRows());
+    ASSERT_EQ(expected, rows);
+    ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestPkQueryReaderClose) {
+    CreatePkTable();
+    auto state = std::make_shared<CloseTrackingReaderState>();
+    auto factory = MakeDecoratingFactory<CloseTrackingRealtimeStore>(state);
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+                         RealtimeContext::Create(factory));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+                         CreateRealtimeWriter(realtime_context));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> batch,
+                         MakeBatch({Row{1, "one", "p0"}}, 
/*partitioned=*/false));
+    ASSERT_OK(writer->Write(std::move(batch)));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<BatchReader> reader, 
CreateQueryReader(realtime_context));
+    reader->Close();
+    ASSERT_EQ(1, state->query_close_count->load(std::memory_order_acquire));
+    ASSERT_OK(writer->Close());

Review Comment:
   Could this behavior be covered in a focused reader adapter unit test 
instead? Verifying that `Close()` propagates to the underlying reader seems 
valuable, but testing it here requires several store/factory wrappers and adds 
considerable integration-test scaffolding.



-- 
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]

Reply via email to