Skip to content

Commit b631683

Browse files
authored
feat(shared_shredding): add read inte test for shared_shredding (alibaba#382)
1 parent ecfa667 commit b631683

19 files changed

Lines changed: 1670 additions & 293 deletions

src/paimon/CMakeLists.txt

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -142,7 +142,7 @@ set(PAIMON_COMMON_SRCS
142142
common/data/shredding/map_shared_shredding_context.cpp
143143
common/data/shredding/map_shared_shredding_batch_converter.cpp
144144
common/data/shredding/map_shared_shredding_column_allocator.cpp
145-
common/data/shredding/shared_shredding_file_reader.cpp
145+
common/data/shredding/map_shared_shredding_file_reader.cpp
146146
common/utils/delta_varint_compressor.cpp
147147
common/utils/fields_comparator.cpp
148148
common/utils/path_util.cpp
@@ -550,7 +550,7 @@ if(PAIMON_BUILD_TESTS)
550550
common/data/shredding/map_shared_shredding_column_allocator_test.cpp
551551
common/data/shredding/map_shared_shredding_field_dict_test.cpp
552552
common/data/shredding/map_shared_shredding_context_test.cpp
553-
common/data/shredding/shared_shredding_file_reader_test.cpp
553+
common/data/shredding/map_shared_shredding_file_reader_test.cpp
554554
STATIC_LINK_LIBS
555555
paimon_shared
556556
test_utils_static

src/paimon/common/data/shredding/shared_shredding_file_reader.cpp renamed to src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp

Lines changed: 50 additions & 128 deletions
Large diffs are not rendered by default.

src/paimon/common/data/shredding/shared_shredding_file_reader.h renamed to src/paimon/common/data/shredding/map_shared_shredding_file_reader.h

Lines changed: 18 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,9 @@
1818

1919
#include <map>
2020
#include <memory>
21+
#include <optional>
2122
#include <string>
23+
#include <utility>
2224
#include <vector>
2325

2426
#include "arrow/api.h"
@@ -28,10 +30,22 @@
2830

2931
namespace paimon {
3032

31-
class SharedShreddingFileReader : public FileBatchReader {
33+
class MapSharedShreddingFileReader : public FileBatchReader {
3234
public:
33-
static Result<std::unique_ptr<SharedShreddingFileReader>> Create(
34-
std::unique_ptr<FileBatchReader>&& reader, const std::shared_ptr<MemoryPool>& pool);
35+
struct SharedShreddingContext {
36+
SharedShreddingContext(const MapSharedShreddingFieldMeta& _meta,
37+
const std::vector<std::string>& _selected_keys,
38+
const std::shared_ptr<arrow::MapType>& _map_type)
39+
: meta(_meta), selected_keys(_selected_keys), map_type(_map_type) {}
40+
MapSharedShreddingFieldMeta meta;
41+
std::vector<std::string> selected_keys;
42+
std::shared_ptr<arrow::MapType> map_type;
43+
};
44+
45+
MapSharedShreddingFileReader(
46+
std::unique_ptr<FileBatchReader>&& reader,
47+
std::map<std::string, SharedShreddingContext>&& shared_shredding_name_to_context,
48+
const std::shared_ptr<MemoryPool>& pool);
3549

3650
Result<std::unique_ptr<::ArrowSchema>> GetFileSchema() const override;
3751

@@ -53,11 +67,6 @@ class SharedShreddingFileReader : public FileBatchReader {
5367
bool SupportPreciseBitmapSelection() const override;
5468

5569
private:
56-
SharedShreddingFileReader(
57-
std::unique_ptr<FileBatchReader>&& reader,
58-
const std::map<std::string, MapSharedShreddingFieldMeta>& shared_shredding_name_to_meta,
59-
const std::shared_ptr<MemoryPool>& pool);
60-
6170
Result<std::shared_ptr<arrow::Array>> RebuildLogicalMapArray(
6271
const std::shared_ptr<arrow::Field>& physical_field,
6372
const std::shared_ptr<arrow::StructArray>& physical_struct_array) const;
@@ -76,9 +85,7 @@ class SharedShreddingFileReader : public FileBatchReader {
7685
private:
7786
std::shared_ptr<arrow::MemoryPool> arrow_pool_;
7887
std::unique_ptr<FileBatchReader> reader_;
79-
std::map<std::string, MapSharedShreddingFieldMeta> shared_shredding_name_to_meta_;
80-
std::map<std::string, std::vector<std::string>> shared_shredding_name_to_selected_keys_;
81-
std::map<std::string, std::shared_ptr<arrow::MapType>> shared_shredding_name_to_map_type_;
88+
std::map<std::string, SharedShreddingContext> shared_shredding_name_to_context_;
8289
};
8390

8491
} // namespace paimon

src/paimon/common/data/shredding/shared_shredding_file_reader_test.cpp renamed to src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp

Lines changed: 81 additions & 50 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@
1414
* limitations under the License.
1515
*/
1616

17-
#include "paimon/common/data/shredding/shared_shredding_file_reader.h"
17+
#include "paimon/common/data/shredding/map_shared_shredding_file_reader.h"
1818

1919
#include <map>
2020
#include <memory>
@@ -47,7 +47,7 @@
4747
#include "paimon/testing/utils/testharness.h"
4848

4949
namespace paimon::test {
50-
class SharedShreddingFileReaderTest : public ::testing::Test {
50+
class MapSharedShreddingFileReaderTest : public ::testing::Test {
5151
public:
5252
void SetUp() override {
5353
pool_ = GetDefaultPool();
@@ -94,9 +94,56 @@ class SharedShreddingFileReaderTest : public ::testing::Test {
9494
.ValueOrDie();
9595
}
9696

97-
std::unique_ptr<SharedShreddingFileReader> CreateReader(
97+
std::unique_ptr<MapSharedShreddingFileReader> WrapReader(
98+
std::unique_ptr<FileBatchReader>&& reader,
99+
const std::optional<std::string>& selected_keys_str = std::nullopt) const {
100+
EXPECT_OK_AND_ASSIGN(auto c_file_schema, reader->GetFileSchema());
101+
auto file_schema = arrow::ImportSchema(c_file_schema.get()).ValueOrDie();
102+
std::map<std::string, MapSharedShreddingFileReader::SharedShreddingContext>
103+
shared_shredding_name_to_context;
104+
for (const auto& field : file_schema->fields()) {
105+
auto metadata = std::const_pointer_cast<arrow::KeyValueMetadata>(field->metadata());
106+
if (!MapSharedShreddingUtils::HasShreddingMetadata(metadata)) {
107+
continue;
108+
}
109+
EXPECT_OK_AND_ASSIGN(auto meta,
110+
MapSharedShreddingUtils::DeserializeMetadata(
111+
metadata, MapSharedShreddingDefine::kDefaultDictCompression));
112+
auto physical_type =
113+
arrow::internal::checked_pointer_cast<arrow::StructType>(field->type());
114+
std::shared_ptr<arrow::Field> item_field;
115+
for (const auto& child : physical_type->fields()) {
116+
if (child->name() != MapSharedShreddingDefine::kFieldMapping &&
117+
child->name() != MapSharedShreddingDefine::kOverflow) {
118+
item_field = child;
119+
break;
120+
}
121+
}
122+
EXPECT_TRUE(item_field);
123+
auto map_type = arrow::internal::checked_pointer_cast<arrow::MapType>(arrow::map(
124+
arrow::utf8(), arrow::field("value", item_field->type(), item_field->nullable())));
125+
std::vector<std::string> selected_keys;
126+
if (selected_keys_str.has_value()) {
127+
selected_keys = StringUtils::Split(selected_keys_str.value(), ",",
128+
/*ignore_empty=*/false);
129+
} else {
130+
selected_keys.reserve(meta.name_to_id.size());
131+
for (const auto& [key_name, _] : meta.name_to_id) {
132+
selected_keys.push_back(key_name);
133+
}
134+
}
135+
shared_shredding_name_to_context.emplace(
136+
field->name(), MapSharedShreddingFileReader::SharedShreddingContext(
137+
meta, selected_keys, map_type));
138+
}
139+
return std::make_unique<MapSharedShreddingFileReader>(
140+
std::move(reader), std::move(shared_shredding_name_to_context), pool_);
141+
}
142+
143+
std::unique_ptr<MapSharedShreddingFileReader> CreateReader(
98144
std::shared_ptr<arrow::Array> physical_array = nullptr,
99-
std::shared_ptr<arrow::Schema> physical_schema = nullptr) const {
145+
std::shared_ptr<arrow::Schema> physical_schema = nullptr,
146+
const std::optional<std::string>& selected_keys = std::nullopt) const {
100147
if (!physical_schema) {
101148
physical_schema = PhysicalSchemaWithMetadata();
102149
}
@@ -106,9 +153,7 @@ class SharedShreddingFileReaderTest : public ::testing::Test {
106153
auto mock_reader = std::make_unique<MockFileBatchReader>(
107154
physical_array, arrow::struct_(physical_schema->fields()), /*read_batch_size=*/10);
108155
mock_reader->EnableRandomizeBatchSize(false);
109-
EXPECT_OK_AND_ASSIGN(auto shared_shredding_reader,
110-
SharedShreddingFileReader::Create(std::move(mock_reader), pool_));
111-
return shared_shredding_reader;
156+
return WrapReader(std::move(mock_reader), selected_keys);
112157
}
113158

114159
std::shared_ptr<arrow::Schema> ReadSchema(
@@ -196,7 +241,7 @@ class SharedShreddingFileReaderTest : public ::testing::Test {
196241
};
197242
};
198243

199-
TEST_F(SharedShreddingFileReaderTest, TestGetFileSchemaReturnsLogicalMapSchema) {
244+
TEST_F(MapSharedShreddingFileReaderTest, TestGetFileSchemaReturnsLogicalMapSchema) {
200245
auto reader = CreateReader();
201246

202247
ASSERT_OK_AND_ASSIGN(auto c_schema, reader->GetFileSchema());
@@ -209,8 +254,9 @@ TEST_F(SharedShreddingFileReaderTest, TestGetFileSchemaReturnsLogicalMapSchema)
209254
ASSERT_FALSE(schema->field(1)->HasMetadata());
210255
}
211256

212-
TEST_F(SharedShreddingFileReaderTest, TestAllExistSelectedKeysWithoutOverflow) {
213-
auto reader = CreateReader();
257+
TEST_F(MapSharedShreddingFileReaderTest, TestAllExistSelectedKeysWithoutOverflow) {
258+
auto reader = CreateReader(/*physical_array=*/nullptr, /*physical_schema=*/nullptr,
259+
/*selected_keys=*/"b");
214260
auto read_schema = ExportSchema(ReadSchema("b"));
215261
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
216262
/*selection_bitmap=*/std::nullopt));
@@ -229,8 +275,9 @@ TEST_F(SharedShreddingFileReaderTest, TestAllExistSelectedKeysWithoutOverflow) {
229275
AssertChunkedArrayEquals(expected, actual);
230276
}
231277

232-
TEST_F(SharedShreddingFileReaderTest, TestAllExistSelectedKeysWithOverflow) {
233-
auto reader = CreateReader();
278+
TEST_F(MapSharedShreddingFileReaderTest, TestAllExistSelectedKeysWithOverflow) {
279+
auto reader = CreateReader(/*physical_array=*/nullptr, /*physical_schema=*/nullptr,
280+
/*selected_keys=*/"a,c");
234281
auto read_schema = ExportSchema(ReadSchema("a,c"));
235282
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
236283
/*selection_bitmap=*/std::nullopt));
@@ -249,8 +296,9 @@ TEST_F(SharedShreddingFileReaderTest, TestAllExistSelectedKeysWithOverflow) {
249296
AssertChunkedArrayEquals(expected, actual);
250297
}
251298

252-
TEST_F(SharedShreddingFileReaderTest, TestPartialExistSelectedKeys) {
253-
auto reader = CreateReader();
299+
TEST_F(MapSharedShreddingFileReaderTest, TestPartialExistSelectedKeys) {
300+
auto reader = CreateReader(/*physical_array=*/nullptr, /*physical_schema=*/nullptr,
301+
/*selected_keys=*/"a,c,missing");
254302
auto read_schema = ExportSchema(ReadSchema("a,c,missing"));
255303
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
256304
/*selection_bitmap=*/std::nullopt));
@@ -270,15 +318,7 @@ TEST_F(SharedShreddingFileReaderTest, TestPartialExistSelectedKeys) {
270318
AssertChunkedArrayEquals(expected, actual);
271319
}
272320

273-
TEST_F(SharedShreddingFileReaderTest, TestDuplicatedSelectedKeys) {
274-
auto reader = CreateReader();
275-
auto read_schema = ExportSchema(ReadSchema("a,c,a"));
276-
ASSERT_NOK_WITH_MSG(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
277-
/*selection_bitmap=*/std::nullopt),
278-
"duplicate key [a] in paimon.map.selected-keys for field tags");
279-
}
280-
281-
TEST_F(SharedShreddingFileReaderTest, TestMissingSelectedKeysReadsWholeMap) {
321+
TEST_F(MapSharedShreddingFileReaderTest, TestMissingSelectedKeysReadsWholeMap) {
282322
auto reader = CreateReader();
283323
auto read_schema = ExportSchema(ReadSchema(std::nullopt));
284324
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
@@ -298,7 +338,7 @@ TEST_F(SharedShreddingFileReaderTest, TestMissingSelectedKeysReadsWholeMap) {
298338
AssertChunkedArrayEquals(expected, actual);
299339
}
300340

301-
TEST_F(SharedShreddingFileReaderTest, TestSpecialSelectedKeys) {
341+
TEST_F(MapSharedShreddingFileReaderTest, TestSpecialSelectedKeys) {
302342
MapSharedShreddingFieldMeta meta;
303343
meta.name_to_id = {{"", 0}, {" ", 1}, {".", 2}, {"a", 3}};
304344
meta.field_to_columns = {{0, {0}}, {1, {1}}, {2, {0}}, {3, {1}}};
@@ -315,7 +355,7 @@ TEST_F(SharedShreddingFileReaderTest, TestSpecialSelectedKeys) {
315355
.ValueOrDie();
316356

317357
auto assert_read = [&](const std::string& selected_keys, const std::string& expected_json) {
318-
auto reader = CreateReader(physical_array, physical_schema);
358+
auto reader = CreateReader(physical_array, physical_schema, selected_keys);
319359
auto read_schema = ExportSchema(ReadSchema(selected_keys));
320360
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
321361
/*selection_bitmap=*/std::nullopt));
@@ -350,18 +390,9 @@ TEST_F(SharedShreddingFileReaderTest, TestSpecialSelectedKeys) {
350390
])");
351391
}
352392

353-
TEST_F(SharedShreddingFileReaderTest, TestSpecialSelectedKeysWithDuplicatedEmptyKey) {
354-
for (const auto& selected_keys : {",", ",,"}) {
355-
auto reader = CreateReader();
356-
auto read_schema = ExportSchema(ReadSchema(selected_keys));
357-
ASSERT_NOK_WITH_MSG(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
358-
/*selection_bitmap=*/std::nullopt),
359-
"duplicate key [] in paimon.map.selected-keys for field tags");
360-
}
361-
}
362-
363-
TEST_F(SharedShreddingFileReaderTest, TestUnknownSelectedKeyReturnsEmptyMap) {
364-
auto reader = CreateReader();
393+
TEST_F(MapSharedShreddingFileReaderTest, TestUnknownSelectedKeyReturnsEmptyMap) {
394+
auto reader = CreateReader(/*physical_array=*/nullptr, /*physical_schema=*/nullptr,
395+
/*selected_keys=*/"missing");
365396
auto read_schema = ExportSchema(ReadSchema("missing"));
366397
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
367398
/*selection_bitmap=*/std::nullopt));
@@ -381,39 +412,39 @@ TEST_F(SharedShreddingFileReaderTest, TestUnknownSelectedKeyReturnsEmptyMap) {
381412
AssertChunkedArrayEquals(expected, actual);
382413
}
383414

384-
TEST_F(SharedShreddingFileReaderTest, TestInvalidNullFieldMappingField) {
415+
TEST_F(MapSharedShreddingFileReaderTest, TestInvalidNullFieldMappingField) {
385416
auto physical_schema = PhysicalSchemaWithMetadata();
386417
std::string json = R"([
387418
[1, [null, 10, null, null]]
388419
])";
389420
auto physical_array =
390421
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(physical_schema->fields()), json)
391422
.ValueOrDie();
392-
auto reader = CreateReader(physical_array, physical_schema);
423+
auto reader = CreateReader(physical_array, physical_schema, /*selected_keys=*/"a");
393424
auto read_schema = ExportSchema(ReadSchema("a"));
394425
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
395426
/*selection_bitmap=*/std::nullopt));
396427
ASSERT_NOK_WITH_MSG(ReadResultCollector::CollectResult(reader.get()),
397428
"__field_mapping cannot be null");
398429
}
399430

400-
TEST_F(SharedShreddingFileReaderTest, TestInvalidNullFieldMappingFieldElement) {
431+
TEST_F(MapSharedShreddingFileReaderTest, TestInvalidNullFieldMappingFieldElement) {
401432
auto physical_schema = PhysicalSchemaWithMetadata();
402433
std::string json = R"([
403434
[1, [[0, null], 10, null, null]]
404435
])";
405436
auto physical_array =
406437
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(physical_schema->fields()), json)
407438
.ValueOrDie();
408-
auto reader = CreateReader(physical_array, physical_schema);
439+
auto reader = CreateReader(physical_array, physical_schema, /*selected_keys=*/"b");
409440
auto read_schema = ExportSchema(ReadSchema("b"));
410441
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
411442
/*selection_bitmap=*/std::nullopt));
412443
ASSERT_NOK_WITH_MSG(ReadResultCollector::CollectResult(reader.get()),
413444
"__field_mapping element cannot be null");
414445
}
415446

416-
TEST_F(SharedShreddingFileReaderTest, TestListValue) {
447+
TEST_F(MapSharedShreddingFileReaderTest, TestListValue) {
417448
std::shared_ptr<arrow::Schema> logical_schema = arrow::schema({
418449
arrow::field("id", arrow::int32()),
419450
arrow::field("tags", arrow::map(arrow::utf8(), arrow::list(arrow::int32()))),
@@ -443,7 +474,8 @@ TEST_F(SharedShreddingFileReaderTest, TestListValue) {
443474
[4, [[1, 0], [8], [9, 10], [[2, [null]]]]]
444475
])")
445476
.ValueOrDie();
446-
auto reader = CreateReader(physical_array, physical_schema);
477+
auto reader = CreateReader(physical_array, physical_schema,
478+
/*selected_keys=*/"a,c"); // NOLINT(whitespace/comma)
447479

448480
auto read_metadata = std::make_shared<arrow::KeyValueMetadata>();
449481
read_metadata->Append("paimon.map.selected-keys", "a,c");
@@ -467,7 +499,7 @@ TEST_F(SharedShreddingFileReaderTest, TestListValue) {
467499
AssertChunkedArrayEquals(expected, actual);
468500
}
469501

470-
TEST_F(SharedShreddingFileReaderTest, TestOrcDictionaryEncodedStringValue) {
502+
TEST_F(MapSharedShreddingFileReaderTest, TestOrcDictionaryEncodedStringValue) {
471503
std::shared_ptr<arrow::Schema> logical_schema = arrow::schema({
472504
arrow::field("id", arrow::int32()),
473505
arrow::field("tags", arrow::map(arrow::utf8(), arrow::utf8())),
@@ -503,9 +535,8 @@ TEST_F(SharedShreddingFileReaderTest, TestOrcDictionaryEncodedStringValue) {
503535
std::string data_file_path =
504536
path_factory->ToPath(inc.GetNewFilesIncrement().NewFiles()[0]->file_name);
505537
std::map<std::string, std::string> reader_options = {{"orc.read.enable-lazy-decoding", "true"}};
506-
ASSERT_OK_AND_ASSIGN(auto reader,
507-
SharedShreddingFileReader::Create(
508-
OpenFormatReader(data_file_path, format, reader_options), pool_));
538+
auto reader = WrapReader(OpenFormatReader(data_file_path, format, reader_options),
539+
/*selected_keys_str=*/"a,c");
509540

510541
auto read_metadata = std::make_shared<arrow::KeyValueMetadata>();
511542
read_metadata->Append("paimon.map.selected-keys", "a,c");
@@ -528,7 +559,7 @@ TEST_F(SharedShreddingFileReaderTest, TestOrcDictionaryEncodedStringValue) {
528559
AssertChunkedArrayEquals(expected, actual);
529560
}
530561

531-
TEST_F(SharedShreddingFileReaderTest, TestReadsRealFormatFile) {
562+
TEST_F(MapSharedShreddingFileReaderTest, TestReadsRealFormatFile) {
532563
// TODO(lisizhuo.lsz): support other format
533564
auto options = options_;
534565
std::string format = "orc";
@@ -559,8 +590,8 @@ TEST_F(SharedShreddingFileReaderTest, TestReadsRealFormatFile) {
559590

560591
std::string data_file_path =
561592
path_factory->ToPath(inc.GetNewFilesIncrement().NewFiles()[0]->file_name);
562-
ASSERT_OK_AND_ASSIGN(auto reader, SharedShreddingFileReader::Create(
563-
OpenFormatReader(data_file_path, format), pool_));
593+
auto reader = WrapReader(OpenFormatReader(data_file_path, format),
594+
/*selected_keys_str=*/"a,c");
564595

565596
auto read_schema = ExportSchema(ReadSchema("a,c"));
566597
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,

0 commit comments

Comments
 (0)