1818
1919#include < algorithm>
2020#include < cassert>
21- #include < cstddef>
2221#include < unordered_set>
2322#include < utility>
2423
2928#include " paimon/common/table/special_fields.h"
3029#include " paimon/common/utils/arrow/status_utils.h"
3130#include " paimon/common/utils/scope_guard.h"
31+ #include " paimon/core/disk/io_manager.h"
3232#include " paimon/core/io/async_key_value_producer_and_consumer.h"
3333#include " paimon/core/io/compact_increment.h"
3434#include " paimon/core/io/data_file_path_factory.h"
4444#include " paimon/core/utils/commit_increment.h"
4545#include " paimon/format/file_format.h"
4646#include " paimon/format/writer_builder.h"
47- #include " paimon/metrics.h"
4847
4948namespace paimon {
5049class FormatStatsExtractor ;
5150
52- MergeTreeWriter:: MergeTreeWriter (
51+ Result<std::shared_ptr< MergeTreeWriter>> MergeTreeWriter::Create (
5352 int64_t last_sequence_number, const std::vector<std::string>& trimmed_primary_keys,
5453 const std::shared_ptr<DataFilePathFactory>& path_factory,
5554 const std::shared_ptr<FieldsComparator>& key_comparator,
5655 const std::shared_ptr<FieldsComparator>& user_defined_seq_comparator,
5756 const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper,
5857 int64_t schema_id, const std::shared_ptr<arrow::Schema>& value_schema,
5958 const CoreOptions& options, const std::shared_ptr<CompactManager>& compact_manager,
60- const std::shared_ptr<MemoryPool>& pool)
59+ const std::shared_ptr<IOManager>& io_manager, const std::shared_ptr<MemoryPool>& pool) {
60+ auto write_schema = SpecialFields::CompleteSequenceAndValueKindField (value_schema);
61+ PAIMON_ASSIGN_OR_RAISE (
62+ std::unique_ptr<WriteBuffer> write_buffer,
63+ WriteBuffer::Create (last_sequence_number, value_schema, trimmed_primary_keys,
64+ options.GetSequenceField (), key_comparator, user_defined_seq_comparator,
65+ merge_function_wrapper, options, io_manager, pool));
66+ return std::shared_ptr<MergeTreeWriter>(
67+ new MergeTreeWriter (pool, trimmed_primary_keys, options, path_factory, key_comparator,
68+ user_defined_seq_comparator, merge_function_wrapper, schema_id,
69+ write_schema, compact_manager, std::move (write_buffer)));
70+ }
71+
72+ MergeTreeWriter::MergeTreeWriter (
73+ const std::shared_ptr<MemoryPool>& pool, const std::vector<std::string>& trimmed_primary_keys,
74+ const CoreOptions& options, const std::shared_ptr<DataFilePathFactory>& path_factory,
75+ const std::shared_ptr<FieldsComparator>& key_comparator,
76+ const std::shared_ptr<FieldsComparator>& user_defined_seq_comparator,
77+ const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper,
78+ int64_t schema_id, const std::shared_ptr<arrow::Schema>& write_schema,
79+ const std::shared_ptr<CompactManager>& compact_manager,
80+ std::unique_ptr<WriteBuffer>&& write_buffer)
6181 : pool_(pool),
6282 trimmed_primary_keys_ (trimmed_primary_keys),
6383 options_(options),
@@ -66,13 +86,10 @@ MergeTreeWriter::MergeTreeWriter(
6686 user_defined_seq_comparator_(user_defined_seq_comparator),
6787 merge_function_wrapper_(merge_function_wrapper),
6888 schema_id_(schema_id),
89+ write_schema_(write_schema),
6990 compact_manager_(compact_manager),
70- metrics_(std::make_shared<MetricsImpl>()) {
71- write_schema_ = SpecialFields::CompleteSequenceAndValueKindField (value_schema);
72- write_buffer_ = std::make_unique<WriteBuffer>(
73- last_sequence_number, arrow::struct_ (value_schema->fields ()), trimmed_primary_keys_,
74- options_.GetSequenceField (), key_comparator_, merge_function_wrapper_, pool_);
75- }
91+ write_buffer_(std::move(write_buffer)),
92+ metrics_(std::make_shared<MetricsImpl>()) {}
7693
7794Status MergeTreeWriter::DoClose () {
7895 // Request cancellation and wait for running compaction to exit.
@@ -115,16 +132,26 @@ Status MergeTreeWriter::DoClose() {
115132 return Status::OK ();
116133}
117134
135+ Status MergeTreeWriter::FlushMemory () {
136+ PAIMON_ASSIGN_OR_RAISE (bool has_remaining_quota, write_buffer_->FlushMemory ());
137+ if (!has_remaining_quota) {
138+ PAIMON_RETURN_NOT_OK (FlushWriteBuffer (/* wait_for_latest_compaction=*/ false ,
139+ /* forced_full_compaction=*/ false ));
140+ }
141+ return Status::OK ();
142+ }
143+
118144Status MergeTreeWriter::Write (std::unique_ptr<RecordBatch>&& moved_batch) {
119- PAIMON_RETURN_NOT_OK (write_buffer_->Write (std::move (moved_batch)));
120- if (write_buffer_->GetMemoryUsage () >= static_cast <uint64_t >(options_.GetWriteBufferSize ())) {
121- return Flush (/* wait_for_latest_compaction=*/ false , /* forced_full_compaction=*/ false );
145+ PAIMON_ASSIGN_OR_RAISE (bool has_remaining_quota, write_buffer_->Write (std::move (moved_batch)));
146+ if (!has_remaining_quota) {
147+ return FlushWriteBuffer (/* wait_for_latest_compaction=*/ false ,
148+ /* forced_full_compaction=*/ false );
122149 }
123150 return Status::OK ();
124151}
125152
126153Status MergeTreeWriter::Compact (bool full_compaction) {
127- return Flush (/* wait_for_latest_compaction=*/ true , full_compaction);
154+ return FlushWriteBuffer (/* wait_for_latest_compaction=*/ true , full_compaction);
128155}
129156
130157Status MergeTreeWriter::Sync () {
@@ -197,7 +224,7 @@ Status MergeTreeWriter::UpdateCompactDeletionFile(
197224}
198225
199226Result<CommitIncrement> MergeTreeWriter::PrepareCommit (bool wait_compaction) {
200- PAIMON_RETURN_NOT_OK (Flush (wait_compaction, /* forced_full_compaction=*/ false ));
227+ PAIMON_RETURN_NOT_OK (FlushWriteBuffer (wait_compaction, /* forced_full_compaction=*/ false ));
201228 if (options_.CommitForceCompact ()) {
202229 wait_compaction = true ;
203230 }
@@ -218,13 +245,14 @@ Result<bool> MergeTreeWriter::CompactNotCompleted() {
218245 return compact_manager_->CompactNotCompleted ();
219246}
220247
221- Status MergeTreeWriter::Flush (bool wait_for_latest_compaction, bool forced_full_compaction) {
248+ Status MergeTreeWriter::FlushWriteBuffer (bool wait_for_latest_compaction,
249+ bool forced_full_compaction) {
222250 if (!write_buffer_->IsEmpty ()) {
223251 if (compact_manager_->ShouldWaitForLatestCompaction ()) {
224252 wait_for_latest_compaction = true ;
225253 }
226254 auto cleanup_guard = ScopeGuard ([&]() { write_buffer_->Clear (); });
227- // 1. flush write buffer to get in-memory readers
255+ // 1. flush write buffer to get sorted readers
228256 PAIMON_ASSIGN_OR_RAISE (std::vector<std::unique_ptr<KeyValueRecordReader>> readers,
229257 write_buffer_->CreateReaders ());
230258 // 2. prepare loser tree sort merge reader
@@ -257,6 +285,7 @@ Status MergeTreeWriter::Flush(bool wait_for_latest_compaction, bool forced_full_
257285 PAIMON_RETURN_NOT_OK (rolling_writer->Close ());
258286 PAIMON_ASSIGN_OR_RAISE (std::vector<std::shared_ptr<DataFileMeta>> flushed_files,
259287 rolling_writer->GetResult ());
288+ async_key_value_producer_consumer->Close ();
260289 write_guard.Release ();
261290
262291 for (const auto & flushed_file : flushed_files) {
0 commit comments