Skip to content
294 changes: 36 additions & 258 deletions src/iceberg/catalog/rest/json_serde.cc
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
*/

#include <iterator>
#include <limits>
#include <map>
#include <memory>
#include <optional>
Expand Down Expand Up @@ -105,35 +106,28 @@ 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";
constexpr std::string_view kStart = "start";
constexpr std::string_view kLength = "length";

Result<int64_t> GetRangeValueOrDefault(const nlohmann::json& json, std::string_view key,
int64_t default_value) {
if (!json.contains(key) || json.at(key).is_null()) {
return default_value;
}
const auto& value = json.at(key);
if (!value.is_number_integer()) {
return JsonParseError("'{}' must be an integer, but is {}", key, value.type_name());
}
if (value.is_number_unsigned() &&
value.get<uint64_t>() >
static_cast<uint64_t>(std::numeric_limits<int64_t>::max())) {
return JsonParseError("'{}' integer is out of range: {}", key, SafeDumpJson(value));
}
return value.get<int64_t>();
}

Result<nlohmann::json> StorageCredentialToJson(const StorageCredential& credential) {
ICEBERG_RETURN_UNEXPECTED(credential.Validate());
Expand All @@ -152,154 +146,14 @@ Result<StorageCredential> StorageCredentialFromJson(const nlohmann::json& json)
return credential;
}

template <typename Value>
Result<std::map<int32_t, Value>> KeyValueMapFromJson(const nlohmann::json& json,
std::string_view key) {
std::map<int32_t, Value> result;
if (!json.contains(key) || json.at(key).is_null()) {
return result;
}

ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue<nlohmann::json>(json, key));
ICEBERG_ASSIGN_OR_RAISE(auto keys,
GetJsonValue<std::vector<int32_t>>(map_json, kMapKeys));
ICEBERG_ASSIGN_OR_RAISE(auto values,
GetJsonValue<std::vector<Value>>(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 <typename Value>
void SetKeyValueMap(nlohmann::json& json, std::string_view key,
const std::map<int32_t, Value>& map) {
if (map.empty()) {
return;
}

std::vector<int32_t> keys;
std::vector<Value> 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<DataFile> DataFileFromJson(
const nlohmann::json& json,
const std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>>&
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<std::string>(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<std::string>(json, kFilePath));
ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue<std::string>(json, kFileFormat));
ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str));

ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue<int32_t>(json, kSpecId));
data_file.partition_spec_id = spec_id;

ICEBERG_ASSIGN_OR_RAISE(auto partition_vals,
GetJsonValue<nlohmann::json>(json, kPartition));
if (!partition_vals.is_array()) {
return JsonParseError("PartitionValues must be a JSON array: {}",
SafeDumpJson(partition_vals));
}
std::vector<Literal> 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<int64_t>(json, kRecordCount));
ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes,
GetJsonValue<int64_t>(json, kFileSizeInBytes));

ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes,
KeyValueMapFromJson<int64_t>(json, kColumnSizes));
ICEBERG_ASSIGN_OR_RAISE(data_file.value_counts,
KeyValueMapFromJson<int64_t>(json, kValueCounts));
ICEBERG_ASSIGN_OR_RAISE(data_file.null_value_counts,
KeyValueMapFromJson<int64_t>(json, kNullValueCounts));
ICEBERG_ASSIGN_OR_RAISE(data_file.nan_value_counts,
KeyValueMapFromJson<int64_t>(json, kNanValueCounts));
ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds,
KeyValueMapFromJson<std::vector<uint8_t>>(json, kLowerBounds));
ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds,
KeyValueMapFromJson<std::vector<uint8_t>>(json, kUpperBounds));

if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata,
GetJsonValue<std::vector<uint8_t>>(json, kKeyMetadata));
}
if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets,
GetJsonValue<std::vector<int64_t>>(json, kSplitOffsets));
}
if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids,
GetJsonValue<std::vector<int32_t>>(json, kEqualityIds));
}
if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id,
GetJsonValue<int32_t>(json, kSortOrderId));
}
if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id,
GetJsonValue<int64_t>(json, kFirstRowId));
}
if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file,
GetJsonValue<std::string>(json, kReferencedDataFile));
}
if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset,
GetJsonValue<int64_t>(json, kContentOffset));
}
if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) {
ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes,
GetJsonValue<int64_t>(json, kContentSizeInBytes));
}

return data_file;
return iceberg::DataFileFromJson(json, partition_spec_by_id, schema);
}

Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
Expand All @@ -322,7 +176,14 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
ICEBERG_ASSIGN_OR_RAISE(auto data_file_json,
GetJsonValue<nlohmann::json>(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));
ICEBERG_ASSIGN_OR_RAISE(auto start, GetRangeValueOrDefault(task_json, kStart, 0));
ICEBERG_ASSIGN_OR_RAISE(
auto length,
GetRangeValueOrDefault(task_json, kLength, data_file.file_size_in_bytes));
ICEBERG_RETURN_UNEXPECTED(
CheckFileScanTaskNotSplit(start, length, data_file.file_size_in_bytes));
// 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.
Expand Down Expand Up @@ -357,97 +218,12 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
return file_scan_tasks;
}

Result<nlohmann::json> 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<nlohmann::json> ToJson(
const DataFile& data_file,
const std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>>&
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 {
Expand Down Expand Up @@ -495,7 +271,7 @@ Result<nlohmann::json> 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()) {
Expand All @@ -521,7 +297,8 @@ Result<nlohmann::json> 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()) {
Expand Down Expand Up @@ -557,8 +334,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<DataFile>(std::move(delete_file)));
}

Expand Down
Loading
Loading