Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
312 changes: 312 additions & 0 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/AvroSchema.h
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,96 @@ static constexpr const char * manifest_list_v2_schema = R"(
}
)";

static constexpr const char * manifest_list_v3_schema = R"(
{
"type": "record",
"name": "manifest_file",
"fields": [
{"name": "manifest_path", "type": "string", "doc": "Location URI with FS scheme", "field-id": 500},
{"name": "manifest_length", "type": "long", "doc": "Total file size in bytes", "field-id": 501},
{"name": "partition_spec_id", "type": "int", "doc": "Spec ID used to write", "field-id": 502},
{"name": "content", "type": "int", "doc": "Contents of the manifest: 0=data, 1=deletes", "field-id": 517},
{
"name": "sequence_number",
"type": "long",
"doc": "Sequence number when the manifest was added",
"field-id": 515
},
{
"name": "min_sequence_number",
"type": "long",
"doc": "Lowest sequence number in the manifest",
"field-id": 516
},
{"name": "added_snapshot_id", "type": "long", "doc": "Snapshot ID that added the manifest", "field-id": 503},
{"name": "added_files_count", "type": "int", "doc": "Added entry count", "field-id": 504},
{"name": "existing_files_count", "type": "int", "doc": "Existing entry count", "field-id": 505},
{"name": "deleted_files_count", "type": "int", "doc": "Deleted entry count", "field-id": 506},
{"name": "added_rows_count", "type": "long", "doc": "Added rows count", "field-id": 512},
{"name": "existing_rows_count", "type": "long", "doc": "Existing rows count", "field-id": 513},
{"name": "deleted_rows_count", "type": "long", "doc": "Deleted rows count", "field-id": 514},
{
"name": "partitions",
"type": [
"null",
{
"type": "array",
"items":
{
"type": "record",
"name": "r508",
"fields": [
{
"name": "contains_null",
"type": "boolean",
"doc": "True if any file has a null partition value",
"field-id": 509
},
{
"name": "contains_nan",
"type": ["null", "boolean"],
"doc": "True if any file has a nan partition value",
"field-id": 518
},
{
"name": "lower_bound",
"type": ["null", "bytes"],
"doc": "Partition lower bound for all files",
"field-id": 510
},
{
"name": "upper_bound",
"type": ["null", "bytes"],
"doc": "Partition upper bound for all files",
"field-id": 511
}
]
},
"element-id": 508
}
],
"doc": "Summary for each partition",
"field-id": 507
},
{
"name": "key_metadata",
"type": ["null", "bytes"],
"doc": "Encryption key metadata blob",
"default": null,
"field-id": 519
},
{
"name": "first_row_id",
"type": ["null", "long"],
"doc": "Starting row ID to assign to new rows in ADDED data files",
"default": null,
"field-id": 520
}
]
}
)";


/// NOTE: This string is just a template for the actual schema. To use it, you must first replace "#" with the correct value.
static constexpr const char * manifest_entry_v1_schema = R"(
{
Expand Down Expand Up @@ -619,4 +709,226 @@ static constexpr const char * data_file_sidecar_schema = R"(
}
)";

