Skip to content

Commit 01447df

Browse files
feat(serde): add core FileScanTask JSON support
1 parent d38395c commit 01447df

5 files changed

Lines changed: 640 additions & 1 deletion

File tree

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

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
*/
1919

2020
#include <iterator>
21+
#include <limits>
2122
#include <map>
2223
#include <memory>
2324
#include <optional>
@@ -108,6 +109,25 @@ constexpr std::string_view kPlanTask = "plan-task";
108109
constexpr std::string_view kDataFile = "data-file";
109110
constexpr std::string_view kDeleteFileReferences = "delete-file-references";
110111
constexpr std::string_view kResidualFilter = "residual-filter";
112+
constexpr std::string_view kStart = "start";
113+
constexpr std::string_view kLength = "length";
114+
115+
Result<int64_t> GetRangeValueOrDefault(const nlohmann::json& json, std::string_view key,
116+
int64_t default_value) {
117+
if (!json.contains(key) || json.at(key).is_null()) {
118+
return default_value;
119+
}
120+
const auto& value = json.at(key);
121+
if (!value.is_number_integer()) {
122+
return JsonParseError("'{}' must be an integer, but is {}", key, value.type_name());
123+
}
124+
if (value.is_number_unsigned() &&
125+
value.get<uint64_t>() >
126+
static_cast<uint64_t>(std::numeric_limits<int64_t>::max())) {
127+
return JsonParseError("'{}' integer is out of range: {}", key, SafeDumpJson(value));
128+
}
129+
return value.get<int64_t>();
130+
}
111131

