Skip to content
Draft
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
305 changes: 278 additions & 27 deletions src/paimon/core/operation/merge_file_split_read.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
#include "paimon/common/types/data_field.h"
#include "paimon/common/utils/arrow/status_utils.h"
#include "paimon/common/utils/object_utils.h"
#include "paimon/common/utils/scope_guard.h"
#include "paimon/core/core_options.h"
#include "paimon/core/deletionvectors/apply_deletion_vector_batch_reader.h"
#include "paimon/core/deletionvectors/bitmap_deletion_vector.h"
Expand Down Expand Up @@ -78,6 +79,185 @@ struct KeyValue;
template <typename T>
class MergeFunctionWrapper;

namespace {

class SortMergeKeyValueRecordReader : public KeyValueRecordReader {
public:
explicit SortMergeKeyValueRecordReader(std::unique_ptr<SortMergeReader>&& reader)
: reader_(std::move(reader)) {}

class Iterator : public KeyValueRecordReader::Iterator {
public:
explicit Iterator(std::unique_ptr<SortMergeReader::Iterator>&& iterator)
: iterator_(std::move(iterator)) {}

Result<bool> HasNext() const override {
return iterator_->HasNext();
}

Result<KeyValue> Next() override {
return std::move(iterator_->Next());
}

private:
std::unique_ptr<SortMergeReader::Iterator> iterator_;
};

Result<std::unique_ptr<KeyValueRecordReader::Iterator>> NextBatch() override {
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<SortMergeReader::Iterator> iterator,
reader_->NextBatch());
if (!iterator) {
return std::unique_ptr<KeyValueRecordReader::Iterator>();
}
return std::make_unique<Iterator>(std::move(iterator));
}

void Close() override {
reader_->Close();
}

std::shared_ptr<Metrics> GetReaderMetrics() const override {
return reader_->GetReaderMetrics();
}

private:
std::unique_ptr<SortMergeReader> reader_;
};

} // namespace

class MergeFileSplitRead::RealtimeReaderBuilder {
public:
static Result<std::unique_ptr<BatchReader>> Create(
const std::vector<std::shared_ptr<Split>>& disk_splits,
std::vector<std::unique_ptr<KeyValueRecordReader>>&& additional_readers,
MergeFileSplitRead* owner) {
ScopeGuard additional_readers_guard([&additional_readers]() {
for (const std::unique_ptr<KeyValueRecordReader>& reader : additional_readers) {
if (reader) {
reader->Close();
}
}
});
RealtimeReaderBuilder builder(owner);
std::vector<std::unique_ptr<KeyValueRecordReader>> readers;
if (!disk_splits.empty()) {
PAIMON_RETURN_NOT_OK(builder.CollectDiskReaders(disk_splits, &readers));
}
readers.reserve(readers.size() + additional_readers.size());
for (std::unique_ptr<KeyValueRecordReader>& additional_reader : additional_readers) {
readers.push_back(std::move(additional_reader));
}
additional_readers_guard.Release();
return builder.CreateMergedReader(std::move(readers));
}

private:
explicit RealtimeReaderBuilder(MergeFileSplitRead* owner) : owner_(owner) {}

Status CollectDiskReaders(const std::vector<std::shared_ptr<Split>>& disk_splits,
std::vector<std::unique_ptr<KeyValueRecordReader>>* readers) {
std::shared_ptr<DataSplitImpl> first_split;
std::vector<std::shared_ptr<DataFileMeta>> data_files;
std::vector<std::optional<DeletionFile>> deletion_files;
for (const std::shared_ptr<Split>& disk_split : disk_splits) {
std::shared_ptr<DataSplitImpl> data_split =
std::dynamic_pointer_cast<DataSplitImpl>(disk_split);
if (!data_split) {
return Status::Invalid("merge input disk split is not a data split");
}
if (!first_split) {
first_split = data_split;
}
const std::vector<std::shared_ptr<DataFileMeta>>& split_files = data_split->DataFiles();
const std::vector<std::optional<DeletionFile>>& split_deletion_files =
data_split->DeletionFiles();
if (!split_deletion_files.empty() &&
split_deletion_files.size() != split_files.size()) {
return Status::Invalid(
"merge input disk split deletion files must be empty or match data files");
}
data_files.insert(data_files.end(), split_files.begin(), split_files.end());
if (split_deletion_files.empty()) {
deletion_files.insert(deletion_files.end(), split_files.size(), std::nullopt);
} else {
deletion_files.insert(deletion_files.end(), split_deletion_files.begin(),
split_deletion_files.end());
}
}
const BinaryRow& partition = first_split->Partition();
const int32_t bucket = first_split->Bucket();
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<DataFilePathFactory> data_file_path_factory,
owner_->path_factory_->CreateDataFilePathFactory(partition, bucket));

DeletionVector::Factory dv_factory;
std::vector<std::vector<SortedRun>> disk_sections;
PAIMON_RETURN_NOT_OK(
owner_->CreateDiskSections(data_files, deletion_files, &dv_factory, &disk_sections));
if (disk_sections.empty()) {
return Status::OK();
}
std::vector<std::unique_ptr<KeyValueRecordReader>> section_readers;
ScopeGuard section_readers_guard([&section_readers]() {
for (const std::unique_ptr<KeyValueRecordReader>& reader : section_readers) {
reader->Close();
}
});
section_readers.reserve(disk_sections.size());
PAIMON_ASSIGN_OR_RAISE(
std::shared_ptr<MergeFunctionWrapper<KeyValue>> merge_function_wrapper,
MergeFileSplitRead::CreateMergeFunctionWrapper(owner_->options_,
owner_->context_->GetTableSchema(),
owner_->value_schema_, owner_->pool_));
for (const std::vector<SortedRun>& section : disk_sections) {
PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<SortMergeReader> section_reader,
owner_->CreateSortMergeReaderForSection(
section, partition, dv_factory, owner_->predicate_for_keys_,
data_file_path_factory, /*drop_delete=*/false, merge_function_wrapper));
section_readers.push_back(
std::make_unique<SortMergeKeyValueRecordReader>(std::move(section_reader)));
}
std::unique_ptr<KeyValueRecordReader> concat_reader =
std::make_unique<ConcatKeyValueRecordReader>(std::move(section_readers));
section_readers_guard.Release();
readers->push_back(std::move(concat_reader));
return Status::OK();
}

