@@ -1111,7 +1111,8 @@ TEST_F(MergeTreeWriterTest, TestSpillWithSameKeyDeduplicate) {
11111111
11121112 WriteBatch (batch1, /* row_kinds=*/ {}, merge_writer.get ());
11131113 WriteBatch (batch2, /* row_kinds=*/ {}, merge_writer.get ());
1114- // actual_max_fan_in_=2 with 1-byte budget, so 2 spill files compact to 1.
1114+ // WRITE_BUFFER_SIZE=1 causes UpdateSpillParameters() to clamp actual_max_fan_in_ to 2,
1115+ // triggering leveled compaction after 2 spill files are produced, merging them into 1.
11151116 ASSERT_EQ (1u , TestHelper::CountChannelFiles (file_system_, dir->Str () + " /tmp" ));
11161117
11171118 std::shared_ptr<arrow::Array> batch3 =
@@ -1179,13 +1180,15 @@ TEST_F(MergeTreeWriterTest, TestIntermediateMergeSpillFileBound) {
11791180 .ValueOrDie ();
11801181
11811182 WriteBatch (batch1, /* row_kinds=*/ {}, merge_writer.get ());
1183+ // Level 0: [A], total = 1
11821184 ASSERT_EQ (1u , TestHelper::CountChannelFiles (file_system_, dir->Str () + " /tmp" ));
11831185
11841186 WriteBatch (batch2, /* row_kinds=*/ {}, merge_writer.get ());
1187+ // Level 0: [A,B] hits max_fan_in=2, compaction -> Level 0: [], Level 1: [C], total = 1
11851188 ASSERT_EQ (1u , TestHelper::CountChannelFiles (file_system_, dir->Str () + " /tmp" ));
11861189
11871190 WriteBatch (batch3, /* row_kinds=*/ {}, merge_writer.get ());
1188- // After 3rd spill: 1 file at level 0 (new) + 1 file at level 1 (from prior compact) = 2.
1191+ // Level 0: [D], Level 1: [C], total = 2 (no single level exceeds max_fan_in)
11891192 ASSERT_EQ (2u , TestHelper::CountChannelFiles (file_system_, dir->Str () + " /tmp" ));
11901193
11911194 ASSERT_OK_AND_ASSIGN (CommitIncrement commit_increment,
@@ -1430,6 +1433,102 @@ TEST_F(MergeTreeWriterTest, TestMultiplePrepareCommitWithSpill) {
14301433 ASSERT_OK (merge_writer->Close ());
14311434}
14321435
1436+ TEST_P (MergeTreeWriterTest, TestSpillWithIOException) {
1437+ // This test exercises IO error injection across all spill-related code paths:
1438+ // 1. SpillToDisk (spill write)
1439+ // 2. MergeAndReplaceFiles (intermediate compaction triggered by LeveledMerger)
1440+ // 3. RunFinalCleanupIfNeeded (final merge in CreateReaders/PrepareCommit)
1441+ // 4. FlushWriteBuffer (reading merged spill data back + writing output data file)
1442+ //
1443+ // WRITE_BUFFER_SIZE=1 causes actual_max_fan_in_ to be clamped to 2, so every
1444+ // 2 spill files at the same level triggers compaction.
1445+ // Each WriteBatch with 2 rows fills the in-memory buffer (write_buffer_size param
1446+ // in InMemorySortBuffer is controlled via WRITE_BUFFER_SIZE in CoreOptions for
1447+ // MergeTreeWriter::Create), triggering a spill.
1448+ if (!GetParam ()) {
1449+ return ;
1450+ }
1451+
1452+ ASSERT_OK_AND_ASSIGN (CoreOptions options,
1453+ CoreOptions::FromMap ({{Options::FILE_FORMAT , " orc" },
1454+ {Options::WRITE_BUFFER_SIZE , " 1" },
1455+ {Options::WRITE_ONLY , " true" },
1456+ {Options::LOCAL_SORT_MAX_NUM_FILE_HANDLES , " 2" }}));
1457+
1458+ bool run_complete = false ;
1459+ auto io_hook = IOHook::GetInstance ();
1460+ for (size_t i = 0 ; i < 2000 ; i++) {
1461+ auto dir = UniqueTestDirectory::Create ();
1462+ ASSERT_TRUE (dir);
1463+ ScopeGuard guard ([&io_hook]() { io_hook->Clear (); });
1464+ io_hook->Reset (i, IOHook::Mode::RETURN_ERROR );
1465+ auto path_factory = std::make_shared<DataFilePathFactory>();
1466+ ASSERT_OK (path_factory->Init (dir->Str (), " orc" , options.DataFilePrefix (), nullptr ));
1467+
1468+ auto merge_writer_result = CreateMergeWriter (
1469+ /* last_sequence_number=*/ -1 , dir->Str (), path_factory, /* schema_id=*/ 0 , options);
1470+ CHECK_HOOK_STATUS (merge_writer_result.status (), i);
1471+ auto merge_writer = std::move (merge_writer_result).value ();
1472+
1473+ // Write 4 batches, each with 2 rows sharing the same key to exercise deduplication.
1474+ // Batch 1: triggers spill file 1
1475+ std::shared_ptr<arrow::Array> batch1 =
1476+ arrow::ipc::internal::json::ArrayFromJSON (value_type_, R"( [
1477+ ["Alice", 1, 0, 1.0],
1478+ ["Bob", 2, 0, 2.0]
1479+ ])" )
1480+ .ValueOrDie ();
1481+ auto b1 = CreateBatch (batch1, {});
1482+ CHECK_HOOK_STATUS (merge_writer->Write (std::move (b1)), i);
1483+
1484+ // Batch 2: triggers spill file 2 → intermediate compaction (merge 2 files into 1)
1485+ std::shared_ptr<arrow::Array> batch2 =
1486+ arrow::ipc::internal::json::ArrayFromJSON (value_type_, R"( [
1487+ ["Alice", 10, 0, 10.0],
1488+ ["Charlie", 3, 0, 3.0]
1489+ ])" )
1490+ .ValueOrDie ();
1491+ auto b2 = CreateBatch (batch2, {});
1492+ CHECK_HOOK_STATUS (merge_writer->Write (std::move (b2)), i);
1493+
1494+ // Batch 3: triggers spill file at level 0 again
1495+ std::shared_ptr<arrow::Array> batch3 =
1496+ arrow::ipc::internal::json::ArrayFromJSON (value_type_, R"( [
1497+ ["Bob", 20, 0, 20.0],
1498+ ["Dave", 4, 0, 4.0]
1499+ ])" )
1500+ .ValueOrDie ();
1501+ auto b3 = CreateBatch (batch3, {});
1502+ CHECK_HOOK_STATUS (merge_writer->Write (std::move (b3)), i);
1503+
1504+ // Batch 4: triggers spill file at level 0 → another compaction at level 0,
1505+ // then level 1 has 2 files → compaction at level 1 as well.
1506+ std::shared_ptr<arrow::Array> batch4 =
1507+ arrow::ipc::internal::json::ArrayFromJSON (value_type_, R"( [
1508+ ["Charlie", 30, 0, 30.0],
1509+ ["Eve", 5, 0, 5.0]
1510+ ])" )
1511+ .ValueOrDie ();
1512+ auto b4 = CreateBatch (batch4, {});
1513+ CHECK_HOOK_STATUS (merge_writer->Write (std::move (b4)), i);
1514+
1515+ // PrepareCommit: triggers FlushWriteBuffer → CreateReaders (RunFinalCleanupIfNeeded)
1516+ // → sort merge → write output data file
1517+ auto commit_increment = merge_writer->PrepareCommit (/* wait_compaction=*/ false );
1518+ CHECK_HOOK_STATUS (commit_increment.status (), i);
1519+ ASSERT_FALSE (commit_increment.value ().GetNewFilesIncrement ().NewFiles ().empty ());
1520+
1521+ // Verify deduplication: Alice(seq=2), Bob(seq=4), Charlie(seq=5), Dave(seq=6), Eve(seq=7)
1522+ ASSERT_EQ (1 , commit_increment.value ().GetNewFilesIncrement ().NewFiles ().size ());
1523+ ASSERT_EQ (5 , commit_increment.value ().GetNewFilesIncrement ().NewFiles ()[0 ]->row_count );
1524+
1525+ ASSERT_OK (merge_writer->Close ());
1526+ run_complete = true ;
1527+ break ;
1528+ }
1529+ ASSERT_TRUE (run_complete);
1530+ }
1531+
14331532INSTANTIATE_TEST_SUITE_P (WithOptionalIOManager, MergeTreeWriterTest,
14341533 ::testing::Values (false , true ));
14351534
0 commit comments