Skip to content

Commit c263a0b

Browse files
refactor(serde): move DataFile JSON serde into core
1 parent 3c5715c commit c263a0b

4 files changed

Lines changed: 326 additions & 253 deletions

File tree

‎src/iceberg/catalog/rest/json_serde.cc‎

Lines changed: 2 additions & 253 deletions
Original file line numberDiff line numberDiff line change
@@ -105,35 +105,9 @@ constexpr std::string_view kEndSnapshotId = "end-snapshot-id";
105105
constexpr std::string_view kStatsFields = "stats-fields";
106106
constexpr std::string_view kMinRowsRequested = "min-rows-requested";
107107
constexpr std::string_view kPlanTask = "plan-task";
108-
constexpr std::string_view kContent = "content";
109-
constexpr std::string_view kContentData = "data";
110-
constexpr std::string_view kContentPositionDeletes = "position-deletes";
111-
constexpr std::string_view kContentEqualityDeletes = "equality-deletes";
112-
constexpr std::string_view kFilePath = "file-path";
113-
constexpr std::string_view kFileFormat = "file-format";
114-
constexpr std::string_view kSpecId = "spec-id";
115-
constexpr std::string_view kPartition = "partition";
116-
constexpr std::string_view kRecordCount = "record-count";
117-
constexpr std::string_view kFileSizeInBytes = "file-size-in-bytes";
118-
constexpr std::string_view kColumnSizes = "column-sizes";
119-
constexpr std::string_view kValueCounts = "value-counts";
120-
constexpr std::string_view kNullValueCounts = "null-value-counts";
121-
constexpr std::string_view kNanValueCounts = "nan-value-counts";
122-
constexpr std::string_view kLowerBounds = "lower-bounds";
123-
constexpr std::string_view kUpperBounds = "upper-bounds";
124-
constexpr std::string_view kKeyMetadata = "key-metadata";
125-
constexpr std::string_view kSplitOffsets = "split-offsets";
126-
constexpr std::string_view kEqualityIds = "equality-ids";
127-
constexpr std::string_view kSortOrderId = "sort-order-id";
128-
constexpr std::string_view kFirstRowId = "first-row-id";
129-
constexpr std::string_view kReferencedDataFile = "referenced-data-file";
130-
constexpr std::string_view kContentOffset = "content-offset";
131-
constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes";
132108
constexpr std::string_view kDataFile = "data-file";
133109
constexpr std::string_view kDeleteFileReferences = "delete-file-references";
134110
constexpr std::string_view kResidualFilter = "residual-filter";
135-
constexpr std::string_view kMapKeys = "keys";
136-
constexpr std::string_view kMapValues = "values";
137111

138112
Result<nlohmann::json> StorageCredentialToJson(const StorageCredential& credential) {
139113
ICEBERG_RETURN_UNEXPECTED(credential.Validate());
@@ -152,154 +126,14 @@ Result<StorageCredential> StorageCredentialFromJson(const nlohmann::json& json)
152126
return credential;
153127
}
154128

155-
template <typename Value>
156-
Result<std::map<int32_t, Value>> KeyValueMapFromJson(const nlohmann::json& json,
157-
std::string_view key) {
158-
std::map<int32_t, Value> result;
159-
if (!json.contains(key) || json.at(key).is_null()) {
160-
return result;
161-
}
162-
163-
ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue<nlohmann::json>(json, key));
164-
ICEBERG_ASSIGN_OR_RAISE(auto keys,
165-
GetJsonValue<std::vector<int32_t>>(map_json, kMapKeys));
166-
ICEBERG_ASSIGN_OR_RAISE(auto values,
167-
GetJsonValue<std::vector<Value>>(map_json, kMapValues));
168-
if (keys.size() != values.size()) {
169-
return JsonParseError("'{}' map keys and values have different lengths", key);
170-
}
171-
172-
for (size_t i = 0; i < keys.size(); ++i) {
173-
result[keys[i]] = std::move(values[i]);
174-
}
175-
return result;
176-
}
177-
178-
template <typename Value>
179-
void SetKeyValueMap(nlohmann::json& json, std::string_view key,
180-
const std::map<int32_t, Value>& map) {
181-
if (map.empty()) {
182-
return;
183-
}
184-
185-
std::vector<int32_t> keys;
186-
std::vector<Value> values;
187-
keys.reserve(map.size());
188-
values.reserve(map.size());
189-
for (const auto& [field_id, value] : map) {
190-
keys.push_back(field_id);
191-
values.push_back(value);
192-
}
193-
json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}};
194-
}
195-
196129
} // namespace
197130

