Skip to content

Commit 6f07d34

Browse files
feat: Add manifests and files system tables (alibaba#309)
1 parent e547186 commit 6f07d34

8 files changed

Lines changed: 884 additions & 8 deletions

src/paimon/common/utils/binary_row_partition_computer.cpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -136,13 +136,13 @@ Result<arrow::Type::type> BinaryRowPartitionComputer::GetTypeFromArrowSchema(
136136

137137
Result<std::string> BinaryRowPartitionComputer::PartToSimpleString(
138138
const std::shared_ptr<arrow::Schema>& partition_type, const BinaryRow& partition,
139-
const std::string& delimiter, int32_t max_length) {
139+
const std::string& delimiter, int32_t max_length, bool legacy_partition_name_enabled) {
140140
std::vector<DataConverterUtils::BinaryRowFieldToStrConverter> partition_converters;
141141
partition_converters.reserve(partition_type->num_fields());
142142
for (const auto& field : partition_type->fields()) {
143143
PAIMON_ASSIGN_OR_RAISE(DataConverterUtils::BinaryRowFieldToStrConverter converter,
144144
DataConverterUtils::CreateBinaryRowFieldToStringConverter(
145-
field->type()->id(), /*legacy_partition_name_enabled=*/true));
145+
field->type()->id(), legacy_partition_name_enabled));
146146
partition_converters.emplace_back(converter);
147147
}
148148
std::vector<std::string> partition_vec;

src/paimon/common/utils/binary_row_partition_computer.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,8 @@ class BinaryRowPartitionComputer {
6060

6161
static Result<std::string> PartToSimpleString(
6262
const std::shared_ptr<arrow::Schema>& partition_type, const BinaryRow& partition,
63-
const std::string& delimiter, int32_t max_length);
63+
const std::string& delimiter, int32_t max_length,
64+
bool legacy_partition_name_enabled = true);
6465

6566
private:
6667
BinaryRowPartitionComputer(const std::vector<std::string>& partition_keys,

src/paimon/core/catalog/file_system_catalog_test.cpp

Lines changed: 51 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -300,8 +300,8 @@ TEST(FileSystemCatalogTest, TestMetadataSystemTableCatalog) {
300300
/*ignore_if_exists=*/false));
301301
ArrowSchemaRelease(&schema);
302302

303-
std::vector<std::string> metadata_tables = {"snapshots", "schemas", "tags", "branches",
304-
"consumers"};
303+
std::vector<std::string> metadata_tables = {"snapshots", "schemas", "tags", "branches",
304+
"consumers", "manifests", "files"};
305305
for (const auto& table_name : metadata_tables) {
306306
Identifier system_identifier("db1", "tbl1$" + table_name);
307307
ASSERT_OK_AND_ASSIGN(bool exists, catalog.TableExists(system_identifier));
@@ -363,6 +363,55 @@ TEST(FileSystemCatalogTest, TestMetadataSystemTableCatalog) {
363363
(std::vector<std::string>{"consumer_id", "next_snapshot_id"}));
364364
ASSERT_FALSE(consumers_arrow_schema->field(1)->nullable());
365365

366+
ASSERT_OK_AND_ASSIGN(std::shared_ptr<Schema> manifests_schema,
367+
catalog.LoadTableSchema(Identifier("db1", "tbl1$manifests")));
368+
ASSERT_OK_AND_ASSIGN(auto manifests_c_schema, manifests_schema->GetArrowSchema());
369+
auto manifests_arrow_schema = arrow::ImportSchema(manifests_c_schema.get()).ValueUnsafe();
370+
ASSERT_EQ(manifests_arrow_schema->field_names(),
371+
(std::vector<std::string>{"file_name", "file_size", "num_added_files",
372+
"num_deleted_files", "schema_id", "min_partition_stats",
373+
"max_partition_stats", "min_row_id", "max_row_id"}));
374+
ASSERT_FALSE(manifests_arrow_schema->field(0)->nullable());
375+
ASSERT_EQ(manifests_arrow_schema->field(1)->type()->id(), arrow::Type::INT64);
376+
ASSERT_FALSE(manifests_arrow_schema->field(4)->nullable());
377+
ASSERT_TRUE(manifests_arrow_schema->field(5)->nullable());
378+
ASSERT_TRUE(manifests_arrow_schema->field(8)->nullable());
379+
380+
ASSERT_OK_AND_ASSIGN(std::shared_ptr<Schema> files_schema,
381+
catalog.LoadTableSchema(Identifier("db1", "tbl1$files")));
382+
ASSERT_OK_AND_ASSIGN(auto files_c_schema, files_schema->GetArrowSchema());
383+
auto files_arrow_schema = arrow::ImportSchema(files_c_schema.get()).ValueUnsafe();
384+
ASSERT_EQ(files_arrow_schema->field_names(), (std::vector<std::string>{"partition",
385+
"bucket",
386+
"file_path",
387+
"file_format",
388+
"schema_id",
389+
"level",
390+
"record_count",
391+
"file_size_in_bytes",
392+
"min_key",
393+
"max_key",
394+
"null_value_counts",
395+
"min_value_stats",
396+
"max_value_stats",
397+
"min_sequence_number",
398+
"max_sequence_number",
399+
"creation_time",
400+
"deleteRowCount",
401+
"file_source",
402+
"first_row_id",
403+
"write_cols"}));
404+
ASSERT_TRUE(files_arrow_schema->field(0)->nullable());
405+
ASSERT_FALSE(files_arrow_schema->field(1)->nullable());
406+
ASSERT_FALSE(files_arrow_schema->field(2)->nullable());
407+
ASSERT_FALSE(files_arrow_schema->field(10)->nullable());
408+
ASSERT_EQ(files_arrow_schema->field(15)->type()->id(), arrow::Type::TIMESTAMP);
409+
ASSERT_EQ(files_arrow_schema->field(19)->type()->id(), arrow::Type::LIST);
410+
auto write_cols_type =
411+
std::dynamic_pointer_cast<arrow::ListType>(files_arrow_schema->field(19)->type());
412+
ASSERT_TRUE(write_cols_type);
413+
ASSERT_EQ(write_cols_type->value_type()->id(), arrow::Type::STRING);
414+
366415
Identifier snapshots_identifier("db1", "tbl1$snapshots");
367416
::ArrowSchema system_create_schema;
368417
ASSERT_TRUE(arrow::ExportSchema(*typed_schema, &system_create_schema).ok());

src/paimon/core/table/system/in_memory_system_table.cpp

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,9 @@ class InMemorySystemTableBatchReader : public BatchReader {
4646
emitted_ = true;
4747
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Schema> schema, table_->ArrowSchema());
4848
PAIMON_ASSIGN_OR_RAISE(std::vector<GenericRow> rows, table_->BuildRows());
49+
if (rows.empty()) {
50+
return BatchReader::MakeEofBatch();
51+
}
4952
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<GenericRowToArrowArrayConverter> converter,
5053
GenericRowToArrowArrayConverter::Create(schema, arrow_pool_.get()));
5154
return converter->NextBatch(rows);

0 commit comments

Comments
 (0)