Result<std::unique_ptr<BatchReader>> CreateMergedReader(
std::vector<std::unique_ptr<KeyValueRecordReader>>&& record_readers) {
ScopeGuard record_readers_guard([&record_readers]() {
for (const std::unique_ptr<KeyValueRecordReader>& reader : record_readers) {
if (reader) {
reader->Close();
}
}
});
if (record_readers.empty()) {
record_readers_guard.Release();
return std::make_unique<ConcatBatchReader>(std::vector<std::unique_ptr<BatchReader>>{},
owner_->pool_);
}
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<SortMergeReader> sort_merge_reader,
owner_->CreateSortMergeReader(std::move(record_readers)));
record_readers_guard.Release();
ScopeGuard sort_merge_reader_guard([&sort_merge_reader]() {
if (sort_merge_reader) {
sort_merge_reader->Close();
}
});
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<BatchReader> result,
owner_->CreateProjectedReader(std::move(sort_merge_reader),
owner_->context_->GetPredicate(),
/*complete_row_kind=*/true));
sort_merge_reader_guard.Release();
return result;
}

MergeFileSplitRead* owner_;
};

Result<std::unique_ptr<MergeFileSplitRead>> MergeFileSplitRead::Create(
const std::shared_ptr<FileStorePathFactory>& path_factory,
const std::shared_ptr<InternalReadContext>& context,
Expand Down Expand Up @@ -158,6 +338,12 @@ Result<std::unique_ptr<BatchReader>> MergeFileSplitRead::CreateReader(
return std::make_unique<CompleteRowKindBatchReader>(std::move(batch_reader), pool_);
}

Result<std::unique_ptr<BatchReader>> MergeFileSplitRead::CreateRealtimeReader(
const std::vector<std::shared_ptr<Split>>& disk_splits,
std::vector<std::unique_ptr<KeyValueRecordReader>>&& additional_readers) {
return RealtimeReaderBuilder::Create(disk_splits, std::move(additional_readers), this);
}

void MergeFileSplitRead::SetMergeFunctionWrapper(
const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper) {
merge_function_wrapper_ = merge_function_wrapper;
Expand Down Expand Up @@ -236,13 +422,10 @@ Result<std::unique_ptr<FileBatchReader>> MergeFileSplitRead::ApplyIndexAndDvRead
Result<std::unique_ptr<BatchReader>> MergeFileSplitRead::CreateMergeReader(
const std::shared_ptr<DataSplitImpl>& data_split,
const std::shared_ptr<DataFilePathFactory>& data_file_path_factory) {
auto dv_factory = DeletionVector::CreateFactory(
options_.GetFileSystem(),
DeletionVector::CreateDeletionFileMap(data_split->DataFiles(), data_split->DeletionFiles()),
pool_);

std::vector<std::vector<SortedRun>> sections =
IntervalPartition(data_split->DataFiles(), key_comparator_).Partition();
DeletionVector::Factory dv_factory;
std::vector<std::vector<SortedRun>> sections;
PAIMON_RETURN_NOT_OK(CreateDiskSections(data_split->DataFiles(), data_split->DeletionFiles(),
&dv_factory, &sections));
std::vector<std::unique_ptr<BatchReader>> batch_readers;
batch_readers.reserve(sections.size());
// no overlap through multiple sections
Expand Down Expand Up @@ -453,38 +636,100 @@ Result<std::unique_ptr<BatchReader>> MergeFileSplitRead::CreateReaderForSection(
} else {
predicate = context_->GetPredicate();
}
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<SortMergeReader> sort_merge_reader,
CreateSortMergeReaderForSection(section, partition, dv_factory,
predicate, data_file_path_factory,
/*drop_delete=*/!force_keep_delete_));
// KeyValueProjectionReader converts KeyValue objects to arrow array according to projection
if (!context_->EnableMultiThreadRowToBatch()) {
return KeyValueProjectionReader::Create(std::move(sort_merge_reader), raw_read_schema_,
projection_, options_.GetReadBatchSize(), pool_);
}
int32_t thread_number = context_->GetRowToBatchThreadNumber();
assert(thread_number > 0);
return std::make_unique<AsyncKeyValueProjectionReader>(
std::move(sort_merge_reader), raw_read_schema_, projection_, options_.GetReadBatchSize(),
thread_number, pool_);
PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<SortMergeReader> sort_merge_reader,
CreateSortMergeReaderForSection(section, partition, dv_factory, predicate,
data_file_path_factory, /*drop_delete=*/false));
return CreateProjectedReader(std::move(sort_merge_reader), /*predicate=*/nullptr,
/*complete_row_kind=*/false);
}

