Skip to content

Commit f052b2b

Browse files
authored
refactor: extract WriteBuffer from MergeTreeWriter (alibaba#206)
1 parent c415846 commit f052b2b

34 files changed

Lines changed: 646 additions & 305 deletions

src/paimon/CMakeLists.txt

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -239,6 +239,7 @@ set(PAIMON_CORE_SRCS
239239
core/mergetree/compact/lookup_merge_tree_compact_rewriter.cpp
240240
core/mergetree/compact/changelog_merge_tree_rewriter.cpp
241241
core/mergetree/merge_tree_writer.cpp
242+
core/mergetree/write_buffer.cpp
242243
core/mergetree/levels.cpp
243244
core/mergetree/lookup_levels.cpp
244245
core/migrate/file_meta_utils.cpp
@@ -592,6 +593,7 @@ if(PAIMON_BUILD_TESTS)
592593
core/mergetree/lookup/persist_processor_test.cpp
593594
core/mergetree/drop_delete_reader_test.cpp
594595
core/mergetree/merge_tree_writer_test.cpp
596+
core/mergetree/write_buffer_test.cpp
595597
core/mergetree/sorted_run_test.cpp
596598
core/migrate/file_meta_utils_test.cpp
597599
core/operation/metrics/compaction_metrics_test.cpp

src/paimon/common/data/internal_row_test.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ TEST(InternalRowTest, TestCreateFieldGetter) {
7070
};
7171

7272
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
73-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
73+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
7474
[true, 1, 2, 3, 4, 5.1, 6.12, "abc", "def", "1970-01-02 00:00:01", "1970-01-02 00:00:00.001",
7575
"1970-01-02 00:00:00.000001", "1970-01-02 00:00:00.000000001", "1970-01-02 00:00:02", "1970-01-02 00:00:00.002",
7676
"1970-01-02 00:00:00.000002", "1970-01-02 00:00:00.000000002", "-123456789987654321.45678", 12345,
@@ -138,7 +138,7 @@ TEST(InternalRowTest, TestCreateFieldGetterWithNull) {
138138
arrow::field("f1", arrow::int8())};
139139

140140
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
141-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
141+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
142142
[true, null]
143143
])")
144144
.ValueOrDie());
@@ -161,7 +161,7 @@ TEST(InternalRowTest, TestCreateFieldGetterWithInvalidType) {
161161
arrow::FieldVector fields = {arrow::field("f0", arrow::large_utf8())};
162162

163163
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
164-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
164+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
165165
["hello"]
166166
])")
167167
.ValueOrDie());

src/paimon/common/data/record_batch_test.cpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -87,7 +87,7 @@ TEST(RecordBatchTest, TestAssignAndMove) {
8787
arrow::field("f1", arrow::int8())};
8888
std::map<std::string, std::string> partition = {{"f1", "1"}};
8989
auto old_array = std::dynamic_pointer_cast<arrow::StructArray>(
90-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
90+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
9191
[true, 1]
9292
])")
9393
.ValueOrDie());
@@ -97,7 +97,7 @@ TEST(RecordBatchTest, TestAssignAndMove) {
9797
&old_arrow_array);
9898

9999
auto new_array = std::dynamic_pointer_cast<arrow::StructArray>(
100-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
100+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
101101
[false, 1]
102102
])")
103103
.ValueOrDie());

src/paimon/common/global_index/complete_index_score_batch_reader_test.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ TEST_F(CompleteIndexScoreBatchReaderTest, TestSimple) {
7272
arrow::field("_ROW_ID", arrow::int64()),
7373
};
7474

75-
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
75+
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
7676
["Alice", 10, null, 0],
7777
["Bob", 11, null, 1],
7878
["Cathy", 12, null, 2]
@@ -106,7 +106,7 @@ TEST_F(CompleteIndexScoreBatchReaderTest, TestWithBitmap) {
106106
arrow::field("_ROW_ID", arrow::int64()),
107107
};
108108

109-
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
109+
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
110110
["Alice", 10, null, 0],
111111
["Bob", 11, null, 1],
112112
["Cathy", 12, null, 2],
@@ -141,7 +141,7 @@ TEST_F(CompleteIndexScoreBatchReaderTest, TestReadWithNullScores) {
141141
arrow::field("_ROW_ID", arrow::int64()),
142142
};
143143

144-
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
144+
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
145145
["Alice", 10, null, 0],
146146
["Bob", 11, null, 1],
147147
["Cathy", 12, null, 2]

