@@ -218,8 +218,7 @@ TEST_F(KeyValueFileStoreWriteTest,
218218 ASSERT_EQ (commit_messages.size (), 1 );
219219}
220220
221- TEST_F (KeyValueFileStoreWriteTest,
222- TestWriteShouldSpillLargestWriterWhenGlobalMemoryLimitIsExceeded) {
221+ TEST_F (KeyValueFileStoreWriteTest, TestSpillSimple) {
223222 auto fields = {arrow::field (" f0" , arrow::utf8 (), /* nullable=*/ false )};
224223 arrow::Schema typed_schema (fields);
225224 ::ArrowSchema schema;
@@ -241,17 +240,38 @@ TEST_F(KeyValueFileStoreWriteTest,
241240
242241 ASSERT_OK_AND_ASSIGN (std::unique_ptr<WriteContext> write_context, context_builder.Finish ());
243242 ASSERT_OK_AND_ASSIGN (auto file_store_write, FileStoreWrite::Create (std::move (write_context)));
243+ auto key_value_file_store_write = dynamic_cast <KeyValueFileStoreWrite*>(file_store_write.get ());
244+ auto get_writer = [&](int32_t bucket) -> std::shared_ptr<paimon::BatchWriter> {
245+ auto partition_iter = key_value_file_store_write->writers_ .find (BinaryRow::EmptyRow ());
246+ if (partition_iter != key_value_file_store_write->writers_ .end ()) {
247+ auto & buckets = partition_iter->second ;
248+ auto bucket_iter = buckets.find (bucket);
249+ if (PAIMON_LIKELY (bucket_iter != buckets.end ())) {
250+ return bucket_iter->second .writer ;
251+ }
252+ }
253+ assert (false );
254+ return nullptr ;
255+ };
244256
257+ // write bucket 0, not trigger spill
245258 ASSERT_OK (WriteSingleStringRow (file_store_write.get (), /* bucket=*/ 0 , std::string (48 , ' a' )));
246259 ASSERT_EQ (TestHelper::CountChannelFiles (dir->GetFileSystem (), dir->Str ()), 0 );
260+ ASSERT_GT (get_writer (0 )->GetMemoryUsage (), 0 );
247261
248- ASSERT_OK (WriteSingleStringRow (file_store_write.get (), /* bucket=*/ 1 , std::string (48 , ' b' )));
262+ // write bucket 1, spill bucket 0 (pick largest writer)
263+ ASSERT_OK (WriteSingleStringRow (file_store_write.get (), /* bucket=*/ 1 , std::string (32 , ' b' )));
249264 ASSERT_GT (TestHelper::CountChannelFiles (dir->GetFileSystem (), dir->Str ()), 0 );
265+ ASSERT_EQ (get_writer (0 )->GetMemoryUsage (), 0 );
266+ ASSERT_GT (get_writer (1 )->GetMemoryUsage (), 0 );
250267
268+ // prepare commit, clean all spill files and memory buffers
251269 ASSERT_OK_AND_ASSIGN (auto commit_messages,
252270 file_store_write->PrepareCommit (/* wait_compaction=*/ true ));
253271 ASSERT_EQ (commit_messages.size (), 2 );
254272 ASSERT_EQ (TestHelper::CountChannelFiles (dir->GetFileSystem (), dir->Str ()), 0 );
273+ ASSERT_EQ (get_writer (0 )->GetMemoryUsage (), 0 );
274+ ASSERT_EQ (get_writer (1 )->GetMemoryUsage (), 0 );
255275}
256276
257277} // namespace paimon::test
0 commit comments