Result<std::unique_ptr<SortMergeReader>> MergeFileSplitRead::CreateSortMergeReaderForSection(
Status MergeFileSplitRead::CreateDiskSections(
const std::vector<std::shared_ptr<DataFileMeta>>& data_files,
const std::vector<std::optional<DeletionFile>>& deletion_files,
DeletionVector::Factory* dv_factory, std::vector<std::vector<SortedRun>>* sections) const {
*dv_factory = DeletionVector::CreateFactory(
options_.GetFileSystem(), DeletionVector::CreateDeletionFileMap(data_files, deletion_files),
pool_);
*sections = IntervalPartition(data_files, key_comparator_).Partition();
return Status::OK();
}

Result<std::vector<std::unique_ptr<KeyValueRecordReader>>>
MergeFileSplitRead::CreateRecordReadersForSection(
const std::vector<SortedRun>& section, const BinaryRow& partition,
DeletionVector::Factory dv_factory, const std::shared_ptr<Predicate>& predicate,
const std::shared_ptr<DataFilePathFactory>& data_file_path_factory, bool drop_delete) {
// with overlap in one section
const std::shared_ptr<DataFilePathFactory>& data_file_path_factory) const {
std::vector<std::unique_ptr<KeyValueRecordReader>> record_readers;
record_readers.reserve(section.size());
for (const auto& run : section) {
for (const SortedRun& run : section) {
// no overlap in a run
PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<KeyValueRecordReader> run_reader,
CreateReaderForRun(partition, run, dv_factory, predicate, data_file_path_factory));
record_readers.emplace_back(std::move(run_reader));
}
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<SortMergeReader> sort_merge_reader,
CreateSortMergeReader(std::move(record_readers)));
return record_readers;
}

