Skip to content

Commit 8a67374

Browse files
committed
test(spill): add unit tests and integration tests for spill-to-disk)
1 parent 56bb617 commit 8a67374

5 files changed

Lines changed: 1693 additions & 33 deletions

File tree

src/paimon/core/mergetree/merge_tree_writer_test.cpp

Lines changed: 353 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@
5252
#include "paimon/testing/utils/binary_row_generator.h"
5353
#include "paimon/testing/utils/io_exception_helper.h"
5454
#include "paimon/testing/utils/read_result_collector.h"
55+
#include "paimon/testing/utils/test_helper.h"
5556
#include "paimon/testing/utils/testharness.h"
5657

5758
namespace paimon {
@@ -1075,6 +1076,358 @@ TEST_P(MergeTreeWriterTest, TestCloseSkipsDeleteForUpgradedFilesInCompactAfter)
10751076
<< "Intermediate file should be deleted because it's not in compact_before_";
10761077
}
10771078

1079+
TEST_P(MergeTreeWriterTest, TestSpillWithSameKeyDeduplicate) {
1080+
ASSERT_OK_AND_ASSIGN(CoreOptions options,
1081+
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
1082+
{Options::WRITE_BUFFER_SIZE, "1"},
1083+
{Options::WRITE_ONLY, "true"}}));
1084+
auto dir = UniqueTestDirectory::Create();
1085+
ASSERT_TRUE(dir);
1086+
auto path_factory = std::make_shared<DataFilePathFactory>();
1087+
ASSERT_OK(path_factory->Init(dir->Str(), "orc", options.DataFilePrefix(), nullptr));
1088+
std::string uuid = path_factory->uuid_;
1089+
1090+
std::shared_ptr<IOManager> io_manager =
1091+
std::make_shared<IOManager>(dir->Str() + "/tmp", file_system_);
1092+
ASSERT_OK_AND_ASSIGN(
1093+
auto merge_writer,
1094+
MergeTreeWriter::Create(/*last_sequence_number=*/-1, primary_keys_, path_factory,
1095+
key_comparator_, /*user_defined_seq_comparator=*/nullptr,
1096+
merge_function_wrapper_, /*schema_id=*/0, value_schema_, options,
1097+
noop_compact_manager_, io_manager, pool_));
1098+
1099+
std::shared_ptr<arrow::Array> batch1 =
1100+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1101+
["Alice", 1, 0, 1.0],
1102+
["Bob", 2, 0, 2.0]
1103+
])")
1104+
.ValueOrDie();
1105+
std::shared_ptr<arrow::Array> batch2 =
1106+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1107+
["Alice", 10, 0, 10.0],
1108+
["Charlie", 3, 0, 3.0]
1109+
])")
1110+
.ValueOrDie();
1111+
1112+
WriteBatch(batch1, /*row_kinds=*/{}, merge_writer.get());
1113+
WriteBatch(batch2, /*row_kinds=*/{}, merge_writer.get());
1114+
ASSERT_EQ(2u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1115+
1116+
std::shared_ptr<arrow::Array> batch3 =
1117+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1118+
["Bob", 20, 0, 20.0],
1119+
["Charlie", 30, 0, 30.0]
1120+
])")
1121+
.ValueOrDie();
1122+
WriteBatch(batch3, /*row_kinds=*/{}, merge_writer.get());
1123+
1124+
ASSERT_OK_AND_ASSIGN(CommitIncrement commit_increment,
1125+
merge_writer->PrepareCommit(/*wait_compaction=*/false));
1126+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1127+
ASSERT_OK(merge_writer->Close());
1128+
1129+
// All three keys deduplicated: Alice(seq=2), Bob(seq=4), Charlie(seq=5).
1130+
ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size());
1131+
std::string expected_data_file_path = dir->Str() + "/data-" + uuid + "-0.orc";
1132+
std::shared_ptr<arrow::ChunkedArray> expected_array;
1133+
auto array_status = arrow::ipc::internal::json::ChunkedArrayFromJSON(write_type_, {R"([
1134+
[2, 0, "Alice", 10, 0, 10.0],
1135+
[4, 0, "Bob", 20, 0, 20.0],
1136+
[5, 0, "Charlie", 30, 0, 30.0]
1137+
])"},
1138+
&expected_array);
1139+
ASSERT_TRUE(array_status.ok());
1140+
CheckFileContent(expected_data_file_path, expected_array);
1141+
}
1142+
1143+
TEST_P(MergeTreeWriterTest, TestIntermediateMergeSpillFileBound) {
1144+
ASSERT_OK_AND_ASSIGN(CoreOptions options,
1145+
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
1146+
{Options::WRITE_BUFFER_SIZE, "1"},
1147+
{Options::LOCAL_SORT_MAX_NUM_FILE_HANDLES, "2"},
1148+
{Options::WRITE_ONLY, "true"}}));
1149+
auto dir = UniqueTestDirectory::Create();
1150+
ASSERT_TRUE(dir);
1151+
auto path_factory = std::make_shared<DataFilePathFactory>();
1152+
ASSERT_OK(path_factory->Init(dir->Str(), "orc", options.DataFilePrefix(), nullptr));
1153+
std::string uuid = path_factory->uuid_;
1154+
1155+
std::shared_ptr<IOManager> io_manager =
1156+
std::make_shared<IOManager>(dir->Str() + "/tmp", file_system_);
1157+
ASSERT_OK_AND_ASSIGN(
1158+
auto merge_writer,
1159+
MergeTreeWriter::Create(/*last_sequence_number=*/-1, primary_keys_, path_factory,
1160+
key_comparator_, /*user_defined_seq_comparator=*/nullptr,
1161+
merge_function_wrapper_, /*schema_id=*/0, value_schema_, options,
1162+
noop_compact_manager_, io_manager, pool_));
1163+
1164+
std::shared_ptr<arrow::Array> batch1 =
1165+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1166+
["Alice", 1, 0, 1.0]
1167+
])")
1168+
.ValueOrDie();
1169+
std::shared_ptr<arrow::Array> batch2 =
1170+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1171+
["Bob", 2, 0, 2.0]
1172+
])")
1173+
.ValueOrDie();
1174+
std::shared_ptr<arrow::Array> batch3 =
1175+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1176+
["Alice", 3, 0, 3.0]
1177+
])")
1178+
.ValueOrDie();
1179+
1180+
WriteBatch(batch1, /*row_kinds=*/{}, merge_writer.get());
1181+
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1182+
1183+
WriteBatch(batch2, /*row_kinds=*/{}, merge_writer.get());
1184+
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1185+
1186+
WriteBatch(batch3, /*row_kinds=*/{}, merge_writer.get());
1187+
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1188+
1189+
ASSERT_OK_AND_ASSIGN(CommitIncrement commit_increment,
1190+
merge_writer->PrepareCommit(/*wait_compaction=*/false));
1191+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1192+
ASSERT_OK(merge_writer->Close());
1193+
1194+
ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size());
1195+
std::string expected_data_file_path = dir->Str() + "/data-" + uuid + "-0.orc";
1196+
std::shared_ptr<arrow::ChunkedArray> expected_array;
1197+
auto array_status = arrow::ipc::internal::json::ChunkedArrayFromJSON(write_type_, {R"([
1198+
[2, 0, "Alice", 3, 0, 3.0],
1199+
[1, 0, "Bob", 2, 0, 2.0]
1200+
])"},
1201+
&expected_array);
1202+
ASSERT_TRUE(array_status.ok());
1203+
CheckFileContent(expected_data_file_path, expected_array);
1204+
}
1205+
1206+
TEST_P(MergeTreeWriterTest, TestDiskQuotaExhaustedFallsBackToFlushWriteBuffer) {
1207+
ASSERT_OK_AND_ASSIGN(CoreOptions options,
1208+
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
1209+
{Options::WRITE_BUFFER_SIZE, "1"},
1210+
{Options::WRITE_BUFFER_SPILL_MAX_DISK_SIZE, "1"},
1211+
{Options::WRITE_ONLY, "true"}}));
1212+
auto dir = UniqueTestDirectory::Create();
1213+
ASSERT_TRUE(dir);
1214+
auto path_factory = std::make_shared<DataFilePathFactory>();
1215+
ASSERT_OK(path_factory->Init(dir->Str(), "orc", options.DataFilePrefix(), nullptr));
1216+
1217+
std::shared_ptr<IOManager> io_manager =
1218+
std::make_shared<IOManager>(dir->Str() + "/tmp", file_system_);
1219+
ASSERT_OK_AND_ASSIGN(
1220+
auto merge_writer,
1221+
MergeTreeWriter::Create(/*last_sequence_number=*/-1, primary_keys_, path_factory,
1222+
key_comparator_, /*user_defined_seq_comparator=*/nullptr,
1223+
merge_function_wrapper_, /*schema_id=*/0, value_schema_, options,
1224+
noop_compact_manager_, io_manager, pool_));
1225+
1226+
// Phase 1: Manual FlushMemory path — disk quota exhausted causes fallback.
1227+
std::shared_ptr<arrow::Array> array1 =
1228+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1229+
["Alice", 1, 0, 1.0],
1230+
["Bob", 2, 0, 2.0],
1231+
["Charlie", 3, 0, 3.0]
1232+
])")
1233+
.ValueOrDie();
1234+
WriteBatch(array1, /*row_kinds=*/{}, merge_writer.get());
1235+
1236+
ASSERT_EQ(merge_writer->GetMemoryUsage(), 0);
1237+
1238+
ASSERT_OK_AND_ASSIGN(CommitIncrement commit1,
1239+
merge_writer->PrepareCommit(/*wait_compaction=*/false));
1240+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1241+
ASSERT_EQ(1, commit1.GetNewFilesIncrement().NewFiles().size());
1242+
ASSERT_EQ(3, commit1.GetNewFilesIncrement().NewFiles()[0]->row_count);
1243+
1244+
// Phase 2: Auto-spill path — WRITE_BUFFER_SIZE=1 triggers spill on each WriteBatch.
1245+
// batch1 spills successfully, but disk quota is now exhausted.
1246+
std::shared_ptr<arrow::Array> batch1 =
1247+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1248+
["Dave", 4, 0, 4.0]
1249+
])")
1250+
.ValueOrDie();
1251+
WriteBatch(batch1, /*row_kinds=*/{}, merge_writer.get());
1252+
1253+
// batch2: spill -> quota exhausted -> FlushWriteBuffer produces a data file.
1254+
std::shared_ptr<arrow::Array> batch2 =
1255+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1256+
["Eve", 5, 0, 5.0]
1257+
])")
1258+
.ValueOrDie();
1259+
WriteBatch(batch2, /*row_kinds=*/{}, merge_writer.get());
1260+
1261+
// batch3: another round after flush, accumulates into a fresh buffer.
1262+
std::shared_ptr<arrow::Array> batch3 =
1263+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1264+
["Frank", 6, 0, 6.0]
1265+
])")
1266+
.ValueOrDie();
1267+
WriteBatch(batch3, /*row_kinds=*/{}, merge_writer.get());
1268+
1269+
ASSERT_OK_AND_ASSIGN(CommitIncrement commit2,
1270+
merge_writer->PrepareCommit(/*wait_compaction=*/false));
1271+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1272+
ASSERT_OK(merge_writer->Close());
1273+
1274+
ASSERT_EQ(3, commit2.GetNewFilesIncrement().NewFiles().size());
1275+
for (const auto& file_meta : commit2.GetNewFilesIncrement().NewFiles()) {
1276+
ASSERT_EQ(1, file_meta->row_count);
1277+
}
1278+
}
1279+
1280+
TEST_P(MergeTreeWriterTest, TestFlushMemoryQuotaExhaustedFallsBackToFlushWriteBuffer) {
1281+
// WRITE_BUFFER_SIZE is large enough so WriteBatch does NOT auto-spill.
1282+
// SPILL_MAX_DISK_SIZE is tiny so the first FlushMemory() exhausts the quota,
1283+
// triggering the fallback path: FlushMemory() -> quota exhausted -> FlushWriteBuffer.
1284+
ASSERT_OK_AND_ASSIGN(CoreOptions options,
1285+
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
1286+
{Options::WRITE_BUFFER_SIZE, "4096000"},
1287+
{Options::WRITE_BUFFER_SPILL_MAX_DISK_SIZE, "1b"},
1288+
{Options::WRITE_ONLY, "true"}}));
1289+
auto dir = UniqueTestDirectory::Create();
1290+
ASSERT_TRUE(dir);
1291+
auto path_factory = std::make_shared<DataFilePathFactory>();
1292+
ASSERT_OK(path_factory->Init(dir->Str(), "orc", options.DataFilePrefix(), nullptr));
1293+
1294+
std::shared_ptr<IOManager> io_manager =
1295+
std::make_shared<IOManager>(dir->Str() + "/tmp", file_system_);
1296+
ASSERT_OK_AND_ASSIGN(
1297+
auto merge_writer,
1298+
MergeTreeWriter::Create(/*last_sequence_number=*/-1, primary_keys_, path_factory,
1299+
key_comparator_, /*user_defined_seq_comparator=*/nullptr,
1300+
merge_function_wrapper_, /*schema_id=*/0, value_schema_, options,
1301+
noop_compact_manager_, io_manager, pool_));
1302+
1303+
std::shared_ptr<arrow::Array> array =
1304+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1305+
["Alice", 1, 0, 1.0],
1306+
["Bob", 2, 0, 2.0]
1307+
])")
1308+
.ValueOrDie();
1309+
WriteBatch(array, /*row_kinds=*/{}, merge_writer.get());
1310+
1311+
ASSERT_GT(merge_writer->GetMemoryUsage(), 0);
1312+
1313+
// FlushMemory: spill succeeds but disk quota is exhausted -> falls back to FlushWriteBuffer.
1314+
ASSERT_OK(merge_writer->FlushMemory());
1315+
ASSERT_EQ(merge_writer->GetMemoryUsage(), 0);
1316+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1317+
1318+
// PrepareCommit should produce a data file (from FlushWriteBuffer fallback).
1319+
ASSERT_OK_AND_ASSIGN(CommitIncrement commit,
1320+
merge_writer->PrepareCommit(/*wait_compaction=*/false));
1321+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1322+
ASSERT_EQ(1, commit.GetNewFilesIncrement().NewFiles().size());
1323+
ASSERT_EQ(2, commit.GetNewFilesIncrement().NewFiles()[0]->row_count);
1324+
ASSERT_OK(merge_writer->Close());
1325+
}
1326+
1327+
TEST_P(MergeTreeWriterTest, TestCloseDeletesSpillTempFiles) {
1328+
ASSERT_OK_AND_ASSIGN(CoreOptions options,
1329+
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
1330+
{Options::WRITE_BUFFER_SIZE, "1"},
1331+
{Options::WRITE_ONLY, "true"}}));
1332+
auto dir = UniqueTestDirectory::Create();
1333+
ASSERT_TRUE(dir);
1334+
auto path_factory = std::make_shared<DataFilePathFactory>();
1335+
ASSERT_OK(path_factory->Init(dir->Str(), "orc", options.DataFilePrefix(), nullptr));
1336+
1337+
std::shared_ptr<IOManager> io_manager =
1338+
std::make_shared<IOManager>(dir->Str() + "/tmp", file_system_);
1339+
ASSERT_OK_AND_ASSIGN(
1340+
auto merge_writer,
1341+
MergeTreeWriter::Create(/*last_sequence_number=*/-1, primary_keys_, path_factory,
1342+
key_comparator_, /*user_defined_seq_comparator=*/nullptr,
1343+
merge_function_wrapper_, /*schema_id=*/0, value_schema_, options,
1344+
noop_compact_manager_, io_manager, pool_));
1345+
1346+
std::shared_ptr<arrow::Array> array =
1347+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1348+
["Alice", 1, 0, 1.0],
1349+
["Bob", 2, 0, 2.0]
1350+
])")
1351+
.ValueOrDie();
1352+
WriteBatch(array, /*row_kinds=*/{}, merge_writer.get());
1353+
1354+
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1355+
ASSERT_OK(merge_writer->Close());
1356+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1357+
}
1358+
1359+
TEST_P(MergeTreeWriterTest, TestMultiplePrepareCommitWithSpill) {
1360+
ASSERT_OK_AND_ASSIGN(
1361+
CoreOptions options,
1362+
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"}, {Options::WRITE_ONLY, "true"}}));
1363+
auto dir = UniqueTestDirectory::Create();
1364+
ASSERT_TRUE(dir);
1365+
auto path_factory = std::make_shared<DataFilePathFactory>();
1366+
ASSERT_OK(path_factory->Init(dir->Str(), "orc", options.DataFilePrefix(), nullptr));
1367+
std::string uuid = path_factory->uuid_;
1368+
1369+
std::shared_ptr<IOManager> io_manager =
1370+
std::make_shared<IOManager>(dir->Str() + "/tmp", file_system_);
1371+
ASSERT_OK_AND_ASSIGN(
1372+
auto merge_writer,
1373+
MergeTreeWriter::Create(/*last_sequence_number=*/-1, primary_keys_, path_factory,
1374+
key_comparator_, /*user_defined_seq_comparator=*/nullptr,
1375+
merge_function_wrapper_, /*schema_id=*/0, value_schema_, options,
1376+
noop_compact_manager_, io_manager, pool_));
1377+
1378+
std::shared_ptr<arrow::Array> array1 =
1379+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1380+
["Alice", 1, 0, 1.0],
1381+
["Bob", 2, 0, 2.0]
1382+
])")
1383+
.ValueOrDie();
1384+
WriteBatch(array1, /*row_kinds=*/{}, merge_writer.get());
1385+
ASSERT_OK(merge_writer->FlushMemory());
1386+
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1387+
1388+
ASSERT_OK_AND_ASSIGN(CommitIncrement commit1,
1389+
merge_writer->PrepareCommit(/*wait_compaction=*/false));
1390+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1391+
ASSERT_EQ(1, commit1.GetNewFilesIncrement().NewFiles().size());
1392+
1393+
std::string expected_path1 = dir->Str() + "/data-" + uuid + "-0.orc";
1394+
std::shared_ptr<arrow::ChunkedArray> expected_array1;
1395+
auto status1 = arrow::ipc::internal::json::ChunkedArrayFromJSON(write_type_, {R"([
1396+
[0, 0, "Alice", 1, 0, 1.0],
1397+
[1, 0, "Bob", 2, 0, 2.0]
1398+
])"},
1399+
&expected_array1);
1400+
ASSERT_TRUE(status1.ok());
1401+
CheckFileContent(expected_path1, expected_array1);
1402+
1403+
std::shared_ptr<arrow::Array> array2 =
1404+
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
1405+
["Dave", 4, 0, 4.0],
1406+
["Eve", 5, 0, 5.0]
1407+
])")
1408+
.ValueOrDie();
1409+
WriteBatch(array2, /*row_kinds=*/{}, merge_writer.get());
1410+
ASSERT_OK(merge_writer->FlushMemory());
1411+
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1412+
1413+
ASSERT_OK_AND_ASSIGN(CommitIncrement commit2,
1414+
merge_writer->PrepareCommit(/*wait_compaction=*/false));
1415+
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1416+
ASSERT_EQ(1, commit2.GetNewFilesIncrement().NewFiles().size());
1417+
1418+
std::string expected_path2 = dir->Str() + "/data-" + uuid + "-1.orc";
1419+
std::shared_ptr<arrow::ChunkedArray> expected_array2;
1420+
auto status2 = arrow::ipc::internal::json::ChunkedArrayFromJSON(write_type_, {R"([
1421+
[2, 0, "Dave", 4, 0, 4.0],
1422+
[3, 0, "Eve", 5, 0, 5.0]
1423+
])"},
1424+
&expected_array2);
1425+
ASSERT_TRUE(status2.ok());
1426+
CheckFileContent(expected_path2, expected_array2);
1427+
1428+
ASSERT_OK(merge_writer->Close());
1429+
}
1430+
10781431
INSTANTIATE_TEST_SUITE_P(WithOptionalIOManager, MergeTreeWriterTest,
10791432
::testing::Values(false, true));
10801433

0 commit comments

Comments
 (0)