112132
Result<nlohmann::json> StorageCredentialToJson(const StorageCredential& credential) {
113133
ICEBERG_RETURN_UNEXPECTED(credential.Validate());
@@ -158,6 +178,12 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
158178
ICEBERG_ASSIGN_OR_RAISE(
159179
auto data_file,
160180
iceberg::rest::DataFileFromJson(data_file_json, partition_spec_by_id, schema));
181+
ICEBERG_ASSIGN_OR_RAISE(auto start, GetRangeValueOrDefault(task_json, kStart, 0));
182+
ICEBERG_ASSIGN_OR_RAISE(
183+
auto length,
184+
GetRangeValueOrDefault(task_json, kLength, data_file.file_size_in_bytes));
185+
ICEBERG_RETURN_UNEXPECTED(
186+
CheckFileScanTaskNotSplit(start, length, data_file.file_size_in_bytes));
161187
// FIXME: REST scan-task DataFile JSON currently carries first-row-id,
162188
// but not the manifest-entry data sequence number. Until the REST API exposes
163189
// it, REST-planned tasks cannot inherit _last_updated_sequence_number.

‎src/iceberg/json_serde.cc‎

Lines changed: 150 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@
3131
#include <nlohmann/json.hpp>
3232

3333
#include "iceberg/constants.h"
34+
#include "iceberg/expression/expressions.h"
3435
#include "iceberg/expression/json_serde_internal.h"
3536
#include "iceberg/expression/literal.h"
3637
#include "iceberg/file_format.h"
@@ -50,6 +51,7 @@
5051
#include "iceberg/table_metadata.h"
5152
#include "iceberg/table_properties.h"
5253
#include "iceberg/table_requirement.h"
54+
#include "iceberg/table_scan.h"
5355
#include "iceberg/table_update.h"
5456
#include "iceberg/transform.h"
5557
#include "iceberg/type.h"
@@ -248,7 +250,7 @@ constexpr std::string_view kRequirementAssertDefaultSortOrderID =
248250
constexpr std::string_view kLastAssignedFieldId = "last-assigned-field-id";
249251
constexpr std::string_view kLastAssignedPartitionId = "last-assigned-partition-id";
250252

251-
// DataFile JSON (Iceberg REST ContentFile field names)
253+
// DataFile / FileScanTask JSON (Iceberg REST ContentFile field names)
252254
constexpr std::string_view kContent = "content";
253255
constexpr std::string_view kContentData = "data";
254256
constexpr std::string_view kContentPositionDeletes = "position-deletes";
@@ -269,6 +271,14 @@ constexpr std::string_view kEqualityIds = "equality-ids";
269271
constexpr std::string_view kReferencedDataFile = "referenced-data-file";
270272
constexpr std::string_view kContentOffset = "content-offset";
271273
constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes";
274+
constexpr std::string_view kDataFile = "data-file";
275+
constexpr std::string_view kDeleteFiles = "delete-files";
276+
constexpr std::string_view kDeleteFileReferences = "delete-file-references";
277+
constexpr std::string_view kResidualFilter = "residual-filter";
278+
constexpr std::string_view kTaskType = "task-type";
279+
constexpr std::string_view kFileScanTaskType = "file-scan-task";
280+
constexpr std::string_view kStart = "start";
281+
constexpr std::string_view kLength = "length";
272282
constexpr std::string_view kMapKeys = "keys";
273283
constexpr std::string_view kMapValues = "values";
274284

@@ -2432,4 +2442,143 @@ Result<nlohmann::json> ToJson(
24322442
return DataFileToJsonUnchecked(data_file);
24332443
}
24342444

2445+
Result<nlohmann::json> ToJson(
2446+
const FileScanTask& task,
2447+
const std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>>&
2448+
partition_specs_by_id,
2449+
const Schema& schema) {
2450+
if (!task.data_file()) {
2451+
return ValidationFailed("Cannot serialize FileScanTask without data-file");
2452+
}
2453+
if (task.data_file()->content != DataFile::Content::kData) {
2454+
return ValidationFailed("FileScanTask data-file must have data content");
2455+
}
2456+
2457+
nlohmann::json json;
2458+
json[kTaskType] = kFileScanTaskType;
2459+
ICEBERG_ASSIGN_OR_RAISE(auto schema_json, ToJson(schema));
2460+
json[kSchema] = std::move(schema_json);
2461+
2462+
ICEBERG_ASSIGN_OR_RAISE(auto data_file_json,
2463+
ToJson(*task.data_file(), partition_specs_by_id, schema));
2464+
const auto spec_id = task.data_file()->partition_spec_id.value();
2465+
json[kSpec] = ToJson(*partition_specs_by_id.at(spec_id));
2466+
json[kDataFile] = std::move(data_file_json);
2467+
json[kStart] = 0;
2468+
json[kLength] = task.data_file()->file_size_in_bytes;
2469+
2470+
nlohmann::json delete_files_json = nlohmann::json::array();
2471+
for (const auto& delete_file : task.delete_files()) {
2472+
if (!delete_file) {
2473+
return ValidationFailed("FileScanTask delete-files must not contain null");
2474+
}
2475+
if (delete_file->content == DataFile::Content::kData) {
2476+
return ValidationFailed("FileScanTask delete-file must have delete content");
2477+
}
2478+
if (delete_file->partition_spec_id != task.data_file()->partition_spec_id) {
2479+
return ValidationFailed(
2480+
"Invalid partition spec id from content file: expected = {}, actual = {}",
2481+
spec_id,
2482+
delete_file->partition_spec_id.has_value()
2483+
? std::to_string(delete_file->partition_spec_id.value())
2484+
: "null");
2485+
}
2486+
ICEBERG_ASSIGN_OR_RAISE(auto delete_file_json,
2487+
ToJson(*delete_file, partition_specs_by_id, schema));
2488+
delete_files_json.push_back(std::move(delete_file_json));
2489+
}
2490+
json[kDeleteFiles] = std::move(delete_files_json);
2491+
2492+
if (task.residual_filter()) {
2493+
ICEBERG_ASSIGN_OR_RAISE(auto residual_json, ToJson(*task.residual_filter()));
2494+
json[kResidualFilter] = std::move(residual_json);
2495+
}
2496+
return json;
2497+
}
2498+
2499+
Status CheckFileScanTaskNotSplit(int64_t start, int64_t length,
2500+
int64_t file_size_in_bytes) {
2501+
if (start != 0 || length != file_size_in_bytes) {
2502+
return NotSupported(
2503+
"Split FileScanTask is not supported: start={}, length={}, "
2504+
"file-size-in-bytes={}",
2505+
start, length, file_size_in_bytes);
2506+
}
2507+
return {};
2508+
}
2509+
2510+
Result<std::shared_ptr<FileScanTask>> FileScanTaskFromJson(const nlohmann::json& json) {
2511+
if (!json.is_object()) {
2512+
return JsonParseError("Cannot parse file scan task from a non-object: {}",
2513+
SafeDumpJson(json));
2514+
}
2515+
ICEBERG_ASSIGN_OR_RAISE(auto task_type,
2516+
GetJsonValueOptional<std::string>(json, kTaskType));
2517+
if (task_type.has_value() &&
2518+
!StringUtils::EqualsIgnoreCase(*task_type, kFileScanTaskType)) {
2519+
return JsonParseError("Unsupported scan task type: {}", *task_type);
2520+
}
2521+
if (json.contains(kDeleteFileReferences) && !json.at(kDeleteFileReferences).is_null()) {
2522+
return JsonParseError(
2523+
"Cannot parse FileScanTask with 'delete-file-references'; use "
2524+
"rest::FileScanTasksFromJson for REST scan responses");
2525+
}
2526+
2527+
ICEBERG_ASSIGN_OR_RAISE(auto schema_json, GetJsonValue<nlohmann::json>(json, kSchema));
2528+
ICEBERG_ASSIGN_OR_RAISE(auto schema, SchemaFromJson(schema_json));
2529+
auto shared_schema = std::shared_ptr<Schema>(std::move(schema));
2530+
2531+
ICEBERG_ASSIGN_OR_RAISE(auto spec_json, GetJsonValue<nlohmann::json>(json, kSpec));
2532+
ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonInteger<int32_t>(spec_json, kSpecId));
2533+
ICEBERG_ASSIGN_OR_RAISE(auto spec,
2534+
PartitionSpecFromJson(shared_schema, spec_json, spec_id));
2535+
auto shared_spec = std::shared_ptr<PartitionSpec>(std::move(spec));
2536+
const std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>> partition_spec_by_id{
2537+
{shared_spec->spec_id(), shared_spec}};
2538+
2539+
ICEBERG_ASSIGN_OR_RAISE(auto data_file_json,
2540+
GetJsonValue<nlohmann::json>(json, kDataFile));
2541+
ICEBERG_ASSIGN_OR_RAISE(
2542+
auto data_file,
2543+
DataFileFromJson(data_file_json, partition_spec_by_id, *shared_schema));
2544+
if (data_file.content != DataFile::Content::kData) {
2545+
return JsonParseError("FileScanTask data-file must have data content");
2546+
}
2547+
2548+
ICEBERG_ASSIGN_OR_RAISE(auto start, GetJsonInteger<int64_t>(json, kStart));
2549+
ICEBERG_ASSIGN_OR_RAISE(auto length, GetJsonInteger<int64_t>(json, kLength));
2550+
ICEBERG_RETURN_UNEXPECTED(
2551+
CheckFileScanTaskNotSplit(start, length, data_file.file_size_in_bytes));
2552+
2553+
std::vector<std::shared_ptr<DataFile>> delete_files;
2554+
if (json.contains(kDeleteFiles)) {
2555+
ICEBERG_ASSIGN_OR_RAISE(auto delete_files_json,
2556+
GetJsonValue<nlohmann::json>(json, kDeleteFiles));
2557+
if (!delete_files_json.is_array()) {
2558+
return JsonParseError("Cannot parse delete files from non-array: {}",
2559+
SafeDumpJson(delete_files_json));
2560+
}
2561+
for (const auto& delete_file_json : delete_files_json) {
2562+
ICEBERG_ASSIGN_OR_RAISE(
2563+
auto delete_file,
2564+
DataFileFromJson(delete_file_json, partition_spec_by_id, *shared_schema));
2565+
if (delete_file.content == DataFile::Content::kData) {
2566+
return JsonParseError("FileScanTask delete-file must have delete content");
2567+
}
2568+
delete_files.push_back(std::make_shared<DataFile>(std::move(delete_file)));
2569+
}
2570+
}
2571+
2572+
std::shared_ptr<Expression> residual_filter = Expressions::AlwaysTrue();
2573+
if (json.contains(kResidualFilter)) {
2574+
ICEBERG_ASSIGN_OR_RAISE(auto filter_json,
2575+
GetJsonValue<nlohmann::json>(json, kResidualFilter));
2576+
ICEBERG_ASSIGN_OR_RAISE(residual_filter, ExpressionFromJson(filter_json));
2577+
}
2578+
2579+
return std::make_shared<FileScanTask>(std::make_shared<DataFile>(std::move(data_file)),
2580+
std::move(delete_files),
2581+
std::move(residual_filter));
2582+
}
2583+
24352584
} // namespace iceberg

