diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 3ce753f18..46593f0b5 100644 --- a/src/iceberg/catalog/rest/json_serde.cc +++ b/src/iceberg/catalog/rest/json_serde.cc @@ -105,35 +105,9 @@ constexpr std::string_view kEndSnapshotId = "end-snapshot-id"; constexpr std::string_view kStatsFields = "stats-fields"; constexpr std::string_view kMinRowsRequested = "min-rows-requested"; constexpr std::string_view kPlanTask = "plan-task"; -constexpr std::string_view kContent = "content"; -constexpr std::string_view kContentData = "data"; -constexpr std::string_view kContentPositionDeletes = "position-deletes"; -constexpr std::string_view kContentEqualityDeletes = "equality-deletes"; -constexpr std::string_view kFilePath = "file-path"; -constexpr std::string_view kFileFormat = "file-format"; -constexpr std::string_view kSpecId = "spec-id"; -constexpr std::string_view kPartition = "partition"; -constexpr std::string_view kRecordCount = "record-count"; -constexpr std::string_view kFileSizeInBytes = "file-size-in-bytes"; -constexpr std::string_view kColumnSizes = "column-sizes"; -constexpr std::string_view kValueCounts = "value-counts"; -constexpr std::string_view kNullValueCounts = "null-value-counts"; -constexpr std::string_view kNanValueCounts = "nan-value-counts"; -constexpr std::string_view kLowerBounds = "lower-bounds"; -constexpr std::string_view kUpperBounds = "upper-bounds"; -constexpr std::string_view kKeyMetadata = "key-metadata"; -constexpr std::string_view kSplitOffsets = "split-offsets"; -constexpr std::string_view kEqualityIds = "equality-ids"; -constexpr std::string_view kSortOrderId = "sort-order-id"; -constexpr std::string_view kFirstRowId = "first-row-id"; -constexpr std::string_view kReferencedDataFile = "referenced-data-file"; -constexpr std::string_view kContentOffset = "content-offset"; -constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes"; constexpr std::string_view kDataFile = "data-file"; constexpr std::string_view kDeleteFileReferences = "delete-file-references"; constexpr std::string_view kResidualFilter = "residual-filter"; -constexpr std::string_view kMapKeys = "keys"; -constexpr std::string_view kMapValues = "values"; Result StorageCredentialToJson(const StorageCredential& credential) { ICEBERG_RETURN_UNEXPECTED(credential.Validate()); @@ -152,47 +126,6 @@ Result StorageCredentialFromJson(const nlohmann::json& json) return credential; } -template -Result> KeyValueMapFromJson(const nlohmann::json& json, - std::string_view key) { - std::map result; - if (!json.contains(key) || json.at(key).is_null()) { - return result; - } - - ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue(json, key)); - ICEBERG_ASSIGN_OR_RAISE(auto keys, - GetJsonValue>(map_json, kMapKeys)); - ICEBERG_ASSIGN_OR_RAISE(auto values, - GetJsonValue>(map_json, kMapValues)); - if (keys.size() != values.size()) { - return JsonParseError("'{}' map keys and values have different lengths", key); - } - - for (size_t i = 0; i < keys.size(); ++i) { - result[keys[i]] = std::move(values[i]); - } - return result; -} - -template -void SetKeyValueMap(nlohmann::json& json, std::string_view key, - const std::map& map) { - if (map.empty()) { - return; - } - - std::vector keys; - std::vector values; - keys.reserve(map.size()); - values.reserve(map.size()); - for (const auto& [field_id, value] : map) { - keys.push_back(field_id); - values.push_back(value); - } - json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; -} - } // namespace Result DataFileFromJson( @@ -200,106 +133,7 @@ Result DataFileFromJson( const std::unordered_map>& partition_spec_by_id, const Schema& schema) { - if (!json.is_object()) { - return JsonParseError("DataFile must be a JSON object: {}", SafeDumpJson(json)); - } - DataFile data_file; - - ICEBERG_ASSIGN_OR_RAISE(auto content_str, GetJsonValue(json, kContent)); - if (content_str == kContentData) { - data_file.content = DataFile::Content::kData; - } else if (content_str == kContentPositionDeletes) { - data_file.content = DataFile::Content::kPositionDeletes; - } else if (content_str == kContentEqualityDeletes) { - data_file.content = DataFile::Content::kEqualityDeletes; - } else { - return JsonParseError("Unknown data file content: {}", content_str); - } - - ICEBERG_ASSIGN_OR_RAISE(data_file.file_path, - GetJsonValue(json, kFilePath)); - ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue(json, kFileFormat)); - ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str)); - - ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); - data_file.partition_spec_id = spec_id; - - ICEBERG_ASSIGN_OR_RAISE(auto partition_vals, - GetJsonValue(json, kPartition)); - if (!partition_vals.is_array()) { - return JsonParseError("PartitionValues must be a JSON array: {}", - SafeDumpJson(partition_vals)); - } - std::vector literals; - auto it = partition_spec_by_id.find(spec_id); - if (it == partition_spec_by_id.end()) { - return JsonParseError("Invalid partition spec id: {}", spec_id); - } - ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema)); - auto fields = struct_type->fields(); - if (partition_vals.size() != fields.size()) { - return JsonParseError("Invalid partition data size: expected = {}, actual = {}", - fields.size(), partition_vals.size()); - } - for (size_t pos = 0; pos < fields.size(); ++pos) { - ICEBERG_ASSIGN_OR_RAISE( - auto literal, LiteralFromJson(partition_vals[pos], fields[pos].type().get())); - literals.push_back(std::move(literal)); - } - data_file.partition = PartitionValues(std::move(literals)); - - ICEBERG_ASSIGN_OR_RAISE(data_file.record_count, - GetJsonValue(json, kRecordCount)); - ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes, - GetJsonValue(json, kFileSizeInBytes)); - - ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes, - KeyValueMapFromJson(json, kColumnSizes)); - ICEBERG_ASSIGN_OR_RAISE(data_file.value_counts, - KeyValueMapFromJson(json, kValueCounts)); - ICEBERG_ASSIGN_OR_RAISE(data_file.null_value_counts, - KeyValueMapFromJson(json, kNullValueCounts)); - ICEBERG_ASSIGN_OR_RAISE(data_file.nan_value_counts, - KeyValueMapFromJson(json, kNanValueCounts)); - ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds, - KeyValueMapFromJson>(json, kLowerBounds)); - ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds, - KeyValueMapFromJson>(json, kUpperBounds)); - - if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata, - GetJsonValue>(json, kKeyMetadata)); - } - if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets, - GetJsonValue>(json, kSplitOffsets)); - } - if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids, - GetJsonValue>(json, kEqualityIds)); - } - if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id, - GetJsonValue(json, kSortOrderId)); - } - if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id, - GetJsonValue(json, kFirstRowId)); - } - if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file, - GetJsonValue(json, kReferencedDataFile)); - } - if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset, - GetJsonValue(json, kContentOffset)); - } - if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes, - GetJsonValue(json, kContentSizeInBytes)); - } - - return data_file; + return iceberg::DataFileFromJson(json, partition_spec_by_id, schema); } Result>> FileScanTasksFromJson( @@ -322,7 +156,8 @@ Result>> FileScanTasksFromJson( ICEBERG_ASSIGN_OR_RAISE(auto data_file_json, GetJsonValue(task_json, kDataFile)); ICEBERG_ASSIGN_OR_RAISE( - auto data_file, DataFileFromJson(data_file_json, partition_spec_by_id, schema)); + auto data_file, + iceberg::rest::DataFileFromJson(data_file_json, partition_spec_by_id, schema)); // FIXME: REST scan-task DataFile JSON currently carries first-row-id, // but not the manifest-entry data sequence number. Until the REST API exposes // it, REST-planned tasks cannot inherit _last_updated_sequence_number. @@ -357,97 +192,12 @@ Result>> FileScanTasksFromJson( return file_scan_tasks; } -Result DataFileToJsonUnchecked(const DataFile& data_file) { - nlohmann::json json; - switch (data_file.content) { - case DataFile::Content::kData: - json[kContent] = kContentData; - break; - case DataFile::Content::kPositionDeletes: - json[kContent] = kContentPositionDeletes; - break; - case DataFile::Content::kEqualityDeletes: - json[kContent] = kContentEqualityDeletes; - break; - } - json[kFilePath] = data_file.file_path; - json[kFileFormat] = ToString(data_file.file_format); - - if (!data_file.partition_spec_id.has_value()) { - return ValidationFailed("Cannot serialize REST content file without 'spec-id'"); - } - json[kSpecId] = data_file.partition_spec_id.value(); - - nlohmann::json partition_json = nlohmann::json::array(); - for (const auto& literal : data_file.partition.values()) { - ICEBERG_ASSIGN_OR_RAISE(auto lit_json, iceberg::ToJson(literal)); - partition_json.push_back(std::move(lit_json)); - } - json[kPartition] = std::move(partition_json); - - json[kRecordCount] = data_file.record_count; - json[kFileSizeInBytes] = data_file.file_size_in_bytes; - - SetKeyValueMap(json, kColumnSizes, data_file.column_sizes); - SetKeyValueMap(json, kValueCounts, data_file.value_counts); - SetKeyValueMap(json, kNullValueCounts, data_file.null_value_counts); - SetKeyValueMap(json, kNanValueCounts, data_file.nan_value_counts); - SetKeyValueMap(json, kLowerBounds, data_file.lower_bounds); - SetKeyValueMap(json, kUpperBounds, data_file.upper_bounds); - - if (!data_file.key_metadata.empty()) { - json[kKeyMetadata] = data_file.key_metadata; - } - if (!data_file.split_offsets.empty()) { - json[kSplitOffsets] = data_file.split_offsets; - } - if (!data_file.equality_ids.empty()) { - json[kEqualityIds] = data_file.equality_ids; - } - if (data_file.sort_order_id.has_value()) { - json[kSortOrderId] = data_file.sort_order_id.value(); - } - if (data_file.first_row_id.has_value()) { - json[kFirstRowId] = data_file.first_row_id.value(); - } - if (data_file.referenced_data_file.has_value()) { - json[kReferencedDataFile] = data_file.referenced_data_file.value(); - } - if (data_file.content_offset.has_value()) { - json[kContentOffset] = data_file.content_offset.value(); - } - if (data_file.content_size_in_bytes.has_value()) { - json[kContentSizeInBytes] = data_file.content_size_in_bytes.value(); - } - - return json; -} - Result ToJson( const DataFile& data_file, const std::unordered_map>& partition_specs_by_id, const Schema& schema) { - if (!data_file.partition_spec_id.has_value()) { - return ValidationFailed("Invalid partition spec id from content file: null"); - } - auto it = partition_specs_by_id.find(data_file.partition_spec_id.value()); - if (it == partition_specs_by_id.end() || !it->second) { - return ValidationFailed("Invalid partition spec: null"); - } - if (data_file.partition_spec_id.value() != it->second->spec_id()) { - return ValidationFailed( - "Invalid partition spec id from content file: expected = {}, actual = {}", - it->second->spec_id(), data_file.partition_spec_id.value()); - } - ICEBERG_ASSIGN_OR_RAISE(auto partition_type, it->second->PartitionType(schema)); - if (data_file.partition.num_fields() != partition_type->fields().size()) { - return ValidationFailed( - "Invalid partition data from content file: expected = {}, actual = {}", - partition_type->fields().empty() ? "unpartitioned" : "partitioned", - data_file.partition.num_fields() == 0 ? "unpartitioned" : "partitioned"); - } - return DataFileToJsonUnchecked(data_file); + return iceberg::ToJson(data_file, partition_specs_by_id, schema); } namespace { @@ -495,7 +245,7 @@ Result ScanTaskFieldsToJson( if (task->data_file()) { ICEBERG_ASSIGN_OR_RAISE( auto data_file_json, - ToJson(*task->data_file(), partition_specs_by_id, schema)); + iceberg::rest::ToJson(*task->data_file(), partition_specs_by_id, schema)); task_json[kDataFile] = std::move(data_file_json); } if (!task->delete_files().empty()) { @@ -521,7 +271,8 @@ Result ScanTaskFieldsToJson( } nlohmann::json delete_files_json = nlohmann::json::array(); for (const auto& file : delete_files) { - ICEBERG_ASSIGN_OR_RAISE(auto df_json, ToJson(*file, partition_specs_by_id, schema)); + ICEBERG_ASSIGN_OR_RAISE(auto df_json, + iceberg::rest::ToJson(*file, partition_specs_by_id, schema)); delete_files_json.push_back(std::move(df_json)); } if (!delete_files_json.empty()) { @@ -557,8 +308,9 @@ Status ScanTaskFieldsFromJson( SafeDumpJson(delete_files_json)); } for (const auto& entry_json : delete_files_json) { - ICEBERG_ASSIGN_OR_RAISE(auto delete_file, - DataFileFromJson(entry_json, partition_specs_by_id, schema)); + ICEBERG_ASSIGN_OR_RAISE( + auto delete_file, + iceberg::rest::DataFileFromJson(entry_json, partition_specs_by_id, schema)); response.delete_files.push_back(std::make_shared(std::move(delete_file))); } diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index f322996ec..1f5caf8f0 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -20,6 +20,7 @@ #include #include #include +#include #include #include #include @@ -29,12 +30,16 @@ #include "iceberg/constants.h" #include "iceberg/expression/json_serde_internal.h" #include "iceberg/expression/literal.h" +#include "iceberg/file_format.h" #include "iceberg/json_serde_internal.h" +#include "iceberg/manifest/manifest_entry.h" #include "iceberg/name_mapping.h" #include "iceberg/partition_field.h" #include "iceberg/partition_spec.h" #include "iceberg/result.h" +#include "iceberg/row/partition_values.h" #include "iceberg/schema.h" +#include "iceberg/schema_field.h" #include "iceberg/snapshot.h" #include "iceberg/sort_order.h" #include "iceberg/statistics_file.h" @@ -240,6 +245,137 @@ constexpr std::string_view kRequirementAssertDefaultSortOrderID = constexpr std::string_view kLastAssignedFieldId = "last-assigned-field-id"; constexpr std::string_view kLastAssignedPartitionId = "last-assigned-partition-id"; +// DataFile JSON (Iceberg REST ContentFile field names) +constexpr std::string_view kContent = "content"; +constexpr std::string_view kContentData = "data"; +constexpr std::string_view kContentPositionDeletes = "position-deletes"; +constexpr std::string_view kContentEqualityDeletes = "equality-deletes"; +constexpr std::string_view kFilePath = "file-path"; +constexpr std::string_view kFileFormat = "file-format"; +constexpr std::string_view kPartition = "partition"; +constexpr std::string_view kRecordCount = "record-count"; +constexpr std::string_view kColumnSizes = "column-sizes"; +constexpr std::string_view kValueCounts = "value-counts"; +constexpr std::string_view kNullValueCounts = "null-value-counts"; +constexpr std::string_view kNanValueCounts = "nan-value-counts"; +constexpr std::string_view kLowerBounds = "lower-bounds"; +constexpr std::string_view kUpperBounds = "upper-bounds"; +constexpr std::string_view kKeyMetadata = "key-metadata"; +constexpr std::string_view kSplitOffsets = "split-offsets"; +constexpr std::string_view kEqualityIds = "equality-ids"; +constexpr std::string_view kReferencedDataFile = "referenced-data-file"; +constexpr std::string_view kContentOffset = "content-offset"; +constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes"; +constexpr std::string_view kMapKeys = "keys"; +constexpr std::string_view kMapValues = "values"; + +template +Result> KeyValueMapFromJson(const nlohmann::json& json, + std::string_view key) { + std::map result; + if (!json.contains(key) || json.at(key).is_null()) { + return result; + } + + ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue(json, key)); + ICEBERG_ASSIGN_OR_RAISE(auto keys, + GetJsonValue>(map_json, kMapKeys)); + ICEBERG_ASSIGN_OR_RAISE(auto values, + GetJsonValue>(map_json, kMapValues)); + if (keys.size() != values.size()) { + return JsonParseError("'{}' map keys and values have different lengths", key); + } + + for (size_t i = 0; i < keys.size(); ++i) { + result[keys[i]] = std::move(values[i]); + } + return result; +} + +template +void SetKeyValueMap(nlohmann::json& json, std::string_view key, + const std::map& map) { + if (map.empty()) { + return; + } + + std::vector keys; + std::vector values; + keys.reserve(map.size()); + values.reserve(map.size()); + for (const auto& [field_id, value] : map) { + keys.push_back(field_id); + values.push_back(value); + } + json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; +} + +Result DataFileToJsonUnchecked(const DataFile& data_file) { + nlohmann::json json; + switch (data_file.content) { + case DataFile::Content::kData: + json[kContent] = kContentData; + break; + case DataFile::Content::kPositionDeletes: + json[kContent] = kContentPositionDeletes; + break; + case DataFile::Content::kEqualityDeletes: + json[kContent] = kContentEqualityDeletes; + break; + } + json[kFilePath] = data_file.file_path; + json[kFileFormat] = ToString(data_file.file_format); + + if (!data_file.partition_spec_id.has_value()) { + return ValidationFailed("Cannot serialize REST content file without 'spec-id'"); + } + json[kSpecId] = data_file.partition_spec_id.value(); + + nlohmann::json partition_json = nlohmann::json::array(); + for (const auto& literal : data_file.partition.values()) { + ICEBERG_ASSIGN_OR_RAISE(auto lit_json, ToJson(literal)); + partition_json.push_back(std::move(lit_json)); + } + json[kPartition] = std::move(partition_json); + + json[kRecordCount] = data_file.record_count; + json[kFileSizeInBytes] = data_file.file_size_in_bytes; + + SetKeyValueMap(json, kColumnSizes, data_file.column_sizes); + SetKeyValueMap(json, kValueCounts, data_file.value_counts); + SetKeyValueMap(json, kNullValueCounts, data_file.null_value_counts); + SetKeyValueMap(json, kNanValueCounts, data_file.nan_value_counts); + SetKeyValueMap(json, kLowerBounds, data_file.lower_bounds); + SetKeyValueMap(json, kUpperBounds, data_file.upper_bounds); + + if (!data_file.key_metadata.empty()) { + json[kKeyMetadata] = data_file.key_metadata; + } + if (!data_file.split_offsets.empty()) { + json[kSplitOffsets] = data_file.split_offsets; + } + if (!data_file.equality_ids.empty()) { + json[kEqualityIds] = data_file.equality_ids; + } + if (data_file.sort_order_id.has_value()) { + json[kSortOrderId] = data_file.sort_order_id.value(); + } + if (data_file.first_row_id.has_value()) { + json[kFirstRowId] = data_file.first_row_id.value(); + } + if (data_file.referenced_data_file.has_value()) { + json[kReferencedDataFile] = data_file.referenced_data_file.value(); + } + if (data_file.content_offset.has_value()) { + json[kContentOffset] = data_file.content_offset.value(); + } + if (data_file.content_size_in_bytes.has_value()) { + json[kContentSizeInBytes] = data_file.content_size_in_bytes.value(); + } + + return json; +} + } // namespace nlohmann::json ToJson(const SortField& sort_field) { @@ -1985,4 +2121,138 @@ Result> TableRequirementFromJson( return JsonParseError("Unknown table requirement type: {}", type); } +Result DataFileFromJson( + const nlohmann::json& json, + const std::unordered_map>& + partition_spec_by_id, + const Schema& schema) { + if (!json.is_object()) { + return JsonParseError("DataFile must be a JSON object: {}", SafeDumpJson(json)); + } + DataFile data_file; + + ICEBERG_ASSIGN_OR_RAISE(auto content_str, GetJsonValue(json, kContent)); + if (content_str == kContentData) { + data_file.content = DataFile::Content::kData; + } else if (content_str == kContentPositionDeletes) { + data_file.content = DataFile::Content::kPositionDeletes; + } else if (content_str == kContentEqualityDeletes) { + data_file.content = DataFile::Content::kEqualityDeletes; + } else { + return JsonParseError("Unknown data file content: {}", content_str); + } + + ICEBERG_ASSIGN_OR_RAISE(data_file.file_path, + GetJsonValue(json, kFilePath)); + ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue(json, kFileFormat)); + ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str)); + + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); + data_file.partition_spec_id = spec_id; + + ICEBERG_ASSIGN_OR_RAISE(auto partition_vals, + GetJsonValue(json, kPartition)); + if (!partition_vals.is_array()) { + return JsonParseError("PartitionValues must be a JSON array: {}", + SafeDumpJson(partition_vals)); + } + std::vector literals; + auto it = partition_spec_by_id.find(spec_id); + if (it == partition_spec_by_id.end()) { + return JsonParseError("Invalid partition spec id: {}", spec_id); + } + ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema)); + auto fields = struct_type->fields(); + if (partition_vals.size() != fields.size()) { + return JsonParseError("Invalid partition data size: expected = {}, actual = {}", + fields.size(), partition_vals.size()); + } + for (size_t pos = 0; pos < fields.size(); ++pos) { + ICEBERG_ASSIGN_OR_RAISE( + auto literal, LiteralFromJson(partition_vals[pos], fields[pos].type().get())); + literals.push_back(std::move(literal)); + } + data_file.partition = PartitionValues(std::move(literals)); + + ICEBERG_ASSIGN_OR_RAISE(data_file.record_count, + GetJsonValue(json, kRecordCount)); + ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes, + GetJsonValue(json, kFileSizeInBytes)); + + ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes, + KeyValueMapFromJson(json, kColumnSizes)); + ICEBERG_ASSIGN_OR_RAISE(data_file.value_counts, + KeyValueMapFromJson(json, kValueCounts)); + ICEBERG_ASSIGN_OR_RAISE(data_file.null_value_counts, + KeyValueMapFromJson(json, kNullValueCounts)); + ICEBERG_ASSIGN_OR_RAISE(data_file.nan_value_counts, + KeyValueMapFromJson(json, kNanValueCounts)); + ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds, + KeyValueMapFromJson>(json, kLowerBounds)); + ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds, + KeyValueMapFromJson>(json, kUpperBounds)); + + if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata, + GetJsonValue>(json, kKeyMetadata)); + } + if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets, + GetJsonValue>(json, kSplitOffsets)); + } + if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids, + GetJsonValue>(json, kEqualityIds)); + } + if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id, + GetJsonValue(json, kSortOrderId)); + } + if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id, + GetJsonValue(json, kFirstRowId)); + } + if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file, + GetJsonValue(json, kReferencedDataFile)); + } + if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset, + GetJsonValue(json, kContentOffset)); + } + if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes, + GetJsonValue(json, kContentSizeInBytes)); + } + + return data_file; +} + +Result ToJson( + const DataFile& data_file, + const std::unordered_map>& + partition_specs_by_id, + const Schema& schema) { + if (!data_file.partition_spec_id.has_value()) { + return ValidationFailed("Invalid partition spec id from content file: null"); + } + auto it = partition_specs_by_id.find(data_file.partition_spec_id.value()); + if (it == partition_specs_by_id.end() || !it->second) { + return ValidationFailed("Invalid partition spec: null"); + } + if (data_file.partition_spec_id.value() != it->second->spec_id()) { + return ValidationFailed( + "Invalid partition spec id from content file: expected = {}, actual = {}", + it->second->spec_id(), data_file.partition_spec_id.value()); + } + ICEBERG_ASSIGN_OR_RAISE(auto partition_type, it->second->PartitionType(schema)); + if (data_file.partition.num_fields() != partition_type->fields().size()) { + return ValidationFailed( + "Invalid partition data from content file: expected = {}, actual = {}", + partition_type->fields().empty() ? "unpartitioned" : "partitioned", + data_file.partition.num_fields() == 0 ? "unpartitioned" : "partitioned"); + } + return DataFileToJsonUnchecked(data_file); +} + } // namespace iceberg diff --git a/src/iceberg/json_serde_internal.h b/src/iceberg/json_serde_internal.h index 1a30e8e8f..3aaced5c9 100644 --- a/src/iceberg/json_serde_internal.h +++ b/src/iceberg/json_serde_internal.h @@ -29,6 +29,7 @@ #include #include "iceberg/encryption/encrypted_key.h" +#include "iceberg/manifest/manifest_entry.h" #include "iceberg/result.h" #include "iceberg/statistics_file.h" #include "iceberg/table_metadata.h" @@ -426,4 +427,22 @@ ICEBERG_EXPORT nlohmann::json ToJson(const TableRequirement& requirement); ICEBERG_EXPORT Result> TableRequirementFromJson( const nlohmann::json& json); +/// \brief Serializes a `DataFile` (content file) to JSON. +/// +/// The JSON object uses the Iceberg REST ContentFile field names (`spec-id`, +/// `file-path`, `file-format`, maps as `{keys, values}`, and so on). `spec-id` is +/// required. Partition values are encoded with the given partition spec and schema. +ICEBERG_EXPORT Result ToJson( + const DataFile& data_file, + const std::unordered_map>& + partition_specs_by_id, + const Schema& schema); + +/// \brief Deserializes a JSON object into a `DataFile`. +ICEBERG_EXPORT Result DataFileFromJson( + const nlohmann::json& json, + const std::unordered_map>& + partition_spec_by_id, + const Schema& schema); + } // namespace iceberg diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 22b0a08de..c7779d703 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -18,15 +18,19 @@ */ #include +#include #include #include #include #include "iceberg/expression/literal.h" +#include "iceberg/file_format.h" #include "iceberg/json_serde_internal.h" +#include "iceberg/manifest/manifest_entry.h" #include "iceberg/name_mapping.h" #include "iceberg/partition_spec.h" +#include "iceberg/row/partition_values.h" #include "iceberg/schema.h" #include "iceberg/schema_field.h" #include "iceberg/snapshot.h" @@ -1044,4 +1048,37 @@ TEST(TableRequirementJsonTest, TableRequirementUnknownType) { EXPECT_THAT(result, HasErrorMessage("Unknown table requirement type")); } +std::unordered_map> UnpartitionedSpecs() { + return {{PartitionSpec::kInitialSpecId, PartitionSpec::Unpartitioned()}}; +} + +DataFile MakeUnpartitionedDataFile(std::string path, int64_t record_count = 100, + int64_t file_size = 12345) { + DataFile data_file; + data_file.content = DataFile::Content::kData; + data_file.file_path = std::move(path); + data_file.file_format = FileFormatType::kParquet; + data_file.partition_spec_id = PartitionSpec::kInitialSpecId; + data_file.partition = PartitionValues{}; + data_file.record_count = record_count; + data_file.file_size_in_bytes = file_size; + return data_file; +} + +TEST(DataFileJsonTest, RoundTripRequiredFields) { + auto data_file = MakeUnpartitionedDataFile("s3://bucket/data/file.parquet"); + Schema schema({}, 0); + ICEBERG_UNWRAP_OR_FAIL(auto json, ToJson(data_file, UnpartitionedSpecs(), schema)); + EXPECT_EQ(json["content"], "data"); + EXPECT_EQ(json["file-path"], "s3://bucket/data/file.parquet"); + EXPECT_EQ(json["spec-id"], 0); + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, + DataFileFromJson(json, UnpartitionedSpecs(), schema)); + EXPECT_EQ(parsed.file_path, data_file.file_path); + EXPECT_EQ(parsed.record_count, data_file.record_count); + EXPECT_EQ(parsed.file_size_in_bytes, data_file.file_size_in_bytes); + EXPECT_EQ(parsed.partition_spec_id, data_file.partition_spec_id); +} + } // namespace iceberg