198131
Result<DataFile> DataFileFromJson(
199132
const nlohmann::json& json,
200133
const std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>>&
201134
partition_spec_by_id,
202135
const Schema& schema) {
203-
if (!json.is_object()) {
204-
return JsonParseError("DataFile must be a JSON object: {}", SafeDumpJson(json));
205-
}
206-
DataFile data_file;
207-
208-
ICEBERG_ASSIGN_OR_RAISE(auto content_str, GetJsonValue<std::string>(json, kContent));
209-
if (content_str == kContentData) {
210-
data_file.content = DataFile::Content::kData;
211-
} else if (content_str == kContentPositionDeletes) {
212-
data_file.content = DataFile::Content::kPositionDeletes;
213-
} else if (content_str == kContentEqualityDeletes) {
214-
data_file.content = DataFile::Content::kEqualityDeletes;
215-
} else {
216-
return JsonParseError("Unknown data file content: {}", content_str);
217-
}
218-
219-
ICEBERG_ASSIGN_OR_RAISE(data_file.file_path,
220-
GetJsonValue<std::string>(json, kFilePath));
221-
ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue<std::string>(json, kFileFormat));
222-
ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str));
223-
224-
ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue<int32_t>(json, kSpecId));
225-
data_file.partition_spec_id = spec_id;
226-
227-
ICEBERG_ASSIGN_OR_RAISE(auto partition_vals,
228-
GetJsonValue<nlohmann::json>(json, kPartition));
229-
if (!partition_vals.is_array()) {
230-
return JsonParseError("PartitionValues must be a JSON array: {}",
231-
SafeDumpJson(partition_vals));
232-
}
233-
std::vector<Literal> literals;
234-
auto it = partition_spec_by_id.find(spec_id);
235-
if (it == partition_spec_by_id.end()) {
236-
return JsonParseError("Invalid partition spec id: {}", spec_id);
237-
}
238-
ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema));
239-
auto fields = struct_type->fields();
240-
if (partition_vals.size() != fields.size()) {
241-
return JsonParseError("Invalid partition data size: expected = {}, actual = {}",
242-
fields.size(), partition_vals.size());
243-
}
244-
for (size_t pos = 0; pos < fields.size(); ++pos) {
245-
ICEBERG_ASSIGN_OR_RAISE(
246-
auto literal, LiteralFromJson(partition_vals[pos], fields[pos].type().get()));
247-
literals.push_back(std::move(literal));
248-
}
249-
data_file.partition = PartitionValues(std::move(literals));
250-
251-
ICEBERG_ASSIGN_OR_RAISE(data_file.record_count,
252-
GetJsonValue<int64_t>(json, kRecordCount));
253-
ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes,
254-
GetJsonValue<int64_t>(json, kFileSizeInBytes));
255-
256-
ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes,
257-
KeyValueMapFromJson<int64_t>(json, kColumnSizes));
258-
ICEBERG_ASSIGN_OR_RAISE(data_file.value_counts,
259-
KeyValueMapFromJson<int64_t>(json, kValueCounts));
260-
ICEBERG_ASSIGN_OR_RAISE(data_file.null_value_counts,
261-
KeyValueMapFromJson<int64_t>(json, kNullValueCounts));
262-
ICEBERG_ASSIGN_OR_RAISE(data_file.nan_value_counts,
263-
KeyValueMapFromJson<int64_t>(json, kNanValueCounts));
264-
ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds,
265-
KeyValueMapFromJson<std::vector<uint8_t>>(json, kLowerBounds));
266-
ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds,
267-
KeyValueMapFromJson<std::vector<uint8_t>>(json, kUpperBounds));
268-
269-
if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) {
270-
ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata,
271-
GetJsonValue<std::vector<uint8_t>>(json, kKeyMetadata));
272-
}
273-
if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) {
274-
ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets,
275-
GetJsonValue<std::vector<int64_t>>(json, kSplitOffsets));
276-
}
277-
if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) {
278-
ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids,
279-
GetJsonValue<std::vector<int32_t>>(json, kEqualityIds));
280-
}
281-
if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) {
282-
ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id,
283-
GetJsonValue<int32_t>(json, kSortOrderId));
284-
}
285-
if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) {
286-
ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id,
287-
GetJsonValue<int64_t>(json, kFirstRowId));
288-
}
289-
if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) {
290-
ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file,
291-
GetJsonValue<std::string>(json, kReferencedDataFile));
292-
}
293-
if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) {
294-
ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset,
295-
GetJsonValue<int64_t>(json, kContentOffset));
296-
}
297-
if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) {
298-
ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes,
299-
GetJsonValue<int64_t>(json, kContentSizeInBytes));
300-
}
301-
302-
return data_file;
136+
return iceberg::DataFileFromJson(json, partition_spec_by_id, schema);
303137
}
304138