‎src/iceberg/json_serde_internal.h‎

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -445,4 +445,36 @@ ICEBERG_EXPORT Result<DataFile> DataFileFromJson(
445445
partition_spec_by_id,
446446
const Schema& schema);
447447

448+
/// \brief Serializes a `FileScanTask` to a self-contained JSON object.
449+
///
450+
/// Unlike REST scan responses, delete files are inlined under `delete-files` rather
451+
/// than encoded as `delete-file-references` into a sibling array.
452+
///
453+
/// JSON fields:
454+
/// - `task-type`: `file-scan-task`
455+
/// - `schema` (required): schema used by this task
456+
/// - `spec` (required): partition spec used by this task
457+
/// - `data-file` (required): ContentFile JSON
458+
/// - `start` (required): always 0 because split tasks are unsupported
459+
/// - `length` (required): always the data file size
460+
/// - `delete-files` (required): array of ContentFile JSON
461+
/// - `residual-filter` (optional): Expression JSON
462+
ICEBERG_EXPORT Result<nlohmann::json> ToJson(
463+
const FileScanTask& task,
464+
const std::unordered_map<int32_t, std::shared_ptr<PartitionSpec>>&
465+
partition_specs_by_id,
466+
const Schema& schema);
467+
468+
/// \brief Deserializes a self-contained FileScanTask JSON object.
469+
///
470+
/// Reads the embedded schema and partition spec. Rejects REST
471+
/// `delete-file-references` and byte-range splits (`start` != 0 or `length` != file
472+
/// size). Both range fields are required, matching Java core.
473+
ICEBERG_EXPORT Result<std::shared_ptr<FileScanTask>> FileScanTaskFromJson(
474+
const nlohmann::json& json);
475+
476+
/// Rejects `start != 0 || length != file_size_in_bytes`.
477+
ICEBERG_EXPORT Status CheckFileScanTaskNotSplit(int64_t start, int64_t length,
478+
int64_t file_size_in_bytes);
479+
448480
} // namespace iceberg

0 commit comments

Comments
 (0)