Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 42 additions & 3 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 @@ -132,9 +133,24 @@ 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 kStart = "start";
constexpr std::string_view kLength = "length";
constexpr std::string_view kMapKeys = "keys";
constexpr std::string_view kMapValues = "values";

Result<int64_t> GetRangeValue(const nlohmann::json& json, std::string_view key) {
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());
nlohmann::json json;
Expand Down Expand Up @@ -323,6 +339,17 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
GetJsonValue<nlohmann::json>(task_json, kDataFile));
ICEBERG_ASSIGN_OR_RAISE(
auto data_file, DataFileFromJson(data_file_json, partition_spec_by_id, schema));
const bool has_start = task_json.contains(kStart);
const bool has_length = task_json.contains(kLength);
if (has_start != has_length) {
return JsonParseError("'start' and 'length' must be provided together");
}
int64_t start = 0;
int64_t length = data_file.file_size_in_bytes;
if (has_start) {
ICEBERG_ASSIGN_OR_RAISE(start, GetRangeValue(task_json, kStart));
ICEBERG_ASSIGN_OR_RAISE(length, GetRangeValue(task_json, kLength));
}
// 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 @@ -350,9 +377,13 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> FileScanTasksFromJson(
ICEBERG_ASSIGN_OR_RAISE(residual_filter, ExpressionFromJson(filter_json));
}