Result<std::unique_ptr<BatchReader>> MergeFileSplitRead::CreateProjectedReader(
std::unique_ptr<SortMergeReader>&& sort_merge_reader,
const std::shared_ptr<Predicate>& predicate, bool complete_row_kind) {
if (!force_keep_delete_) {
sort_merge_reader = std::make_unique<DropDeleteReader>(std::move(sort_merge_reader));
}
// KeyValueProjectionReader converts KeyValue objects to arrow array according to projection
std::unique_ptr<BatchReader> projection_reader;
if (!context_->EnableMultiThreadRowToBatch()) {
PAIMON_ASSIGN_OR_RAISE(
projection_reader,
KeyValueProjectionReader::Create(std::move(sort_merge_reader), raw_read_schema_,
projection_, options_.GetReadBatchSize(), pool_));
} else {
const int32_t thread_number = context_->GetRowToBatchThreadNumber();
assert(thread_number > 0);
projection_reader = std::make_unique<AsyncKeyValueProjectionReader>(
std::move(sort_merge_reader), raw_read_schema_, projection_,
options_.GetReadBatchSize(), thread_number, pool_);
}
ScopeGuard projection_reader_guard([&projection_reader]() {
if (projection_reader) {
projection_reader->Close();
}
});
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<BatchReader> filtered_reader,
ApplyPredicateFilterIfNeeded(std::move(projection_reader), predicate));
projection_reader_guard.Release();
projection_reader = std::move(filtered_reader);
if (complete_row_kind) {
return std::make_unique<CompleteRowKindBatchReader>(std::move(projection_reader), pool_);
}
return projection_reader;
}

Result<std::unique_ptr<SortMergeReader>> MergeFileSplitRead::CreateSortMergeReaderForSection(
const std::vector<SortedRun>& section, const BinaryRow& partition,
DeletionVector::Factory dv_factory, const std::shared_ptr<Predicate>& predicate,
const std::shared_ptr<DataFilePathFactory>& data_file_path_factory, bool drop_delete) {
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<MergeFunctionWrapper<KeyValue>> merge_function_wrapper,
GetMergeFunctionWrapper());
return CreateSortMergeReaderForSection(section, partition, dv_factory, predicate,
data_file_path_factory, drop_delete,
merge_function_wrapper);
}

Result<std::unique_ptr<SortMergeReader>> MergeFileSplitRead::CreateSortMergeReaderForSection(
const std::vector<SortedRun>& section, const BinaryRow& partition,
DeletionVector::Factory dv_factory, const std::shared_ptr<Predicate>& predicate,
const std::shared_ptr<DataFilePathFactory>& data_file_path_factory, bool drop_delete,
const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper) {
// with overlap in one section
PAIMON_ASSIGN_OR_RAISE(std::vector<std::unique_ptr<KeyValueRecordReader>> record_readers,
CreateRecordReadersForSection(section, partition, dv_factory, predicate,
data_file_path_factory));
PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<SortMergeReader> sort_merge_reader,
CreateSortMergeReader(std::move(record_readers), merge_function_wrapper));
if (drop_delete) {
sort_merge_reader = std::make_unique<DropDeleteReader>(std::move(sort_merge_reader));
}
Expand Down Expand Up @@ -536,6 +781,12 @@ Result<std::unique_ptr<SortMergeReader>> MergeFileSplitRead::CreateSortMergeRead
std::vector<std::unique_ptr<KeyValueRecordReader>>&& record_readers) {
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<MergeFunctionWrapper<KeyValue>> merge_function_wrapper,
GetMergeFunctionWrapper());
return CreateSortMergeReader(std::move(record_readers), merge_function_wrapper);
}

Result<std::unique_ptr<SortMergeReader>> MergeFileSplitRead::CreateSortMergeReader(
std::vector<std::unique_ptr<KeyValueRecordReader>>&& record_readers,
const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper) const {
auto sort_engine = options_.GetSortEngine();
if (sort_engine == SortEngine::MIN_HEAP) {
return std::make_unique<SortMergeReaderWithMinHeap>(
Expand Down
Loading
Loading