From c263a0b9e6bdc1d4d73d62249fab4290afae7dfd Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 07:57:43 +0000 Subject: [PATCH 1/6] refactor(serde): move DataFile JSON serde into core --- src/iceberg/catalog/rest/json_serde.cc | 255 +---------------------- src/iceberg/json_serde.cc | 269 +++++++++++++++++++++++++ src/iceberg/json_serde_internal.h | 19 ++ src/iceberg/test/json_serde_test.cc | 36 ++++ 4 files changed, 326 insertions(+), 253 deletions(-) diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 3ce753f18..83a2dec09 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( @@ -357,97 +191,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 { diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index f322996ec..c2e685ef4 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,137 @@ 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..b15d08790 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,36 @@ 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 From ece593036953b7ed15f02ae2fc28a97b068e4ec4 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 08:54:58 +0000 Subject: [PATCH 2/6] fix(serde): disambiguate REST DataFile calls --- src/iceberg/catalog/rest/json_serde.cc | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 83a2dec09..46593f0b5 100644 --- a/src/iceberg/catalog/rest/json_serde.cc +++ b/src/iceberg/catalog/rest/json_serde.cc @@ -156,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. @@ -244,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()) { @@ -270,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()) { @@ -306,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))); } From 5304e282d76b70e84366f5efd32defa71bdb7f4e Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 09:16:33 +0000 Subject: [PATCH 3/6] fix clang error --- src/iceberg/json_serde.cc | 3 ++- src/iceberg/test/json_serde_test.cc | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index c2e685ef4..1f5caf8f0 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -2142,7 +2142,8 @@ Result DataFileFromJson( return JsonParseError("Unknown data file content: {}", content_str); } - ICEBERG_ASSIGN_OR_RAISE(data_file.file_path, GetJsonValue(json, kFilePath)); + 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)); diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index b15d08790..c7779d703 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -1073,7 +1073,8 @@ TEST(DataFileJsonTest, RoundTripRequiredFields) { 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)); + 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); From 2c48f026a40f384187e698f80247f1c7ed355a77 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 08:02:10 +0000 Subject: [PATCH 4/6] fix(serde): align DataFile JSON with Java --- src/iceberg/json_serde.cc | 273 +++++++++++++++++++---- src/iceberg/test/json_serde_test.cc | 108 +++++++++ src/iceberg/test/rest_json_serde_test.cc | 10 +- 3 files changed, 339 insertions(+), 52 deletions(-) diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index 1f5caf8f0..45dfa29f5 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -18,10 +18,13 @@ */ #include +#include #include #include +#include #include #include +#include #include #include @@ -269,7 +272,58 @@ constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes"; constexpr std::string_view kMapKeys = "keys"; constexpr std::string_view kMapValues = "values"; -template +template +Result IntegerFromJson(const nlohmann::json& json, + std::string_view description) { + if (!json.is_number_integer()) { + return JsonParseError("{} must be an integer, but is {}", description, + json.type_name()); + } + if (json.is_number_unsigned()) { + const auto value = json.get(); + if (value > static_cast(std::numeric_limits::max())) { + return JsonParseError("{} integer is out of range: {}", description, + SafeDumpJson(json)); + } + return static_cast(value); + } + + const auto value = json.get(); + if (value < static_cast(std::numeric_limits::min()) || + value > static_cast(std::numeric_limits::max())) { + return JsonParseError("{} integer is out of range: {}", description, + SafeDumpJson(json)); + } + return static_cast(value); +} + +template +Result GetJsonInteger(const nlohmann::json& json, std::string_view key) { + ICEBERG_ASSIGN_OR_RAISE(auto value, GetJsonValue(json, key)); + return IntegerFromJson(value, std::format("'{}'", key)); +} + +template +Result> IntegerVectorFromJson(const nlohmann::json& json, + std::string_view key) { + ICEBERG_ASSIGN_OR_RAISE(auto values_json, GetJsonValue(json, key)); + if (!values_json.is_array()) { + return JsonParseError("'{}' must be an array, but is {}", key, + values_json.type_name()); + } + + std::vector values; + values.reserve(values_json.size()); + for (const auto& value_json : values_json) { + ICEBERG_ASSIGN_OR_RAISE( + auto value, + IntegerFromJson(value_json, std::format("'{}' element", key))); + values.push_back(value); + } + return values; +} + +template Result> KeyValueMapFromJson(const nlohmann::json& json, std::string_view key) { std::map result; @@ -278,16 +332,27 @@ Result> KeyValueMapFromJson(const nlohmann::json& json, } 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()) { + ICEBERG_ASSIGN_OR_RAISE(auto keys_json, + GetJsonValue(map_json, kMapKeys)); + ICEBERG_ASSIGN_OR_RAISE(auto values_json, + GetJsonValue(map_json, kMapValues)); + if (!keys_json.is_array() || !values_json.is_array()) { + return JsonParseError("'{}' map keys and values must be arrays", key); + } + if (keys_json.size() != values_json.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]); + for (size_t i = 0; i < keys_json.size(); ++i) { + ICEBERG_ASSIGN_OR_RAISE( + auto field_id, + IntegerFromJson(keys_json[i], std::format("'{}' map key", key))); + ICEBERG_ASSIGN_OR_RAISE( + auto value, + IntegerFromJson(values_json[i], std::format("'{}' map value", key))); + if (!result.emplace(field_id, value).second) { + return JsonParseError("'{}' map contains duplicate key {}", key, field_id); + } } return result; } @@ -310,6 +375,118 @@ void SetKeyValueMap(nlohmann::json& json, std::string_view key, json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; } +std::string BytesToHex(const std::vector& bytes) { + std::string hex; + hex.reserve(bytes.size() * 2); + for (uint8_t byte : bytes) { + hex += std::format("{:02X}", byte); + } + return hex; +} + +Result>> BytesMapFromJson( + 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_json, + GetJsonValue(map_json, kMapKeys)); + ICEBERG_ASSIGN_OR_RAISE(auto values_json, + GetJsonValue(map_json, kMapValues)); + if (!keys_json.is_array() || !values_json.is_array()) { + return JsonParseError("'{}' map keys and values must be arrays", key); + } + if (keys_json.size() != values_json.size()) { + return JsonParseError("'{}' map keys and values have different lengths", key); + } + + for (size_t i = 0; i < keys_json.size(); ++i) { + ICEBERG_ASSIGN_OR_RAISE( + auto field_id, + IntegerFromJson(keys_json[i], std::format("'{}' map key", key))); + if (!values_json[i].is_string()) { + return JsonParseError("'{}' map value must be a string, but is {}", key, + values_json[i].type_name()); + } + ICEBERG_ASSIGN_OR_RAISE( + auto bytes, StringUtils::HexStringToBytes(values_json[i].get())); + if (!result.emplace(field_id, std::move(bytes)).second) { + return JsonParseError("'{}' map contains duplicate key {}", key, field_id); + } + } + return result; +} + +void SetBytesMap(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(BytesToHex(value)); + } + json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; +} + +Result PartitionLiteralFromJson(const nlohmann::json* value, + const SchemaField& field) { + if (!field.type() || !field.type()->is_primitive()) { + return InvalidSchema("Partition field {} must have a primitive type", + field.field_id()); + } + + auto primitive_type = internal::checked_pointer_cast(field.type()); + if (value == nullptr || value->is_null()) { + return Literal::Null(std::move(primitive_type)); + } + return LiteralFromJson(*value, primitive_type.get()); +} + +Result PartitionValuesFromJson(const nlohmann::json& json, + const StructType& partition_type) { + const auto& fields = partition_type.fields(); + std::vector values; + values.reserve(fields.size()); + + if (json.is_array()) { + if (json.size() != fields.size()) { + return JsonParseError("Invalid partition data size: expected = {}, actual = {}", + fields.size(), json.size()); + } + for (size_t pos = 0; pos < fields.size(); ++pos) { + ICEBERG_ASSIGN_OR_RAISE(auto value, + PartitionLiteralFromJson(&json[pos], fields[pos])); + values.push_back(std::move(value)); + } + } else if (json.is_object()) { + if (json.size() > fields.size()) { + return JsonParseError("Invalid partition data size: expected <= {}, actual = {}", + fields.size(), json.size()); + } + for (const auto& field : fields) { + const auto it = json.find(std::to_string(field.field_id())); + const nlohmann::json* value = it == json.end() ? nullptr : &it.value(); + ICEBERG_ASSIGN_OR_RAISE(auto literal, PartitionLiteralFromJson(value, field)); + values.push_back(std::move(literal)); + } + } else { + return JsonParseError( + "Invalid partition data for content file: expected array or object ({})", + SafeDumpJson(json)); + } + + return PartitionValues(std::move(values)); +} + Result DataFileToJsonUnchecked(const DataFile& data_file) { nlohmann::json json; switch (data_file.content) { @@ -327,7 +504,7 @@ Result DataFileToJsonUnchecked(const DataFile& data_file) { 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'"); + return ValidationFailed("Cannot serialize content file without 'spec-id'"); } json[kSpecId] = data_file.partition_spec_id.value(); @@ -345,11 +522,11 @@ Result DataFileToJsonUnchecked(const DataFile& data_file) { 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); + SetBytesMap(json, kLowerBounds, data_file.lower_bounds); + SetBytesMap(json, kUpperBounds, data_file.upper_bounds); if (!data_file.key_metadata.empty()) { - json[kKeyMetadata] = data_file.key_metadata; + json[kKeyMetadata] = BytesToHex(data_file.key_metadata); } if (!data_file.split_offsets.empty()) { json[kSplitOffsets] = data_file.split_offsets; @@ -2132,11 +2309,13 @@ Result DataFileFromJson( DataFile data_file; ICEBERG_ASSIGN_OR_RAISE(auto content_str, GetJsonValue(json, kContent)); - if (content_str == kContentData) { + if (content_str == kContentData || content_str == "DATA") { data_file.content = DataFile::Content::kData; - } else if (content_str == kContentPositionDeletes) { + } else if (content_str == kContentPositionDeletes || + content_str == "POSITION_DELETES") { data_file.content = DataFile::Content::kPositionDeletes; - } else if (content_str == kContentEqualityDeletes) { + } else if (content_str == kContentEqualityDeletes || + content_str == "EQUALITY_DELETES") { data_file.content = DataFile::Content::kEqualityDeletes; } else { return JsonParseError("Unknown data file content: {}", content_str); @@ -2147,37 +2326,25 @@ Result DataFileFromJson( 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)); + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonInteger(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()) { + if (it == partition_spec_by_id.end() || !it->second) { 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)); + if (json.contains(kPartition)) { + ICEBERG_ASSIGN_OR_RAISE(auto partition_vals, + GetJsonValue(json, kPartition)); + ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema)); + ICEBERG_ASSIGN_OR_RAISE(data_file.partition, + PartitionValuesFromJson(partition_vals, *struct_type)); } - data_file.partition = PartitionValues(std::move(literals)); ICEBERG_ASSIGN_OR_RAISE(data_file.record_count, - GetJsonValue(json, kRecordCount)); + GetJsonInteger(json, kRecordCount)); ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes, - GetJsonValue(json, kFileSizeInBytes)); + GetJsonInteger(json, kFileSizeInBytes)); ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes, KeyValueMapFromJson(json, kColumnSizes)); @@ -2187,30 +2354,30 @@ Result DataFileFromJson( 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)); + ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds, BytesMapFromJson(json, kLowerBounds)); + ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds, BytesMapFromJson(json, kUpperBounds)); if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(auto key_metadata, + GetJsonValue(json, kKeyMetadata)); ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata, - GetJsonValue>(json, kKeyMetadata)); + StringUtils::HexStringToBytes(key_metadata)); } if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets, - GetJsonValue>(json, kSplitOffsets)); + IntegerVectorFromJson(json, kSplitOffsets)); } if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids, - GetJsonValue>(json, kEqualityIds)); + IntegerVectorFromJson(json, kEqualityIds)); } if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id, - GetJsonValue(json, kSortOrderId)); + GetJsonInteger(json, kSortOrderId)); } if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id, - GetJsonValue(json, kFirstRowId)); + GetJsonInteger(json, kFirstRowId)); } if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file, @@ -2218,11 +2385,11 @@ Result DataFileFromJson( } if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset, - GetJsonValue(json, kContentOffset)); + GetJsonInteger(json, kContentOffset)); } if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes, - GetJsonValue(json, kContentSizeInBytes)); + GetJsonInteger(json, kContentSizeInBytes)); } return data_file; @@ -2252,6 +2419,16 @@ Result ToJson( partition_type->fields().empty() ? "unpartitioned" : "partitioned", data_file.partition.num_fields() == 0 ? "unpartitioned" : "partitioned"); } + for (size_t pos = 0; pos < partition_type->fields().size(); ++pos) { + const auto& literal = data_file.partition.values()[pos]; + const auto& expected_type = partition_type->fields()[pos].type(); + if (!literal.IsNull() && (!literal.type() || *literal.type() != *expected_type)) { + return ValidationFailed( + "Invalid partition value type at position {}: expected = {}, actual = {}", pos, + expected_type->ToString(), + literal.type() ? literal.type()->ToString() : "unknown"); + } + } return DataFileToJsonUnchecked(data_file); } diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index c7779d703..0353f1ced 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -1081,4 +1081,112 @@ TEST(DataFileJsonTest, RoundTripRequiredFields) { EXPECT_EQ(parsed.partition_spec_id, data_file.partition_spec_id); } +class DataFileJavaJsonTest : public ::testing::Test { + protected: + void SetUp() override { + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(schema_, /*spec_id=*/0, + {PartitionField(1, 1000, "id", Transform::Identity())}, + /*allow_missing_fields=*/false)); + spec_ = std::shared_ptr(std::move(spec)); + specs_.emplace(spec_->spec_id(), spec_); + } + + DataFile DataFileForTest() const { + DataFile data_file; + data_file.content = DataFile::Content::kData; + data_file.file_path = "/path/to/data.parquet"; + data_file.file_format = FileFormatType::kParquet; + data_file.partition_spec_id = 0; + data_file.partition = PartitionValues({Literal::Int(7)}); + data_file.record_count = 1; + data_file.file_size_in_bytes = 10; + data_file.lower_bounds = {{1, {0x01, 0x00, 0x00, 0x00}}}; + data_file.upper_bounds = {{1, {0x05, 0x00, 0x00, 0x00}}}; + data_file.key_metadata = {0x0A, 0x0B}; + data_file.sort_order_id = 0; + return data_file; + } + + nlohmann::json JavaGolden() const { + return R"({ + "spec-id": 0, + "content": "data", + "file-path": "/path/to/data.parquet", + "file-format": "parquet", + "partition": [7], + "file-size-in-bytes": 10, + "record-count": 1, + "lower-bounds": {"keys": [1], "values": ["01000000"]}, + "upper-bounds": {"keys": [1], "values": ["05000000"]}, + "key-metadata": "0A0B", + "sort-order-id": 0 + })"_json; + } + + Schema schema_{{SchemaField::MakeRequired(1, "id", int32())}, 0}; + std::shared_ptr spec_; + std::unordered_map> specs_; +}; + +TEST_F(DataFileJavaJsonTest, SerializesAndParsesJavaEncoding) { + ICEBERG_UNWRAP_OR_FAIL(auto json, ToJson(DataFileForTest(), specs_, schema_)); + EXPECT_EQ(json, JavaGolden()); + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed, DataFileForTest()); +} + +TEST_F(DataFileJavaJsonTest, ParsesPartitionObjectAndLegacyContentName) { + auto json = JavaGolden(); + json["content"] = "DATA"; + json["partition"] = {{"1000", 7}}; + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed.content, DataFile::Content::kData); + ASSERT_EQ(parsed.partition.num_fields(), 1U); + EXPECT_EQ(parsed.partition.values()[0], Literal::Int(7)); +} + +TEST_F(DataFileJavaJsonTest, RejectsInvalidIntegersAndDuplicateMetricKeys) { + for (std::string_view field : + {"spec-id", "file-size-in-bytes", "record-count", "sort-order-id"}) { + SCOPED_TRACE(field); + auto json = JavaGolden(); + json[field] = 0.5; + EXPECT_THAT(DataFileFromJson(json, specs_, schema_), + IsError(ErrorKind::kJsonParseError)); + } + + auto duplicate_keys = JavaGolden(); + duplicate_keys["column-sizes"] = {{"keys", {1, 1}}, {"values", {10, 20}}}; + auto result = DataFileFromJson(duplicate_keys, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("duplicate key")); +} + +TEST_F(DataFileJavaJsonTest, AcceptsMissingPartitionAndNullMetricMaps) { + auto json = JavaGolden(); + json.erase("partition"); + for (std::string_view field : {"column-sizes", "value-counts", "null-value-counts", + "nan-value-counts", "lower-bounds", "upper-bounds"}) { + json[field] = nullptr; + } + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed.partition.num_fields(), 0U); + EXPECT_TRUE(parsed.column_sizes.empty()); + EXPECT_TRUE(parsed.lower_bounds.empty()); +} + +TEST_F(DataFileJavaJsonTest, RejectsInvalidPartitionTypeWhenSerializing) { + auto data_file = DataFileForTest(); + data_file.partition = PartitionValues({Literal::String("not-an-int")}); + + auto result = ToJson(data_file, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(result, HasErrorMessage("partition value type")); +} + } // namespace iceberg diff --git a/src/iceberg/test/rest_json_serde_test.cc b/src/iceberg/test/rest_json_serde_test.cc index ec41e4a66..e1979dbef 100644 --- a/src/iceberg/test/rest_json_serde_test.cc +++ b/src/iceberg/test/rest_json_serde_test.cc @@ -19,7 +19,9 @@ #include #include +#include #include +#include #include #include @@ -2119,7 +2121,7 @@ TEST(DataFileFromJsonTest, MissingSpecId) { EXPECT_THAT(result, HasErrorMessage("Missing 'spec-id'")); } -TEST(DataFileFromJsonTest, MissingPartition) { +TEST(DataFileFromJsonTest, MissingPartitionIsAccepted) { auto json = R"({ "content": "data", "file-path": "s3://bucket/data/file.parquet", @@ -2129,9 +2131,9 @@ TEST(DataFileFromJsonTest, MissingPartition) { "record-count": 10 })"_json; - auto result = DataFileFromJson(json, UnpartitionedSpecs(), Schema({}, 0)); - EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); - EXPECT_THAT(result, HasErrorMessage("Missing 'partition'")); + ICEBERG_UNWRAP_OR_FAIL(auto data_file, + DataFileFromJson(json, UnpartitionedSpecs(), Schema({}, 0))); + EXPECT_EQ(data_file.partition.num_fields(), 0U); } TEST(DataFileFromJsonTest, NotAnObject) { From fef0a43d59986b1af2a32c783ad4dd7cfa17079d Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 10:41:15 +0000 Subject: [PATCH 5/6] test(serde): cover DataFile Java JSON edge cases --- src/iceberg/test/json_serde_test.cc | 140 ++++++++++++++++++++++++++-- 1 file changed, 134 insertions(+), 6 deletions(-) diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 0353f1ced..3f8dbc165 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -18,6 +18,8 @@ */ #include +#include +#include #include #include @@ -1138,32 +1140,158 @@ TEST_F(DataFileJavaJsonTest, SerializesAndParsesJavaEncoding) { EXPECT_EQ(parsed, DataFileForTest()); } -TEST_F(DataFileJavaJsonTest, ParsesPartitionObjectAndLegacyContentName) { +TEST_F(DataFileJavaJsonTest, ParsesPartitionObject) { auto json = JavaGolden(); - json["content"] = "DATA"; json["partition"] = {{"1000", 7}}; ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); - EXPECT_EQ(parsed.content, DataFile::Content::kData); ASSERT_EQ(parsed.partition.num_fields(), 1U); EXPECT_EQ(parsed.partition.values()[0], Literal::Int(7)); } +TEST_F(DataFileJavaJsonTest, ParsesAllLegacyContentNames) { + struct Case { + std::string_view legacy; + DataFile::Content content; + std::string_view canonical; + }; + for (const auto& test_case : + {Case{"DATA", DataFile::Content::kData, "data"}, + Case{"POSITION_DELETES", DataFile::Content::kPositionDeletes, "position-deletes"}, + Case{"EQUALITY_DELETES", DataFile::Content::kEqualityDeletes, + "equality-deletes"}}) { + SCOPED_TRACE(test_case.legacy); + auto json = JavaGolden(); + json["content"] = std::string(test_case.legacy); + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed.content, test_case.content); + ICEBERG_UNWRAP_OR_FAIL(auto serialized, ToJson(parsed, specs_, schema_)); + EXPECT_EQ(serialized["content"], std::string(test_case.canonical)); + } +} + +TEST_F(DataFileJavaJsonTest, ParsesFieldIdPartitionsInSpecOrderAndTypedNulls) { + Schema schema{{SchemaField::MakeRequired(1, "id", int32()), + SchemaField::MakeOptional(2, "region", string())}, + 0}; + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(schema, /*spec_id=*/0, + {PartitionField(2, 1001, "region", Transform::Identity()), + PartitionField(1, 1000, "id", Transform::Identity())}, + /*allow_missing_fields=*/false)); + std::unordered_map> specs{ + {0, std::shared_ptr(std::move(spec))}}; + + auto json = JavaGolden(); + json["partition"] = {{"1000", 7}, {"1001", "north"}}; + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs, schema)); + ASSERT_EQ(parsed.partition.num_fields(), 2U); + EXPECT_EQ(parsed.partition.values()[0], Literal::String("north")); + EXPECT_EQ(parsed.partition.values()[1], Literal::Int(7)); + + json["partition"] = {{"1000", 7}}; + ICEBERG_UNWRAP_OR_FAIL(auto sparse, DataFileFromJson(json, specs, schema)); + ASSERT_EQ(sparse.partition.num_fields(), 2U); + EXPECT_TRUE(sparse.partition.values()[0].IsNull()); + ASSERT_NE(sparse.partition.values()[0].type(), nullptr); + EXPECT_EQ(sparse.partition.values()[0].type()->type_id(), TypeId::kString); + EXPECT_EQ(sparse.partition.values()[1], Literal::Int(7)); + + json["partition"] = nlohmann::json::array({nullptr, 7}); + ICEBERG_UNWRAP_OR_FAIL(auto positional, DataFileFromJson(json, specs, schema)); + ASSERT_EQ(positional.partition.num_fields(), 2U); + EXPECT_TRUE(positional.partition.values()[0].IsNull()); + ASSERT_NE(positional.partition.values()[0].type(), nullptr); + EXPECT_EQ(positional.partition.values()[0].type()->type_id(), TypeId::kString); + EXPECT_EQ(positional.partition.values()[1], Literal::Int(7)); +} + +TEST_F(DataFileJavaJsonTest, RejectsMalformedPartitionShapes) { + for (const auto& partition : + {nlohmann::json::array(), nlohmann::json::object({{"1000", 7}, {"9999", 8}}), + nlohmann::json("invalid")}) { + SCOPED_TRACE(partition.dump()); + auto json = JavaGolden(); + json["partition"] = partition; + EXPECT_THAT(DataFileFromJson(json, specs_, schema_), + IsError(ErrorKind::kJsonParseError)); + } +} + TEST_F(DataFileJavaJsonTest, RejectsInvalidIntegersAndDuplicateMetricKeys) { for (std::string_view field : - {"spec-id", "file-size-in-bytes", "record-count", "sort-order-id"}) { + {"spec-id", "file-size-in-bytes", "record-count", "sort-order-id", "first-row-id", + "content-offset", "content-size-in-bytes"}) { SCOPED_TRACE(field); auto json = JavaGolden(); json[field] = 0.5; - EXPECT_THAT(DataFileFromJson(json, specs_, schema_), - IsError(ErrorKind::kJsonParseError)); + auto result = DataFileFromJson(json, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); } + auto expect_invalid = [&](nlohmann::json json, std::string_view field) { + auto result = DataFileFromJson(json, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); + }; + + auto fractional_key = JavaGolden(); + fractional_key["column-sizes"] = {{"keys", {0.5}}, {"values", {10}}}; + expect_invalid(fractional_key, "column-sizes"); + + auto fractional_value = JavaGolden(); + fractional_value["column-sizes"] = {{"keys", {1}}, {"values", {0.5}}}; + expect_invalid(fractional_value, "column-sizes"); + + auto fractional_bound_key = JavaGolden(); + fractional_bound_key["lower-bounds"]["keys"] = {0.5}; + expect_invalid(fractional_bound_key, "lower-bounds"); + + for (std::string_view field : {"split-offsets", "equality-ids"}) { + auto json = JavaGolden(); + json[field] = {0.5}; + expect_invalid(std::move(json), field); + } + + auto oversized_spec_id = JavaGolden(); + oversized_spec_id["spec-id"] = 2147483648ULL; + expect_invalid(oversized_spec_id, "spec-id"); + + auto oversized_record_count = JavaGolden(); + oversized_record_count["record-count"] = 18446744073709551615ULL; + expect_invalid(oversized_record_count, "record-count"); + auto duplicate_keys = JavaGolden(); duplicate_keys["column-sizes"] = {{"keys", {1, 1}}, {"values", {10, 20}}}; auto result = DataFileFromJson(duplicate_keys, specs_, schema_); EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); EXPECT_THAT(result, HasErrorMessage("duplicate key")); + + auto duplicate_bounds = JavaGolden(); + duplicate_bounds["lower-bounds"] = {{"keys", {1, 1}}, {"values", {"01", "02"}}}; + auto bounds_result = DataFileFromJson(duplicate_bounds, specs_, schema_); + EXPECT_THAT(bounds_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(bounds_result, HasErrorMessage("duplicate key")); +} + +TEST_F(DataFileJavaJsonTest, RejectsMalformedHexMetadata) { + auto bad_bound = JavaGolden(); + bad_bound["lower-bounds"]["values"] = {"0G"}; + EXPECT_THAT(DataFileFromJson(bad_bound, specs_, schema_), + IsError(ErrorKind::kInvalidArgument)); + + auto non_string_bound = JavaGolden(); + non_string_bound["lower-bounds"]["values"] = {7}; + EXPECT_THAT(DataFileFromJson(non_string_bound, specs_, schema_), + IsError(ErrorKind::kJsonParseError)); + + auto bad_key_metadata = JavaGolden(); + bad_key_metadata["key-metadata"] = "0G"; + EXPECT_THAT(DataFileFromJson(bad_key_metadata, specs_, schema_), + IsError(ErrorKind::kInvalidArgument)); } TEST_F(DataFileJavaJsonTest, AcceptsMissingPartitionAndNullMetricMaps) { From d4a1c887eaec1debd33e9bace944ef8fc2fd4e94 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 20:31:27 +0000 Subject: [PATCH 6/6] fix(serde): use designated initializers in DataFile tests --- src/iceberg/test/json_serde_test.cc | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 3f8dbc165..7bb867290 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -1156,10 +1156,13 @@ TEST_F(DataFileJavaJsonTest, ParsesAllLegacyContentNames) { std::string_view canonical; }; for (const auto& test_case : - {Case{"DATA", DataFile::Content::kData, "data"}, - Case{"POSITION_DELETES", DataFile::Content::kPositionDeletes, "position-deletes"}, - Case{"EQUALITY_DELETES", DataFile::Content::kEqualityDeletes, - "equality-deletes"}}) { + {Case{.legacy = "DATA", .content = DataFile::Content::kData, .canonical = "data"}, + Case{.legacy = "POSITION_DELETES", + .content = DataFile::Content::kPositionDeletes, + .canonical = "position-deletes"}, + Case{.legacy = "EQUALITY_DELETES", + .content = DataFile::Content::kEqualityDeletes, + .canonical = "equality-deletes"}}) { SCOPED_TRACE(test_case.legacy); auto json = JavaGolden(); json["content"] = std::string(test_case.legacy);