file_scan_tasks.push_back(std::make_shared<FileScanTask>(
std::make_shared<DataFile>(std::move(data_file)), std::move(task_delete_files),
std::move(residual_filter)));
auto task = FileScanTask::MakeSplit(std::make_shared<DataFile>(std::move(data_file)),
start, length, std::move(task_delete_files),
std::move(residual_filter));
if (!task) {
return JsonParseError("Invalid FileScanTask range: {}", task.error().message);
}
file_scan_tasks.push_back(std::move(task.value()));
}
return file_scan_tasks;
}
Expand Down Expand Up @@ -493,10 +524,18 @@ Result<nlohmann::json> ScanTaskFieldsToJson(
if (!task) continue;
nlohmann::json task_json;
if (task->data_file()) {
auto range_result =
FileScanTask::MakeSplit(task->data_file(), task->start(), task->length());
if (!range_result) {
return ValidationFailed("Invalid FileScanTask range: {}",
range_result.error().message);
}
ICEBERG_ASSIGN_OR_RAISE(
auto data_file_json,
ToJson(*task->data_file(), partition_specs_by_id, schema));
task_json[kDataFile] = std::move(data_file_json);
task_json[kStart] = task->start();
task_json[kLength] = task->length();
}
if (!task->delete_files().empty()) {
std::vector<int32_t> refs;
Expand Down
1 change: 1 addition & 0 deletions src/iceberg/catalog/rest/types.cc
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,7 @@ bool FileScanTaskEqual(const std::shared_ptr<FileScanTask>& lhs,
return false;
}
return SharedPtrEqual(lhs->data_file(), rhs->data_file()) &&
lhs->start() == rhs->start() && lhs->length() == rhs->length() &&
SharedPtrVectorEqual(lhs->delete_files(), rhs->delete_files()) &&
ExpressionPtrEqual(lhs->residual_filter(), rhs->residual_filter());
}
Expand Down
6 changes: 6 additions & 0 deletions src/iceberg/data/file_scan_task_reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,12 @@ class FileScanTaskReader::Impl {
ICEBERG_PRECHECK(data_file->file_size_in_bytes >= 0,
"Data file size must not be negative: {}",
data_file->file_size_in_bytes);
if (task.is_split()) {
return NotSupported(
"Reading a partial FileScanTask is not supported yet: {} "
"(start={}, length={})",
data_file->file_path, task.start(), task.length());
}

if (task.delete_files().empty()) {
auto options = MakeReaderOptions(
Expand Down
3 changes: 3 additions & 0 deletions src/iceberg/data/file_scan_task_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,9 @@ class ICEBERG_DATA_EXPORT FileScanTaskReader {
~FileScanTaskReader();

/// \brief Open a task and return an Arrow C stream for its projected live rows.
///
/// Partial tasks are rejected until range-aware reading is implemented. The task and
/// its shared data/delete-file metadata must not be mutated while the stream uses them.
Result<ArrowArrayStream> Open(const FileScanTask& task);

FileScanTaskReader(const FileScanTaskReader&) = delete;
Expand Down
70 changes: 68 additions & 2 deletions src/iceberg/table_scan.cc
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
#include "iceberg/table.h"
#include "iceberg/table_metadata.h"
#include "iceberg/util/content_file_util.h"
#include "iceberg/util/int128.h"
#include "iceberg/util/macros.h"
#include "iceberg/util/snapshot_util.h"
#include "iceberg/util/timepoint.h"
Expand Down Expand Up @@ -271,11 +272,76 @@ FileScanTask::FileScanTask(std::shared_ptr<DataFile> data_file,
ICEBERG_DCHECK(data_file_ != nullptr, "Data file cannot be null for FileScanTask");
}

int64_t FileScanTask::size_bytes() const { return data_file_->file_size_in_bytes; }
FileScanTask::FileScanTask(std::shared_ptr<DataFile> data_file,
std::vector<std::shared_ptr<DataFile>> delete_files,
std::shared_ptr<Expression> residual_filter, Range range)
: FileScanTask(std::move(data_file), std::move(delete_files),
std::move(residual_filter)) {
range_ = range;
}

Result<std::shared_ptr<FileScanTask>> FileScanTask::MakeSplit(
std::shared_ptr<DataFile> data_file, int64_t start, int64_t length,
std::vector<std::shared_ptr<DataFile>> delete_files,
std::shared_ptr<Expression> filter) {
if (data_file == nullptr) {
return InvalidArgument("Cannot create a file scan task without a data file");
}

const int64_t file_size = data_file->file_size_in_bytes;
if (file_size < 0) {
return InvalidArgument("Cannot split file {} with negative size {}",
data_file->file_path, file_size);
}
if (start < 0 || start > file_size) {
return InvalidArgument("Split start {} is outside file {} of size {}", start,
data_file->file_path, file_size);
}
// Subtraction is safe after validating start, even when file_size is INT64_MAX.
if (length < 0 || length > file_size - start) {
return InvalidArgument("Split length {} from start {} is outside file {} of size {}",
length, start, data_file->file_path, file_size);
}
if (start == 0 && length == file_size) {
return std::make_shared<FileScanTask>(std::move(data_file), std::move(delete_files),
std::move(filter));
}
if (length == 0) {
return InvalidArgument("A partial split of file {} must have positive length",
data_file->file_path);
}

return std::shared_ptr<FileScanTask>(
new FileScanTask(std::move(data_file), std::move(delete_files), std::move(filter),
Range{.start = start, .length = length, .file_size = file_size}));
}

int64_t FileScanTask::length() const {
return range_ ? range_->length : data_file_->file_size_in_bytes;
}

int64_t FileScanTask::size_bytes() const { return length(); }

int32_t FileScanTask::files_count() const { return 1; }

int64_t FileScanTask::estimated_row_count() const { return data_file_->record_count; }
int64_t FileScanTask::estimated_row_count() const {
if (!range_) {
return data_file_->record_count;
}

const int64_t record_count = data_file_->record_count;
if (record_count <= 0) {
return 0;
}

// Cumulative boundaries make estimates additive across contiguous splits, and
// 128-bit products avoid overflow for any valid int64 file size and row count.
const auto end = static_cast<int128_t>(range_->start) + range_->length;
const auto count_at_end = end * record_count / range_->file_size;
const auto count_at_start =
static_cast<int128_t>(range_->start) * record_count / range_->file_size;
return static_cast<int64_t>(count_at_end - count_at_start);
}

// ChangelogScanTask implementation

Expand Down
34 changes: 34 additions & 0 deletions src/iceberg/table_scan.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
/// \file iceberg/table_scan.h
/// \brief Define table scan APIs and scan task types.

#include <cstdint>
#include <functional>
#include <memory>
#include <optional>
Expand Down Expand Up @@ -71,10 +72,32 @@ class ICEBERG_EXPORT FileScanTask : public ScanTask {
/// \param data_file The data file to read.
/// \param delete_files Delete files that apply to this data file.
/// \param filter Optional residual filter to apply after reading.
/// \note The data file, delete files, and residual filter must not be mutated while
/// a planning or reading stream uses this task.
explicit FileScanTask(std::shared_ptr<DataFile> data_file,
std::vector<std::shared_ptr<DataFile>> delete_files = {},
std::shared_ptr<Expression> filter = nullptr);

/// \brief Make a task for a validated half-open byte range [start, start + length).
///
/// A range covering the entire file is represented as a whole-file task. An empty
/// range is valid only for an empty whole file. The data file, delete files, and
/// residual filter are shared with the resulting task and must not be mutated
/// while a planning or reading stream uses it.
static Result<std::shared_ptr<FileScanTask>> MakeSplit(
std::shared_ptr<DataFile> data_file, int64_t start, int64_t length,
std::vector<std::shared_ptr<DataFile>> delete_files = {},
std::shared_ptr<Expression> filter = nullptr);

/// \brief The inclusive byte offset where this task starts, or zero for a whole file.
int64_t start() const { return range_ ? range_->start : 0; }

/// \brief The byte length of this task, or the current size of a whole file.
int64_t length() const;

/// \brief Whether this task reads a proper subset of its data file.
bool is_split() const { return range_.has_value(); }

/// \brief The data file that should be read by this scan task.
const std::shared_ptr<DataFile>& data_file() const { return data_file_; }

Expand All @@ -92,9 +115,20 @@ class ICEBERG_EXPORT FileScanTask : public ScanTask {
int64_t estimated_row_count() const override;

private:
struct Range {
int64_t start;
int64_t length;
int64_t file_size;
};

FileScanTask(std::shared_ptr<DataFile> data_file,
std::vector<std::shared_ptr<DataFile>> delete_files,
std::shared_ptr<Expression> filter, Range range);

std::shared_ptr<DataFile> data_file_;
std::vector<std::shared_ptr<DataFile>> delete_files_;
std::shared_ptr<Expression> residual_filter_;
std::optional<Range> range_;
};

/// \brief Stream of file scan tasks.
Expand Down
19 changes: 19 additions & 0 deletions src/iceberg/test/file_scan_task_reader_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -409,6 +409,25 @@ TEST_F(FileScanTaskReaderTest, OpenWithoutDeletesReadsProjectedSchema) {
VerifyStream(&stream, R"([[1, "Foo"], [2, "Bar"], [3, "Baz"]])"));
}

TEST_F(FileScanTaskReaderTest, RejectsPartialTaskUntilRangeReadsAreSupported) {
ICEBERG_UNWRAP_OR_FAIL(
auto data_file,
MakeDataFile(table_schema_,
R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])"));
ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTask::MakeSplit(data_file, 0, 1));

FileScanTaskReader::Options options{
.io = file_io_,
.table_schema = table_schema_,
.schemas = {table_schema_},
.projected_schema = projected_schema_,
};
ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options)));
auto result = reader->Open(*task);
EXPECT_THAT(result, IsError(ErrorKind::kNotSupported));
EXPECT_THAT(result, HasErrorMessage("Reading a partial FileScanTask"));
}

TEST_F(FileScanTaskReaderTest, ReadLastUpdatedFromDataSeq) {
ICEBERG_UNWRAP_OR_FAIL(
auto data_file,
Expand Down
Loading
Loading