Skip to content

Commit 0e16cf3

Browse files
committed
Address shared-shredding schema utils review comments
1 parent 26d59f5 commit 0e16cf3

3 files changed

Lines changed: 93 additions & 17 deletions

File tree

include/paimon/data/shredding/map_shared_shredding_schema_utils.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ namespace paimon {
3434
struct PAIMON_EXPORT MapSharedShreddingFieldMeta {
3535
/// field_name -> field_id
3636
std::map<std::string, int32_t> name_to_id;
37-
/// field_id -> set of physical column indices S
37+
/// field_id -> ordered physical column indices
3838
std::map<int32_t, std::vector<int32_t>> field_to_columns;
3939
/// Set of field_ids that ever spilled into __overflow
4040
std::set<int32_t> overflow_field_set;
@@ -73,6 +73,7 @@ class PAIMON_EXPORT MapSharedShreddingSchemaUtils {
7373
/// @param physical_schema The Arrow C physical schema whose fields should receive metadata.
7474
/// Ownership of schema resources is transferred to this method.
7575
/// @param field_name_to_meta Map from physical field name to its shared-shredding metadata.
76+
/// Existing shared-shredding metadata keys on matching fields are overwritten.
7677
/// @param compression Compression codec name for field_dict serialization.
7778
/// @return A new Arrow C schema with shared-shredding metadata attached to matching fields.
7879
static Result<std::unique_ptr<::ArrowSchema>> AttachMetadataToSchema(

src/paimon/common/data/shredding/map_shared_shredding_schema_utils_test.cpp

Lines changed: 72 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,11 @@ TEST(MapSharedShreddingSchemaUtilsTest, AttachMetadataToSchemaBasic) {
6666
ASSERT_EQ(deserialized, tags_meta);
6767
}
6868

69-
TEST(MapSharedShreddingSchemaUtilsTest, AttachMetadataToSchemaMissingField) {
69+
TEST(MapSharedShreddingSchemaUtilsTest, AttachMetadataToSchemaInvalidInput) {
70+
ASSERT_NOK_WITH_MSG(MapSharedShreddingSchemaUtils::AttachMetadataToSchema(
71+
std::unique_ptr<::ArrowSchema>(), {}, "none"),
72+
"physical schema is null");
73+
7074
MapSharedShreddingFieldMeta tags_meta;
7175
tags_meta.name_to_id = {{"host", 0}};
7276
tags_meta.field_to_columns = {{0, {0}}};
@@ -114,6 +118,51 @@ TEST(MapSharedShreddingSchemaUtilsTest, AttachMetadataToSchemaPreservesExistingF
114118
ASSERT_TRUE(MapSharedShreddingUtils::HasShreddingMetadata(updated_metadata));
115119
}
116120

121+
TEST(MapSharedShreddingSchemaUtilsTest, AttachMetadataToSchemaOverwritesExistingShreddingMetadata) {
122+
MapSharedShreddingFieldMeta old_meta;
123+
old_meta.name_to_id = {{"old", 0}};
124+
old_meta.field_to_columns = {{0, {0}}};
125+
old_meta.num_columns = 1;
126+
old_meta.max_row_width = 1;
127+
128+
MapSharedShreddingFieldMeta tags_meta;
129+
tags_meta.name_to_id = {{"host", 0}, {"region", 1}};
130+
tags_meta.field_to_columns = {{0, {0}}, {1, {1}}};
131+
tags_meta.overflow_field_set = {1};
132+
tags_meta.num_columns = 2;
133+
tags_meta.max_row_width = 2;
134+
135+
auto tags_metadata = std::make_shared<arrow::KeyValueMetadata>();
136+
tags_metadata->Append("paimon.field.id", "7");
137+
ASSERT_OK(MapSharedShreddingUtils::SerializeMetadata(old_meta, "none", tags_metadata.get()));
138+
auto schema = arrow::schema({arrow::field(
139+
"tags",
140+
arrow::struct_({arrow::field("__field_mapping", arrow::list(arrow::int32()), true),
141+
arrow::field("__col_0", arrow::utf8(), true),
142+
arrow::field("__col_1", arrow::utf8(), true)}),
143+
true, tags_metadata)});
144+
145+
auto c_schema = std::make_unique<::ArrowSchema>();
146+
ASSERT_TRUE(arrow::ExportSchema(*schema, c_schema.get()).ok());
147+
ASSERT_OK_AND_ASSIGN(auto c_updated_schema,
148+
MapSharedShreddingSchemaUtils::AttachMetadataToSchema(
149+
std::move(c_schema), {{"tags", tags_meta}}, "none"));
150+
auto updated_schema = arrow::ImportSchema(c_updated_schema.get()).ValueOrDie();
151+
auto updated_metadata = updated_schema->field(0)->metadata()->Copy();
152+
153+
ASSERT_EQ(updated_metadata->value(updated_metadata->FindKey("paimon.field.id")), "7");
154+
int32_t storage_layout_key_count = 0;
155+
for (const auto& key : updated_metadata->keys()) {
156+
if (key == MapShreddingDefine::kStorageLayout) {
157+
++storage_layout_key_count;
158+
}
159+
}
160+
ASSERT_EQ(storage_layout_key_count, 1);
161+
ASSERT_OK_AND_ASSIGN(auto deserialized,
162+
MapSharedShreddingUtils::DeserializeMetadata(updated_metadata, "none"));
163+
ASSERT_EQ(deserialized, tags_meta);
164+
}
165+
117166
TEST(MapSharedShreddingSchemaUtilsTest, ExtractMetadataFromField) {
118167
MapSharedShreddingFieldMeta tags_meta;
119168
tags_meta.name_to_id = {{"host", 0}, {"region", 1}};
@@ -150,10 +199,30 @@ TEST(MapSharedShreddingSchemaUtilsTest, ExtractMetadataFromFieldNoShreddingMetad
150199
"metadata is null or storage layout is not shared-shredding");
151200
}
152201

153-
TEST(MapSharedShreddingSchemaUtilsTest, LogicalToPhysicalSchemaNonMapField) {
202+
TEST(MapSharedShreddingSchemaUtilsTest, ExtractMetadataFromFieldInvalidInput) {
203+
ASSERT_NOK_WITH_MSG(MapSharedShreddingSchemaUtils::ExtractMetadataFromField(
204+
std::unique_ptr<::ArrowSchema>(), "tags", "none"),
205+
"physical schema is null");
206+
207+
auto schema = arrow::schema({arrow::field("id", arrow::int32())});
208+
209+
auto c_schema = std::make_unique<::ArrowSchema>();
210+
ASSERT_TRUE(arrow::ExportSchema(*schema, c_schema.get()).ok());
211+
212+
ASSERT_NOK_WITH_MSG(MapSharedShreddingSchemaUtils::ExtractMetadataFromField(std::move(c_schema),
213+
"tags", "none"),
214+
"Shared-shredding field 'tags' not found in physical schema.");
215+
}
216+
217+
TEST(MapSharedShreddingSchemaUtilsTest, LogicalToPhysicalSchemaInvalidInput) {
218+
std::map<std::string, int32_t> field_to_num_columns = {{"tags", 2}};
219+
220+
ASSERT_NOK_WITH_MSG(MapSharedShreddingSchemaUtils::LogicalToPhysicalSchema(
221+
std::unique_ptr<::ArrowSchema>(), field_to_num_columns),
222+
"logical schema is null");
223+
154224
auto schema =
155225
arrow::schema({arrow::field("ts", arrow::int64()), arrow::field("tags", arrow::int32())});
156-
std::map<std::string, int32_t> field_to_num_columns = {{"tags", 2}};
157226

158227
auto c_schema = std::make_unique<::ArrowSchema>();
159228
ASSERT_TRUE(arrow::ExportSchema(*schema, c_schema.get()).ok());

src/paimon/common/data/shredding/map_shared_shredding_utils.cpp

Lines changed: 19 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
#include "paimon/common/compression/block_decompressor.h"
2828
#include "paimon/common/data/shredding/map_shared_shredding_batch_converter.h"
2929
#include "paimon/common/data/shredding/map_shared_shredding_context.h"
30+
#include "paimon/common/utils/arrow/status_utils.h"
3031
#include "paimon/common/utils/string_utils.h"
3132
#include "paimon/core/core_options.h"
3233
#include "paimon/core/options/map_storage_layout.h"
@@ -357,23 +358,28 @@ Result<std::set<int32_t>> DeserializeOverflowSet(const std::string& json_str) {
357358
Status MapSharedShreddingUtils::SerializeMetadata(const MapSharedShreddingFieldMeta& field_meta,
358359
const std::string& compression,
359360
arrow::KeyValueMetadata* metadata) {
360-
metadata->Append(MapShreddingDefine::kStorageLayout,
361-
MapShreddingDefine::kStorageLayoutSharedShredding);
362-
metadata->Append(MapSharedShreddingDefine::kVersion,
363-
std::to_string(MapSharedShreddingDefine::kCurrentVersion));
361+
PAIMON_RETURN_NOT_OK_FROM_ARROW(metadata->Set(
362+
MapShreddingDefine::kStorageLayout, MapShreddingDefine::kStorageLayoutSharedShredding));
363+
PAIMON_RETURN_NOT_OK_FROM_ARROW(
364+
metadata->Set(MapSharedShreddingDefine::kVersion,
365+
std::to_string(MapSharedShreddingDefine::kCurrentVersion)));
364366

365367
std::string field_dict_json = SerializeFieldDict(field_meta);
366-
metadata->Append(MapSharedShreddingDefine::kFieldDictOriginalSize,
367-
std::to_string(field_dict_json.size()));
368+
PAIMON_RETURN_NOT_OK_FROM_ARROW(metadata->Set(MapSharedShreddingDefine::kFieldDictOriginalSize,
369+
std::to_string(field_dict_json.size())));
368370
PAIMON_ASSIGN_OR_RAISE(std::string compressed_dict,
369371
CompressString(field_dict_json, compression));
370-
metadata->Append(MapSharedShreddingDefine::kFieldDict, std::move(compressed_dict));
371-
372-
metadata->Append(MapSharedShreddingDefine::kFieldColumns, SerializeFieldColumns(field_meta));
373-
metadata->Append(MapSharedShreddingDefine::kOverflowSet, SerializeOverflowSet(field_meta));
374-
metadata->Append(MapSharedShreddingDefine::kNumColumns, std::to_string(field_meta.num_columns));
375-
metadata->Append(MapSharedShreddingDefine::kMaxRowWidth,
376-
std::to_string(field_meta.max_row_width));
372+
PAIMON_RETURN_NOT_OK_FROM_ARROW(
373+
metadata->Set(MapSharedShreddingDefine::kFieldDict, std::move(compressed_dict)));
374+
375+
PAIMON_RETURN_NOT_OK_FROM_ARROW(
376+
metadata->Set(MapSharedShreddingDefine::kFieldColumns, SerializeFieldColumns(field_meta)));
377+
PAIMON_RETURN_NOT_OK_FROM_ARROW(
378+
metadata->Set(MapSharedShreddingDefine::kOverflowSet, SerializeOverflowSet(field_meta)));
379+
PAIMON_RETURN_NOT_OK_FROM_ARROW(metadata->Set(MapSharedShreddingDefine::kNumColumns,
380+
std::to_string(field_meta.num_columns)));
381+
PAIMON_RETURN_NOT_OK_FROM_ARROW(metadata->Set(MapSharedShreddingDefine::kMaxRowWidth,
382+
std::to_string(field_meta.max_row_width)));
377383

378384
return Status::OK();
379385
}

0 commit comments

Comments
 (0)