305139
Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
@@ -357,97 +191,12 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
357191
return file_scan_tasks;
358192
}
359193

360-
Result<nlohmann::json> DataFileToJsonUnchecked(const DataFile& data_file) {
361-
nlohmann::json json;
362-
switch (data_file.content) {
363-
case DataFile::Content::kData:
364-
json[kContent] = kContentData;
365-
break;
366-
case DataFile::Content::kPositionDeletes:
367-
json[kContent] = kContentPositionDeletes;
368-
break;
369-
case DataFile::Content::kEqualityDeletes:
370-
json[kContent] = kContentEqualityDeletes;
371-
break;
372-
}
373-
json[kFilePath] = data_file.file_path;
374-
json[kFileFormat] = ToString(data_file.file_format);
375-
376-
if (!data_file.partition_spec_id.has_value()) {
377-
return ValidationFailed("Cannot serialize REST content file without 'spec-id'");
378-
}
379-
json[kSpecId] = data_file.partition_spec_id.value();
380-
381-
nlohmann::json partition_json = nlohmann::json::array();
382-
for (const auto& literal : data_file.partition.values()) {
383-
ICEBERG_ASSIGN_OR_RAISE(auto lit_json, iceberg::ToJson(literal));
384-
partition_json.push_back(std::move(lit_json));
385-
}
386-
json[kPartition] = std::move(partition_json);
387-
388-
json[kRecordCount] = data_file.record_count;
389-
json[kFileSizeInBytes] = data_file.file_size_in_bytes;
390-
391-
SetKeyValueMap(json, kColumnSizes, data_file.column_sizes);
392-
SetKeyValueMap(json, kValueCounts, data_file.value_counts);
393-
SetKeyValueMap(json, kNullValueCounts, data_file.null_value_counts);
394-
SetKeyValueMap(json, kNanValueCounts, data_file.nan_value_counts);
395-
SetKeyValueMap(json, kLowerBounds, data_file.lower_bounds);
396-
SetKeyValueMap(json, kUpperBounds, data_file.upper_bounds);
397-
398-
if (!data_file.key_metadata.empty()) {
399-
json[kKeyMetadata] = data_file.key_metadata;
400-
}
401-
if (!data_file.split_offsets.empty()) {
402-
json[kSplitOffsets] = data_file.split_offsets;
403-
}
404-
if (!data_file.equality_ids.empty()) {
405-
json[kEqualityIds] = data_file.equality_ids;
406-
}
407-
if (data_file.sort_order_id.has_value()) {
408-
json[kSortOrderId] = data_file.sort_order_id.value();
409-
}
410-
if (data_file.first_row_id.has_value()) {
411-
json[kFirstRowId] = data_file.first_row_id.value();
412-
}
413-
if (data_file.referenced_data_file.has_value()) {
414-
json[kReferencedDataFile] = data_file.referenced_data_file.value();
415-
}
416-
if (data_file.content_offset.has_value()) {
417-
json[kContentOffset] = data_file.content_offset.value();
418-
}
419-
if (data_file.content_size_in_bytes.has_value()) {
420-
json[kContentSizeInBytes] = data_file.content_size_in_bytes.value();
421-
}
422-
423-
return json;
424-
}
425-
426194
Result<nlohmann::json> ToJson(
427195
const DataFile& data_file,
428196
const std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>>&
429197
partition_specs_by_id,
430198
const Schema& schema) {
431-
if (!data_file.partition_spec_id.has_value()) {
432-
return ValidationFailed("Invalid partition spec id from content file: null");
433-
}
434-
auto it = partition_specs_by_id.find(data_file.partition_spec_id.value());
435-
if (it == partition_specs_by_id.end() || !it->second) {
436-
return ValidationFailed("Invalid partition spec: null");
437-
}
438-
if (data_file.partition_spec_id.value() != it->second->spec_id()) {
439-
return ValidationFailed(
440-
"Invalid partition spec id from content file: expected = {}, actual = {}",
441-
it->second->spec_id(), data_file.partition_spec_id.value());
442-
}
443-
ICEBERG_ASSIGN_OR_RAISE(auto partition_type, it->second->PartitionType(schema));
444-
if (data_file.partition.num_fields() != partition_type->fields().size()) {
445-
return ValidationFailed(
446-
"Invalid partition data from content file: expected = {}, actual = {}",
447-
partition_type->fields().empty() ? "unpartitioned" : "partitioned",
448-
data_file.partition.num_fields() == 0 ? "unpartitioned" : "partitioned");
449-
}
450-
return DataFileToJsonUnchecked(data_file);
199+
return iceberg::ToJson(data_file, partition_specs_by_id, schema);
451200
}
452201

453202
namespace {

0 commit comments

Comments
 (0)