src/paimon/common/reader/complete_row_kind_batch_reader_test.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -139,7 +139,7 @@ TEST_F(CompleteRowKindBatchReaderTest, TestNestedType) {
139139
arrow::field("f1", arrow::map(arrow::struct_({field("a", arrow::int64()),
140140
field("b", arrow::boolean())}),
141141
arrow::boolean()))};
142-
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({fields}), R"([
142+
auto src_array = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([
143143
[[null, [1, true], null], [[[1, true], true]]],
144144
[[[2, false], null], null],
145145
[[[2, false], [3, true], [4, null]], [[[1, true], true], [[5, false], null]]],

src/paimon/core/io/concat_key_value_record_reader_test.cpp

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,7 @@ TEST_F(ConcatKeyValueRecordReaderTest, TestSimple) {
8484
arrow::schema(arrow::FieldVector({fields[2], fields[3]}));
8585
std::shared_ptr<arrow::Schema> value_schema =
8686
arrow::schema(arrow::FieldVector({fields[2], fields[3], fields[4], fields[5], fields[6]}));
87-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_({fields});
87+
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(fields);
8888
auto src_array1 = std::dynamic_pointer_cast<arrow::StructArray>(
8989
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
9090
[0, 0, 1, 1, 10, 20, 30],
@@ -130,7 +130,7 @@ TEST_F(ConcatKeyValueRecordReaderTest, TestSingleReaderInConcat) {
130130
arrow::schema(arrow::FieldVector({fields[2], fields[3]}));
131131
std::shared_ptr<arrow::Schema> value_schema =
132132
arrow::schema(arrow::FieldVector({fields[2], fields[3], fields[4], fields[5], fields[6]}));
133-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_({fields});
133+
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(fields);
134134
auto src_array1 = std::dynamic_pointer_cast<arrow::StructArray>(
135135
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
136136
[0, 0, 1, 1, 10, 20, 30],
@@ -159,7 +159,7 @@ TEST_F(ConcatKeyValueRecordReaderTest, TestEmptyResult) {
159159
arrow::schema(arrow::FieldVector({fields[2], fields[3]}));
160160
std::shared_ptr<arrow::Schema> value_schema =
161161
arrow::schema(arrow::FieldVector({fields[2], fields[3], fields[4], fields[5], fields[6]}));
162-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_({fields});
162+
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(fields);
163163
auto src_array1 = std::dynamic_pointer_cast<arrow::StructArray>(
164164
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
165165
])")
@@ -180,7 +180,7 @@ TEST_F(ConcatKeyValueRecordReaderTest, TestEmptyReader) {
180180
arrow::schema(arrow::FieldVector({fields[2], fields[3]}));
181181
std::shared_ptr<arrow::Schema> value_schema =
182182
arrow::schema(arrow::FieldVector({fields[2], fields[3], fields[4], fields[5], fields[6]}));
183-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_({fields});
183+
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(fields);
184184
auto src_array1 = std::dynamic_pointer_cast<arrow::StructArray>(
185185
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
186186
])")

src/paimon/core/io/field_mapping_reader_test.cpp

Lines changed: 10 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020
#include <ostream>
2121
#include <string>
2222
#include <utility>
23-
#include <variant>
2423

2524
#include "arrow/api.h"
2625
#include "arrow/array/array_base.h"
@@ -418,7 +417,7 @@ TEST_F(FieldMappingReaderTest, TestDictionaryTypeWithSchemaEvolution) {
418417
std::shared_ptr<arrow::Schema> data_schema =
419418
DataField::ConvertDataFieldsToArrowSchema(data_fields);
420419
auto data_array = std::dynamic_pointer_cast<arrow::StructArray>(
421-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({data_schema->fields()}), R"([
420+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(data_schema->fields()), R"([
422421
["apple", 4.0, 5.1, 10, 100, 1000, 10000, true],
423422
["banana", 4.1, 6.2, 10, 200, 1000, 20000, null],
424423
[null, 4.2, null, 10, 300, 1000, 30000, true],
@@ -443,7 +442,7 @@ TEST_F(FieldMappingReaderTest, TestDictionaryTypeWithSchemaEvolution) {
443442
std::vector<std::string> partition_keys = {"f3", "f5"};
444443
BinaryRow partition = BinaryRowGenerator::GenerateRow({10, 1000}, pool_.get());
445444
auto expected_array = std::dynamic_pointer_cast<arrow::StructArray>(
446-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({read_schema->fields()}), R"([
445+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(read_schema->fields()), R"([
447446
[10000, 5.1, "apple", 4.0, 1000, 10, null, true, 100],
448447
[20000, 6.2, "banana", 4.1, 1000, 10, null, null, 200],
449448
[30000, null, null, 4.2, 1000, 10, null, true, 300],
@@ -464,7 +463,7 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyType) {
464463
std::shared_ptr<arrow::Schema> data_schema =
465464
DataField::ConvertDataFieldsToArrowSchema(data_fields);
466465
auto data_array = std::dynamic_pointer_cast<arrow::StructArray>(
467-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({data_schema->fields()}), R"([
466+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(data_schema->fields()), R"([
468467
["true", 4.0, 5.1, 10],
469468
["False", 4.1, 6.2, 10],
470469
[null, 4.2, null, 10],
@@ -485,7 +484,7 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyType) {
485484
std::vector<std::string> partition_keys = {"f3"};
486485
BinaryRow partition = BinaryRowGenerator::GenerateRow({10}, pool_.get());
487486
auto expected_array = std::dynamic_pointer_cast<arrow::StructArray>(
488-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({read_schema->fields()}), R"([
487+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(read_schema->fields()), R"([
489488
[true, "4", 5, 10],
490489
[false, "4.1", 6, 10],
491490
[null, "4.2", null, 10],
@@ -506,7 +505,7 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyTypeWithDict) {
506505
std::shared_ptr<arrow::Schema> data_schema =
507506
DataField::ConvertDataFieldsToArrowSchema(data_fields);
508507
auto data_array = std::dynamic_pointer_cast<arrow::StructArray>(
509-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({data_schema->fields()}), R"([
508+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(data_schema->fields()), R"([
510509
["true", 4.0, 5.1, 10],
511510
["false", 4.1, 6.2, 10],
512511
[null, 4.2, null, 10],
@@ -527,7 +526,7 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyTypeWithDict) {
527526
std::vector<std::string> partition_keys = {"f3"};
528527
BinaryRow partition = BinaryRowGenerator::GenerateRow({10}, pool_.get());
529528
auto expected_array = std::dynamic_pointer_cast<arrow::StructArray>(
530-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({read_schema->fields()}), R"([
529+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(read_schema->fields()), R"([
531530
[true, "4", 5, 10],
532531
[false, "4.1", 6, 10],
533532
[null, "4.2", null, 10],
@@ -548,7 +547,7 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyTypeWithPredicate) {
548547
std::shared_ptr<arrow::Schema> data_schema =
549548
DataField::ConvertDataFieldsToArrowSchema(data_fields);
550549
auto data_array = std::dynamic_pointer_cast<arrow::StructArray>(
551-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({data_schema->fields()}), R"([
550+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(data_schema->fields()), R"([
552551
["true", 4.0, 5, 10],
553552
["False", 4.1, 6, 10],
554553
[null, 4.2, null, 10],
@@ -615,7 +614,7 @@ TEST_F(FieldMappingReaderTest, TestReadWithSchemaEvolutionWithRenameAndModifyTyp
615614
/*field_index=*/3, /*field_name=*/"f0", FieldType::STRING,
616615
Literal(FieldType::STRING, literal_str.data(), literal_str.size()));
617616
auto expected_array = std::dynamic_pointer_cast<arrow::StructArray>(
618-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({read_schema->fields()}), R"([
617+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(read_schema->fields()), R"([
619618
[0, "Emily", null, "15.1", 10],
620619
[0, "Bob", null, "12.1", 10],
621620
[0, "Alex", null, "16.1", 10]
@@ -632,7 +631,7 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithDictType) {
632631
std::shared_ptr<arrow::Schema> data_schema =
633632
DataField::ConvertDataFieldsToArrowSchema(data_fields);
634633
auto data_array = std::dynamic_pointer_cast<arrow::StructArray>(
635-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({data_schema->fields()}), R"([
634+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(data_schema->fields()), R"([
636635
["Bob", 4.0, 5.1, 10],
637636
["Emily", 4.1, 6.2, 10],
638637
["Alice", 4.2, null, 10],
@@ -653,7 +652,7 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithDictType) {
653652
std::vector<std::string> partition_keys = {"f3"};
654653
BinaryRow partition = BinaryRowGenerator::GenerateRow({10}, pool_.get());
655654
auto expected_array = std::dynamic_pointer_cast<arrow::StructArray>(
656-
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_({read_schema->fields()}), R"([
655+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(read_schema->fields()), R"([
657656
[10, "4", "Bob", 5],
658657
[10, "4.1", "Emily", 6],
659658
[10, "4.2", "Alice", null],

src/paimon/core/io/key_value_data_file_record_reader_test.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -317,7 +317,7 @@ TEST_F(KeyValueDataFileRecordReaderTest, TestEmptyReader) {
317317
arrow::schema(arrow::FieldVector({fields[2], fields[3]}));
318318
std::shared_ptr<arrow::Schema> value_schema =
319319
arrow::schema(arrow::FieldVector({fields[2], fields[3], fields[4], fields[5], fields[6]}));
320-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_({fields});
320+
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(fields);
321321
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
322322
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
323323
])")
@@ -347,7 +347,7 @@ TEST_F(KeyValueDataFileRecordReaderTest, TestInvalidSequenceNumerColumn) {
347347
arrow::schema(arrow::FieldVector({fields[2], fields[3]}));
348348
std::shared_ptr<arrow::Schema> value_schema =
349349
arrow::schema(arrow::FieldVector({fields[2], fields[3], fields[4], fields[5], fields[6]}));
350-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_({fields});
350+
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(fields);
351351
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
352352
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
353353
["0", 0, 1, 1, 10, 20, 30]
@@ -376,7 +376,7 @@ TEST_F(KeyValueDataFileRecordReaderTest, TestInvalidValueKindColumn) {
376376
arrow::schema(arrow::FieldVector({fields[2], fields[3]}));
377377
std::shared_ptr<arrow::Schema> value_schema =
378378
arrow::schema(arrow::FieldVector({fields[2], fields[3], fields[4], fields[5], fields[6]}));
379-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_({fields});
379+
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(fields);
380380
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
381381
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
382382
[0, 100, 1, 1, 10, 20, 30]

0 commit comments

Comments
 (0)