zjw1111 commented on code in PR #224:
URL: https://github.com/apache/paimon-cpp/pull/224#discussion_r3879949267
##########
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:
In my view, this test can be removed entirely. The components in this reader
chain and most of their `Close()` propagation logic already existed in the
framework before this PR, so they should not be tested specifically as part of
the realtime PK implementation.
--
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]