diff --git a/src/paimon/core/operation/merge_file_split_read.cpp b/src/paimon/core/operation/merge_file_split_read.cpp index 2e517cdff..9b913edfb 100644 --- a/src/paimon/core/operation/merge_file_split_read.cpp +++ b/src/paimon/core/operation/merge_file_split_read.cpp @@ -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" @@ -78,6 +79,185 @@ struct KeyValue; template class MergeFunctionWrapper; +namespace { + +class SortMergeKeyValueRecordReader : public KeyValueRecordReader { + public: + explicit SortMergeKeyValueRecordReader(std::unique_ptr&& reader) + : reader_(std::move(reader)) {} + + class Iterator : public KeyValueRecordReader::Iterator { + public: + explicit Iterator(std::unique_ptr&& iterator) + : iterator_(std::move(iterator)) {} + + Result HasNext() const override { + return iterator_->HasNext(); + } + + Result Next() override { + return std::move(iterator_->Next()); + } + + private: + std::unique_ptr iterator_; + }; + + Result> NextBatch() override { + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr iterator, + reader_->NextBatch()); + if (!iterator) { + return std::unique_ptr(); + } + return std::make_unique(std::move(iterator)); + } + + void Close() override { + reader_->Close(); + } + + std::shared_ptr GetReaderMetrics() const override { + return reader_->GetReaderMetrics(); + } + + private: + std::unique_ptr reader_; +}; + +} // namespace + +class MergeFileSplitRead::RealtimeReaderBuilder { + public: + static Result> Create( + const std::vector>& disk_splits, + std::vector>&& additional_readers, + MergeFileSplitRead* owner) { + ScopeGuard additional_readers_guard([&additional_readers]() { + for (const std::unique_ptr& reader : additional_readers) { + if (reader) { + reader->Close(); + } + } + }); + RealtimeReaderBuilder builder(owner); + std::vector> 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& 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>& disk_splits, + std::vector>* readers) { + std::shared_ptr first_split; + std::vector> data_files; + std::vector> deletion_files; + for (const std::shared_ptr& disk_split : disk_splits) { + std::shared_ptr data_split = + std::dynamic_pointer_cast(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>& split_files = data_split->DataFiles(); + const std::vector>& 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 data_file_path_factory, + owner_->path_factory_->CreateDataFilePathFactory(partition, bucket)); + + DeletionVector::Factory dv_factory; + std::vector> 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> section_readers; + ScopeGuard section_readers_guard([§ion_readers]() { + for (const std::unique_ptr& reader : section_readers) { + reader->Close(); + } + }); + section_readers.reserve(disk_sections.size()); + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr> merge_function_wrapper, + MergeFileSplitRead::CreateMergeFunctionWrapper(owner_->options_, + owner_->context_->GetTableSchema(), + owner_->value_schema_, owner_->pool_)); + for (const std::vector& section : disk_sections) { + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr 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(std::move(section_reader))); + } + std::unique_ptr concat_reader = + std::make_unique(std::move(section_readers)); + section_readers_guard.Release(); + readers->push_back(std::move(concat_reader)); + return Status::OK(); + } + + Result> CreateMergedReader( + std::vector>&& record_readers) { + ScopeGuard record_readers_guard([&record_readers]() { + for (const std::unique_ptr& reader : record_readers) { + if (reader) { + reader->Close(); + } + } + }); + if (record_readers.empty()) { + record_readers_guard.Release(); + return std::make_unique(std::vector>{}, + owner_->pool_); + } + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr 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 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> MergeFileSplitRead::Create( const std::shared_ptr& path_factory, const std::shared_ptr& context, @@ -158,6 +338,12 @@ Result> MergeFileSplitRead::CreateReader( return std::make_unique(std::move(batch_reader), pool_); } +Result> MergeFileSplitRead::CreateRealtimeReader( + const std::vector>& disk_splits, + std::vector>&& additional_readers) { + return RealtimeReaderBuilder::Create(disk_splits, std::move(additional_readers), this); +} + void MergeFileSplitRead::SetMergeFunctionWrapper( const std::shared_ptr>& merge_function_wrapper) { merge_function_wrapper_ = merge_function_wrapper; @@ -236,13 +422,10 @@ Result> MergeFileSplitRead::ApplyIndexAndDvRead Result> MergeFileSplitRead::CreateMergeReader( const std::shared_ptr& data_split, const std::shared_ptr& data_file_path_factory) { - auto dv_factory = DeletionVector::CreateFactory( - options_.GetFileSystem(), - DeletionVector::CreateDeletionFileMap(data_split->DataFiles(), data_split->DeletionFiles()), - pool_); - - std::vector> sections = - IntervalPartition(data_split->DataFiles(), key_comparator_).Partition(); + DeletionVector::Factory dv_factory; + std::vector> sections; + PAIMON_RETURN_NOT_OK(CreateDiskSections(data_split->DataFiles(), data_split->DeletionFiles(), + &dv_factory, §ions)); std::vector> batch_readers; batch_readers.reserve(sections.size()); // no overlap through multiple sections @@ -453,38 +636,100 @@ Result> MergeFileSplitRead::CreateReaderForSection( } else { predicate = context_->GetPredicate(); } - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr 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( - std::move(sort_merge_reader), raw_read_schema_, projection_, options_.GetReadBatchSize(), - thread_number, pool_); + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr 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> MergeFileSplitRead::CreateSortMergeReaderForSection( +Status MergeFileSplitRead::CreateDiskSections( + const std::vector>& data_files, + const std::vector>& deletion_files, + DeletionVector::Factory* dv_factory, std::vector>* 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>> +MergeFileSplitRead::CreateRecordReadersForSection( const std::vector& section, const BinaryRow& partition, DeletionVector::Factory dv_factory, const std::shared_ptr& predicate, - const std::shared_ptr& data_file_path_factory, bool drop_delete) { - // with overlap in one section + const std::shared_ptr& data_file_path_factory) const { std::vector> 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 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 sort_merge_reader, - CreateSortMergeReader(std::move(record_readers))); + return record_readers; +} + +Result> MergeFileSplitRead::CreateProjectedReader( + std::unique_ptr&& sort_merge_reader, + const std::shared_ptr& predicate, bool complete_row_kind) { + if (!force_keep_delete_) { + sort_merge_reader = std::make_unique(std::move(sort_merge_reader)); + } + // KeyValueProjectionReader converts KeyValue objects to arrow array according to projection + std::unique_ptr 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( + 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 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(std::move(projection_reader), pool_); + } + return projection_reader; +} + +Result> MergeFileSplitRead::CreateSortMergeReaderForSection( + const std::vector& section, const BinaryRow& partition, + DeletionVector::Factory dv_factory, const std::shared_ptr& predicate, + const std::shared_ptr& data_file_path_factory, bool drop_delete) { + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr> merge_function_wrapper, + GetMergeFunctionWrapper()); + return CreateSortMergeReaderForSection(section, partition, dv_factory, predicate, + data_file_path_factory, drop_delete, + merge_function_wrapper); +} + +Result> MergeFileSplitRead::CreateSortMergeReaderForSection( + const std::vector& section, const BinaryRow& partition, + DeletionVector::Factory dv_factory, const std::shared_ptr& predicate, + const std::shared_ptr& data_file_path_factory, bool drop_delete, + const std::shared_ptr>& merge_function_wrapper) { + // with overlap in one section + PAIMON_ASSIGN_OR_RAISE(std::vector> record_readers, + CreateRecordReadersForSection(section, partition, dv_factory, predicate, + data_file_path_factory)); + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr sort_merge_reader, + CreateSortMergeReader(std::move(record_readers), merge_function_wrapper)); if (drop_delete) { sort_merge_reader = std::make_unique(std::move(sort_merge_reader)); } @@ -536,6 +781,12 @@ Result> MergeFileSplitRead::CreateSortMergeRead std::vector>&& record_readers) { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr> merge_function_wrapper, GetMergeFunctionWrapper()); + return CreateSortMergeReader(std::move(record_readers), merge_function_wrapper); +} + +Result> MergeFileSplitRead::CreateSortMergeReader( + std::vector>&& record_readers, + const std::shared_ptr>& merge_function_wrapper) const { auto sort_engine = options_.GetSortEngine(); if (sort_engine == SortEngine::MIN_HEAP) { return std::make_unique( diff --git a/src/paimon/core/operation/merge_file_split_read.h b/src/paimon/core/operation/merge_file_split_read.h index bda6e2401..07b5e70b2 100644 --- a/src/paimon/core/operation/merge_file_split_read.h +++ b/src/paimon/core/operation/merge_file_split_read.h @@ -123,10 +123,20 @@ class MergeFileSplitRead : public AbstractSplitRead { return value_schema_; } + std::shared_ptr GetKeySchema() const { + return key_schema_; + } + + Result> CreateRealtimeReader( + const std::vector>& disk_splits, + std::vector>&& additional_readers); + void SetMergeFunctionWrapper( const std::shared_ptr>& merge_function_wrapper); private: + class RealtimeReaderBuilder; + Result> CreateMergeReader( const std::shared_ptr& data_split, const std::shared_ptr& data_file_path_factory); @@ -140,6 +150,26 @@ class MergeFileSplitRead : public AbstractSplitRead { DeletionVector::Factory dv_factory, const std::shared_ptr& data_file_path_factory); + Status CreateDiskSections(const std::vector>& data_files, + const std::vector>& deletion_files, + DeletionVector::Factory* dv_factory, + std::vector>* sections) const; + + Result>> CreateRecordReadersForSection( + const std::vector& section, const BinaryRow& partition, + DeletionVector::Factory dv_factory, const std::shared_ptr& predicate, + const std::shared_ptr& data_file_path_factory) const; + + Result> CreateProjectedReader( + std::unique_ptr&& sort_merge_reader, + const std::shared_ptr& predicate, bool complete_row_kind); + + Result> CreateSortMergeReaderForSection( + const std::vector& section, const BinaryRow& partition, + DeletionVector::Factory dv_factory, const std::shared_ptr& predicate, + const std::shared_ptr& data_file_path_factory, bool drop_delete, + const std::shared_ptr>& merge_function_wrapper); + Result> CreateReaderForRun( const BinaryRow& partition, const SortedRun& sorted_run, DeletionVector::Factory dv_factory, const std::shared_ptr& predicate, @@ -148,6 +178,10 @@ class MergeFileSplitRead : public AbstractSplitRead { Result> CreateSortMergeReader( std::vector>&& record_readers); + Result> CreateSortMergeReader( + std::vector>&& record_readers, + const std::shared_ptr>& merge_function_wrapper) const; + Result>> GetMergeFunctionWrapper(); MergeFileSplitRead(const std::shared_ptr& path_factory, diff --git a/src/paimon/core/operation/merge_file_split_read_test.cpp b/src/paimon/core/operation/merge_file_split_read_test.cpp index f857d4876..d02120de4 100644 --- a/src/paimon/core/operation/merge_file_split_read_test.cpp +++ b/src/paimon/core/operation/merge_file_split_read_test.cpp @@ -40,6 +40,7 @@ #include "paimon/common/utils/scope_guard.h" #include "paimon/core/core_options.h" #include "paimon/core/io/data_file_meta.h" +#include "paimon/core/io/key_value_in_memory_record_reader.h" #include "paimon/core/manifest/file_source.h" #include "paimon/core/operation/internal_read_context.h" #include "paimon/core/schema/schema_manager.h" @@ -51,7 +52,6 @@ #include "paimon/executor.h" #include "paimon/fs/local/local_file_system.h" #include "paimon/memory/memory_pool.h" -#include "paimon/metrics.h" #include "paimon/predicate/literal.h" #include "paimon/predicate/predicate_builder.h" #include "paimon/read_context.h" @@ -66,6 +66,7 @@ class FileSystem; } // namespace paimon namespace paimon::test { + // Parameter: min_heap/loser_tree; enable/disable IO prefetch; enable/disable multi thread row to // batch class MergeFileSplitReadTest : public ::testing::Test, @@ -328,9 +329,8 @@ class MergeFileSplitReadTest : public ::testing::Test, return {data_split1}; } - Result> CreateReader( - const std::shared_ptr& internal_context, - const std::vector>& data_splits) { + Result> CreateMergeFileSplitRead( + const std::shared_ptr& internal_context) { const auto& core_options = internal_context->GetCoreOptions(); const auto& table_schema = internal_context->GetTableSchema(); auto arrow_schema = DataField::ConvertDataFieldsToArrowSchema(table_schema->Fields()); @@ -347,9 +347,14 @@ class MergeFileSplitReadTest : public ::testing::Test, core_options.DataFilePrefix(), core_options.LegacyPartitionNameEnabled(), external_paths, global_index_external_path, core_options.IndexFileInDataFileDir(), pool_)); - PAIMON_ASSIGN_OR_RAISE(auto split_read, - MergeFileSplitRead::Create(path_factory, std::move(internal_context), - pool_, executor_)); + return MergeFileSplitRead::Create(path_factory, internal_context, pool_, executor_); + } + + Result> CreateReader( + const std::shared_ptr& internal_context, + const std::vector>& data_splits) { + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr split_read, + CreateMergeFileSplitRead(internal_context)); std::vector> batch_readers; batch_readers.reserve(data_splits.size()); for (const auto& split : data_splits) { @@ -666,6 +671,73 @@ TEST_P(MergeFileSplitReadTest, TestSimple) { CheckResult(result_array, expected_array, read_schema); } +TEST_P(MergeFileSplitReadTest, TestRealtimeReadConcatenatesOrderedDiskSections) { + std::string path = + paimon::test::GetDataDir() + "/parquet/pk_table_with_mor.db/pk_table_with_mor"; + ReadContextBuilder context_builder(path); + std::vector raw_read_fields = {DataField(0, arrow::field("k0", arrow::int32())), + DataField(1, arrow::field("k1", arrow::int32())), + DataField(5, arrow::field("s1", arrow::utf8())), + DataField(6, arrow::field("v0", arrow::float64()))}; + std::shared_ptr read_schema = + DataField::ConvertDataFieldsToArrowSchema(raw_read_fields); + ASSERT_TRUE(read_schema); + + context_builder.SetReadFieldNames({"k0", "k1", "s1", "v0"}); + context_builder.SetOptions( + {{Options::SEQUENCE_FIELD, "s0,s1"}, {Options::MERGE_ENGINE, "deduplicate"}}); + AddOptions(&context_builder); + ASSERT_OK_AND_ASSIGN(std::shared_ptr read_context, context_builder.Finish()); + std::shared_ptr internal_context = CreateInternalReadContext(read_context); + ASSERT_OK_AND_ASSIGN(std::unique_ptr split_read, + CreateMergeFileSplitRead(internal_context)); + + std::shared_ptr memory_type = + arrow::struct_(split_read->GetValueSchema()->fields()); + std::shared_ptr memory_array = + std::dynamic_pointer_cast( + arrow::ipc::internal::json::ArrayFromJSON(memory_type, R"([ + [100, 200, "memory-late", 10000.0, "zzzz"], + [1, 1, "memory-delete", 1100.0, "zzzz"], + [0, 0, "memory-first", 1000.0, "zzzz"], + [50, 0, "memory-middle", 5000.0, "zzzz"] + ])") + .ValueOrDie()); + std::vector> memory_readers; + memory_readers.push_back(std::make_unique( + /*last_sequence_num=*/9, memory_array, + std::vector( + {RecordBatch::RowKind::UPDATE_AFTER, RecordBatch::RowKind::DELETE, + RecordBatch::RowKind::UPDATE_AFTER, RecordBatch::RowKind::INSERT}), + std::vector({"k0", "k1"}), std::vector({"s0", "s1"}), + /*sequence_fields_ascending=*/true, split_read->GetKeyComparator(), pool_)); + + std::vector> disk_splits = {PrepareDataSplit().front()}; + ASSERT_OK_AND_ASSIGN(std::unique_ptr batch_reader, + split_read->CreateRealtimeReader(disk_splits, std::move(memory_readers))); + ASSERT_OK_AND_ASSIGN(std::shared_ptr result_array, + ReadResultCollector::CollectResult(batch_reader.get())); + + arrow::FieldVector fields_with_row_kind = read_schema->fields(); + fields_with_row_kind.insert(fields_with_row_kind.begin(), + arrow::field("_VALUE_KIND", arrow::int8())); + std::shared_ptr expected_array; + auto expected_status = + arrow::ipc::internal::json::ChunkedArrayFromJSON(arrow::struct_(fields_with_row_kind), {R"([ + [0, 0, 0, "memory-first", 1000.0], + [0, 0, 1, "you", 11.1], + [0, 1, 0, "later", 12.2], + [0, 1, 2, "!", 13.3], + [0, 50, 0, "memory-middle", 5000.0], + [0, 100, 200, "memory-late", 10000.0] + ])"}, + &expected_array); + ASSERT_TRUE(expected_status.ok()); + CheckResult(result_array, expected_array, read_schema); + ASSERT_TRUE(batch_reader->GetReaderMetrics()); + batch_reader->Close(); +} + TEST_P(MergeFileSplitReadTest, TestLookUp) { std::string path = paimon::test::GetDataDir() + "/parquet/pk_table_with_mor.db/pk_table_with_mor"; diff --git a/src/paimon/core/table/source/key_value_table_read.cpp b/src/paimon/core/table/source/key_value_table_read.cpp index 208807493..afb852a6b 100644 --- a/src/paimon/core/table/source/key_value_table_read.cpp +++ b/src/paimon/core/table/source/key_value_table_read.cpp @@ -19,13 +19,31 @@ #include "paimon/core/table/source/key_value_table_read.h" +#include #include +#include +#include "arrow/api.h" +#include "arrow/c/bridge.h" +#include "paimon/common/reader/concat_batch_reader.h" +#include "paimon/common/table/special_fields.h" +#include "paimon/common/utils/arrow/status_utils.h" +#include "paimon/common/utils/scope_guard.h" #include "paimon/core/global_index/indexed_split_impl.h" +#include "paimon/core/io/merged_key_value_record_reader.h" +#include "paimon/core/key_value.h" +#include "paimon/core/mergetree/compact/merge_function.h" +#include "paimon/core/mergetree/compact/reducer_merge_function_wrapper.h" #include "paimon/core/operation/merge_file_split_read.h" #include "paimon/core/operation/raw_file_split_read.h" +#include "paimon/core/realtime/realtime_context_impl.h" +#include "paimon/core/realtime/realtime_primary_key_reader.h" +#include "paimon/core/realtime/realtime_reader.h" #include "paimon/core/table/source/data_split_impl.h" #include "paimon/core/table/source/pk_count_reader.h" +#include "paimon/core/table/source/realtime_split.h" +#include "paimon/core/utils/nested_projection_utils.h" +#include "paimon/core/utils/primary_key_table_utils.h" #include "paimon/status.h" namespace paimon { @@ -34,16 +52,79 @@ class Executor; class FileStorePathFactory; class InternalReadContext; class MemoryPool; +struct ColumnarBatchContext; -KeyValueTableRead::KeyValueTableRead(std::vector>&& split_reads, - const std::shared_ptr& path_factory, - const std::shared_ptr& context, - const std::shared_ptr& memory_pool, - const std::shared_ptr& executor) +namespace { + +Result> CreateRealtimePrimaryKeyQueryTransportSchema( + const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema) { + arrow::FieldVector transport_value_fields; + transport_value_fields.reserve(key_schema->num_fields() + value_schema->num_fields()); + std::unordered_set field_ids; + for (const std::shared_ptr& field : key_schema->fields()) { + PAIMON_ASSIGN_OR_RAISE(int32_t field_id, NestedProjectionUtils::GetPaimonFieldId(field)); + if (field_ids.insert(field_id).second) { + transport_value_fields.push_back(field); + } + } + for (const std::shared_ptr& field : value_schema->fields()) { + PAIMON_ASSIGN_OR_RAISE(int32_t field_id, NestedProjectionUtils::GetPaimonFieldId(field)); + if (field_ids.insert(field_id).second) { + transport_value_fields.push_back(field); + } + } + return RealtimePrimaryKeyLayout::CreateSchema(transport_value_fields); +} + +Result>> CreateMemoryReaders( + const std::shared_ptr& split, const RealtimePartitionBucketView& memory, + const std::shared_ptr& transport_schema, + const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema, + const std::shared_ptr& key_comparator, + const std::shared_ptr& context, + const std::shared_ptr& memory_pool) { + auto c_schema = std::make_unique(); + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*transport_schema, c_schema.get())); + ScopeGuard schema_guard([schema = c_schema.get()]() { ArrowSchemaRelease(schema); }); + RealtimeQueryContext query_context{c_schema.get(), nullptr, false}; + PAIMON_ASSIGN_OR_RAISE(std::vector> batch_readers, + memory.store->CreateQueryReaders(memory.read_view, 0, query_context)); + PAIMON_ASSIGN_OR_RAISE( + std::vector> realtime_primary_key_readers, + RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::move(batch_readers), transport_schema, + OffsetRange(split->CommittedEndOffset(), split->MemoryEndOffset()), key_schema, + value_schema, memory_pool)); + std::vector> result; + result.reserve(realtime_primary_key_readers.size()); + for (std::unique_ptr& realtime_primary_key_reader : + realtime_primary_key_readers) { + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr merge, + PrimaryKeyTableUtils::CreateMergeFunction( + value_schema, context->GetTableSchema()->PrimaryKeys(), + context->GetCoreOptions(), memory_pool)); + result.push_back(std::make_unique( + std::move(realtime_primary_key_reader), key_comparator, + std::make_shared(std::move(merge)))); + } + return result; +} + +} // namespace + +KeyValueTableRead::KeyValueTableRead( + std::vector>&& split_reads, + const std::shared_ptr& path_factory, + const std::shared_ptr& context, + const std::shared_ptr& realtime_primary_key_transport_schema, + const std::shared_ptr& memory_pool, const std::shared_ptr& executor) : TableRead(memory_pool), split_reads_(std::move(split_reads)), path_factory_(path_factory), context_(context), + realtime_primary_key_transport_schema_(realtime_primary_key_transport_schema), executor_(executor) {} Result> KeyValueTableRead::Create( @@ -57,10 +138,18 @@ Result> KeyValueTableRead::Create( PAIMON_ASSIGN_OR_RAISE( std::unique_ptr merge_file_split_read, MergeFileSplitRead::Create(path_factory, context, memory_pool, executor)); + std::shared_ptr realtime_primary_key_transport_schema; + if (context->GetRealtimeContext()) { + PAIMON_ASSIGN_OR_RAISE( + realtime_primary_key_transport_schema, + CreateRealtimePrimaryKeyQueryTransportSchema(merge_file_split_read->GetKeySchema(), + merge_file_split_read->GetValueSchema())); + } split_reads.emplace_back(std::move(merge_file_split_read)); - return std::unique_ptr(new KeyValueTableRead(std::move(split_reads), path_factory, - context, memory_pool, executor)); + return std::unique_ptr( + new KeyValueTableRead(std::move(split_reads), path_factory, context, + realtime_primary_key_transport_schema, memory_pool, executor)); } void KeyValueTableRead::ForceKeepDelete(bool force_keep_delete) { @@ -75,6 +164,11 @@ void KeyValueTableRead::ForceKeepDelete(bool force_keep_delete) { Result> KeyValueTableRead::CreateReader( const std::shared_ptr& split) { + std::shared_ptr realtime_split = std::dynamic_pointer_cast(split); + if (realtime_split) { + return CreateRealtimeReader(realtime_split, /*release_ticket=*/true); + } + std::shared_ptr dispatch_split = split; if (auto indexed_split = std::dynamic_pointer_cast(split)) { PAIMON_RETURN_NOT_OK(indexed_split->Validate()); @@ -126,8 +220,104 @@ Result> KeyValueTableRead::CreateReader( return Status::Invalid("create reader failed, not read match with data split."); } +Result> KeyValueTableRead::CreateReader( + const std::vector>& splits) { + std::vector> readers; + readers.reserve(splits.size()); + std::vector> realtime_splits; + ScopeGuard cleanup_guard([&]() { + for (const std::unique_ptr& reader : readers) { + if (reader) { + reader->Close(); + } + } + }); + for (const std::shared_ptr& split : splits) { + std::shared_ptr realtime_split = + std::dynamic_pointer_cast(split); + if (realtime_split) { + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr reader, + CreateRealtimeReader(realtime_split, + /*release_ticket=*/false)); + readers.push_back(std::move(reader)); + realtime_splits.push_back(std::move(realtime_split)); + } else { + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr reader, CreateReader(split)); + readers.push_back(std::move(reader)); + } + } + + if (!realtime_splits.empty()) { + const std::shared_ptr realtime_context = context_->GetRealtimeContext(); + if (!realtime_context) { + return Status::Invalid("reading a real-time split requires a real-time context"); + } + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr realtime_context_impl, + RealtimeContextImpl::Cast(realtime_context)); + for (const std::shared_ptr& realtime_split : realtime_splits) { + PAIMON_RETURN_NOT_OK( + realtime_context_impl->ReleaseReadView(realtime_split->OpaqueTicket())); + } + } + return std::make_unique(std::move(readers), GetMemoryPool()); +} + +Result> KeyValueTableRead::CreateRealtimeReader( + const std::shared_ptr& realtime_split, bool release_ticket) { + if (realtime_split->Version() != RealtimeSplit::kCurrentVersion) { + return Status::Invalid("unsupported real-time split version"); + } + if (realtime_split->MemoryEndOffset() < realtime_split->CommittedEndOffset()) { + return Status::Invalid("real-time split memory end offset precedes committed end offset"); + } + const std::shared_ptr realtime_context = context_->GetRealtimeContext(); + if (!realtime_context) { + return Status::Invalid("reading a real-time split requires a real-time context"); + } + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr realtime_context_impl, + RealtimeContextImpl::Cast(realtime_context)); + PAIMON_ASSIGN_OR_RAISE(RealtimePartitionBucketView memory, + realtime_context_impl->ResolveReadView(realtime_split->OpaqueTicket())); + const RealtimePartitionBucket expected_partition_bucket(realtime_split->Partition(), + realtime_split->Bucket()); + if (memory.partition_bucket != expected_partition_bucket) { + return Status::Invalid("real-time read-view ticket belongs to another partition-bucket"); + } + const std::optional memory_range = memory.read_view->GetOffsetRange(); + if (!memory_range || memory_range->end != realtime_split->MemoryEndOffset()) { + return Status::Invalid("real-time read-view ticket does not match the split offset range"); + } + for (const std::unique_ptr& read : split_reads_) { + auto* merge_read = dynamic_cast(read.get()); + if (merge_read) { + PAIMON_ASSIGN_OR_RAISE( + std::vector> memory_readers, + CreateMemoryReaders(realtime_split, memory, realtime_primary_key_transport_schema_, + merge_read->GetKeySchema(), merge_read->GetValueSchema(), + merge_read->GetKeyComparator(), context_, GetMemoryPool())); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr reader, + merge_read->CreateRealtimeReader(realtime_split->DiskSplits(), + std::move(memory_readers))); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr realtime_reader, + RealtimeReader::Create(memory.read_view, std::move(reader))); + if (release_ticket) { + PAIMON_RETURN_NOT_OK( + realtime_context_impl->ReleaseReadView(realtime_split->OpaqueTicket())); + } + return std::unique_ptr(std::move(realtime_reader)); + } + } + return Status::Invalid("create reader failed, merge file split read not found"); +} + Result> KeyValueTableRead::CreateCountReader( const std::vector>& splits) { + for (const std::shared_ptr& split : splits) { + if (std::dynamic_pointer_cast(split)) { + return Status::NotImplemented( + "CreateCountReader does not support process-local real-time splits"); + } + } if (context_->GetPredicate() != nullptr) { return Status::NotImplemented( "CreateCountReader with predicate pushdown is not supported yet"); diff --git a/src/paimon/core/table/source/key_value_table_read.h b/src/paimon/core/table/source/key_value_table_read.h index d6a1c83d3..1dd59b016 100644 --- a/src/paimon/core/table/source/key_value_table_read.h +++ b/src/paimon/core/table/source/key_value_table_read.h @@ -22,6 +22,7 @@ #include #include +#include "arrow/type_fwd.h" #include "paimon/core/operation/internal_read_context.h" #include "paimon/core/operation/split_read.h" #include "paimon/core/utils/file_store_path_factory.h" @@ -35,6 +36,7 @@ class Executor; class FileStorePathFactory; class InternalReadContext; class MemoryPool; +class RealtimeSplit; class KeyValueTableRead : public TableRead { public: @@ -45,6 +47,9 @@ class KeyValueTableRead : public TableRead { Result> CreateReader(const std::shared_ptr& split) override; + Result> CreateReader( + const std::vector>& splits) override; + Result> CreateCountReader( const std::vector>& splits) override; @@ -54,12 +59,17 @@ class KeyValueTableRead : public TableRead { KeyValueTableRead(std::vector>&& split_reads, const std::shared_ptr& path_factory, const std::shared_ptr& context, + const std::shared_ptr& realtime_primary_key_transport_schema, const std::shared_ptr& memory_pool, const std::shared_ptr& executor); + Result> CreateRealtimeReader( + const std::shared_ptr& realtime_split, bool release_ticket); + std::vector> split_reads_; std::shared_ptr path_factory_; std::shared_ptr context_; + std::shared_ptr realtime_primary_key_transport_schema_; std::shared_ptr executor_; bool force_keep_delete_ = false; };