static constexpr const char * manifest_entry_v3_schema = R"(
{
"type": "record",
"name": "manifest_entry",
"fields": [
{"name": "status", "type": "int", "field-id": 0},
{"name": "snapshot_id", "type": ["null", "long"], "field-id": 1},
{"name": "sequence_number", "type": ["null", "long"], "field-id": 3},
{"name": "file_sequence_number", "type": ["null", "long"], "field-id": 4},
{
"name": "data_file",
"type":
{
"type": "record",
"name": "r2",
"fields": [
{"name": "content", "type": "int", "doc": "Type of content stored by the data file", "field-id": 134},
{"name": "file_path", "type": "string", "doc": "Location URI with FS scheme", "field-id": 100},
{
"name": "file_format",
"type": "string",
"doc": "File format name: avro, orc, or parquet",
"field-id": 101
},
{
"name": "partition",
"type":
{
"type": "record",
"name": "r102",
"fields": #
},
"field-id": 102
},
{"name": "record_count", "type": "long", "doc": "Number of records in the file", "field-id": 103},
{"name": "file_size_in_bytes", "type": "long", "doc": "Total file size in bytes", "field-id": 104},
{
"name": "column_sizes",
"type": [
"null",
{
"type": "array",
"items":
{
"type": "record",
"name": "k117_v118",
"fields": [
{"name": "key", "type": "int", "field-id": 117},
{"name": "value", "type": "long", "field-id": 118}
]
},
"logicalType": "map"
}
],
"doc": "Map of column id to total size on disk",
"field-id": 108
},
{
"name": "value_counts",
"type": [
"null",
{
"type": "array",
"items":
{
"type": "record",
"name": "k119_v120",
"fields": [
{"name": "key", "type": "int", "field-id": 119},
{"name": "value", "type": "long", "field-id": 120}
]
},
"logicalType": "map"
}
],
"doc": "Map of column id to total count, including null and NaN",
"field-id": 109
},
{
"name": "null_value_counts",
"type": [
"null",
{
"type": "array",
"items":
{
"type": "record",
"name": "k121_v122",
"fields": [
{"name": "key", "type": "int", "field-id": 121},
{"name": "value", "type": "long", "field-id": 122}
]
},
"logicalType": "map"
}
],
"doc": "Map of column id to null value count",
"field-id": 110
},
{
"name": "nan_value_counts",
"type": [
"null",
{
"type": "array",
"items":
{
"type": "record",
"name": "k138_v139",
"fields": [
{"name": "key", "type": "int", "field-id": 138},
{"name": "value", "type": "long", "field-id": 139}
]
},
"logicalType": "map"
}
],
"doc": "Map of column id to number of NaN values in the column",
"field-id": 137
},
{
"name": "lower_bounds",
"type": [
"null",
{
"type": "array",
"items":
{
"type": "record",
"name": "k126_v127",
"fields": [
{"name": "key", "type": "int", "field-id": 126},
{"name": "value", "type": "bytes", "field-id": 127}
]
},
"logicalType": "map"
}
],
"doc": "Map of column id to lower bound",
"field-id": 125
},
{
"name": "upper_bounds",
"type": [
"null",
{
"type": "array",
"items":
{
"type": "record",
"name": "k129_v130",
"fields": [
{"name": "key", "type": "int", "field-id": 129},
{"name": "value", "type": "bytes", "field-id": 130}
]
},
"logicalType": "map"
}
],
"doc": "Map of column id to upper bound",
"field-id": 128
},
{
"name": "key_metadata",
"type": ["null", "bytes"],
"doc": "Encryption key metadata blob",
"field-id": 131
},
{
"name": "split_offsets",
"type": ["null", {"type": "array", "items": "long", "element-id": 133}],
"doc": "Splittable offsets",
"field-id": 132
},
{
"name": "equality_ids",
"type": ["null", {"type": "array", "items": "int", "element-id": 136}],
"doc": "Field ids used to determine row equality for delete files",
"field-id": 135
},
{
"name": "sort_order_id",
"type": ["null", "int"],
"doc": "Sort order ID",
"field-id": 140
},
{
"name": "first_row_id",
"type": ["null", "long"],
"doc": "The _row_id for the first row in the data file",
"default": null,
"field-id": 142
},
{
"name": "referenced_data_file",
"type": ["null", "string"],
"doc": "Fully qualified location (URI with FS scheme) of a data file that all deletes reference",
"default": null,
"field-id": 143
},
{
"name": "content_offset",
"type": ["null", "long"],
"doc": "The offset in the file where the content starts",
"default": null,
"field-id": 144
},
{
"name": "content_size_in_bytes",
"type": ["null", "long"],
"doc": "The length of referenced content stored in the file",
"default": null,
"field-id": 145
}
]
},
"field-id": 2
}
]
}
)";

}
17 changes: 16 additions & 1 deletion src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1041,6 +1041,7 @@ static void writeMetadataFiles(
Poco::JSON::Object::Ptr initial_metadata_object = plan.initial_metadata_object;
std::unordered_map<Iceberg::IcebergPathFromMetadata, Iceberg::IcebergPathFromMetadata> manifest_file_renamings;
std::unordered_map<Iceberg::IcebergPathFromMetadata, Int64> manifest_file_sizes;
std::unordered_map<Iceberg::IcebergPathFromMetadata, Int64> manifest_file_row_counts;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't see this change (and other in this file, in Constant.h, and in IcebergMetadata.cpp) in original PR (https://github.com/ClickHouse/ClickHouse/pull/116171/changes) as well as don't see it in current upstream master (https://github.com/ClickHouse/ClickHouse/blame/master/src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp#L1107).
From where are these additional changes?


{
std::unordered_map<std::shared_ptr<ManifestFilePlan>, std::unordered_set<Iceberg::IcebergPathFromMetadata>> grouped_by_manifest_files_result;
Expand Down Expand Up @@ -1114,6 +1115,10 @@ static void writeMetadataFiles(
file_byte_counts.push_back(0);
}
}
Int64 manifest_rows = 0;
for (const auto rows : file_row_counts)
manifest_rows += static_cast<Int64>(rows);
manifest_file_row_counts[manifest_entry->patched_path] = manifest_rows;
generateManifestFile(
metadata_object,
partition_columns,
Expand Down Expand Up @@ -1170,8 +1175,12 @@ static void writeMetadataFiles(
}
}
std::vector<Int64> per_manifest_sizes;
std::vector<Int64> per_manifest_row_counts;
for (const auto & entry : renamed_manifest_entries)
{
per_manifest_sizes.push_back(manifest_file_sizes[entry]);
per_manifest_row_counts.push_back(manifest_file_row_counts[entry]);
}
auto buffer_manifest_list = object_storage->writeObject(
StoredObject(path_resolver.resolve(renamed_manifest_list)),
WriteMode::Rewrite,
Expand All @@ -1189,7 +1198,13 @@ static void writeMetadataFiles(
per_manifest_sizes,
*buffer_manifest_list,
Iceberg::FileContentType::DATA,
false);
false,
/* per_entry_content_types */ {},
/* existing_entry_counts */ {},
/* carry_forward_manifest_paths */ {},
/* entry_partition_spec_ids */ {},
/* entry_partition_summaries */ {},
per_manifest_row_counts);
buffer_manifest_list->finalize();
}

Expand Down
2 changes: 2 additions & 0 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ DEFINE_ICEBERG_FIELD(name);
DEFINE_ICEBERG_FIELD(required);
DEFINE_ICEBERG_FIELD(schema);
DEFINE_ICEBERG_FIELD(schemas);
/// The Avro file-header metadata key that stores the writer schema (Avro spec, "avro.schema").
DEFINE_ICEBERG_FIELD_ALIAS(avro_schema, avro.schema);
DEFINE_ICEBERG_FIELD(sequence_number);
DEFINE_ICEBERG_FIELD(snapshots);
DEFINE_ICEBERG_FIELD(status);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2005,7 +2005,13 @@ std::optional<IStorage::ExportPartitionCommitInfo> IcebergMetadata::commitImport
try
{
generateManifestList(
resolver, metadata, object_storage, *secondary_storages, context, {manifest_entry_path}, new_snapshot, {manifest_lengths}, *buffer_manifest_list, Iceberg::FileContentType::DATA, true);
resolver, metadata, object_storage, *secondary_storages, context, {manifest_entry_path}, new_snapshot, {manifest_lengths}, *buffer_manifest_list, Iceberg::FileContentType::DATA, true,
/* per_entry_content_types */ {},
/* existing_entry_counts */ {},
/* carry_forward_manifest_paths */ {},
/* entry_partition_spec_ids */ {},
/* entry_partition_summaries */ {},
/* entry_row_counts */ {total_rows});
buffer_manifest_list->finalize();
}
catch (...)
Expand Down
Loading
Loading