From 576a6e11bd3c7243706513f740adf5979f368ea3 Mon Sep 17 00:00:00 2001 From: JeffZhou <17023790+HaHaJeff@users.noreply.github.com> Date: Sat, 29 Aug 2026 21:33:20 +0800 Subject: [PATCH] feat(realtime): adapt primary-key transport readers --- src/paimon/CMakeLists.txt | 2 + .../merged_key_value_record_reader_test.cpp | 1 + .../primary_key_realtime_store_test.cpp | 85 +-- .../realtime/realtime_primary_key_reader.cpp | 510 +++++++++++++ .../realtime/realtime_primary_key_reader.h | 71 ++ .../realtime_primary_key_reader_test.cpp | 688 ++++++++++++++++++ 6 files changed, 1292 insertions(+), 65 deletions(-) create mode 100644 src/paimon/core/realtime/realtime_primary_key_reader.cpp create mode 100644 src/paimon/core/realtime/realtime_primary_key_reader.h create mode 100644 src/paimon/core/realtime/realtime_primary_key_reader_test.cpp diff --git a/src/paimon/CMakeLists.txt b/src/paimon/CMakeLists.txt index b8e5ec3e8..bf57ce9f8 100644 --- a/src/paimon/CMakeLists.txt +++ b/src/paimon/CMakeLists.txt @@ -385,6 +385,7 @@ set(PAIMON_CORE_SRCS core/operation/write_restore.cpp core/realtime/arrow_realtime_store.cpp core/realtime/arrow_realtime_store_factory.cpp + core/realtime/realtime_primary_key_reader.cpp core/realtime/primary_key_realtime_store.cpp core/realtime/realtime_append_only_writer.cpp core/realtime/realtime_context.cpp @@ -793,6 +794,7 @@ if(PAIMON_BUILD_TESTS) core/memory/writer_memory_manager_test.cpp core/realtime/arrow_realtime_store_test.cpp core/realtime/primary_key_realtime_store_test.cpp + core/realtime/realtime_primary_key_reader_test.cpp core/realtime/realtime_context_test.cpp core/realtime/realtime_reader_test.cpp core/mergetree/levels_test.cpp diff --git a/src/paimon/core/io/merged_key_value_record_reader_test.cpp b/src/paimon/core/io/merged_key_value_record_reader_test.cpp index 1b6b71c69..858e302e8 100644 --- a/src/paimon/core/io/merged_key_value_record_reader_test.cpp +++ b/src/paimon/core/io/merged_key_value_record_reader_test.cpp @@ -20,6 +20,7 @@ #include #include +#include #include "arrow/api.h" #include "arrow/array/array_nested.h" diff --git a/src/paimon/core/realtime/primary_key_realtime_store_test.cpp b/src/paimon/core/realtime/primary_key_realtime_store_test.cpp index 7d6fa72f5..ed80db275 100644 --- a/src/paimon/core/realtime/primary_key_realtime_store_test.cpp +++ b/src/paimon/core/realtime/primary_key_realtime_store_test.cpp @@ -18,7 +18,6 @@ #include "paimon/core/realtime/primary_key_realtime_store.h" -#include #include #include #include @@ -29,10 +28,10 @@ #include "arrow/api.h" #include "arrow/c/bridge.h" #include "arrow/ipc/json_simple.h" -#include "paimon/common/table/special_fields.h" #include "paimon/common/types/data_field.h" #include "paimon/common/utils/arrow/status_utils.h" #include "paimon/common/utils/checked_cast.h" +#include "paimon/core/realtime/realtime_primary_key_reader.h" #include "paimon/macros.h" #include "paimon/memory/memory_pool.h" #include "paimon/realtime/arrow_realtime_store_factory.h" @@ -50,28 +49,17 @@ std::shared_ptr FieldWithId(const std::string& name, } std::shared_ptr TransportSchema() { - return arrow::schema( - {DataField::ConvertDataFieldToArrowField(SpecialFields::ValueKind())->WithNullable(false), - DataField::ConvertDataFieldToArrowField(SpecialFields::SequenceNumber()) - ->WithNullable(false), - DataField::ConvertDataFieldToArrowField(SpecialFields::RealtimeOffset()), - DataField::ConvertDataFieldToArrowField(DataField(0, arrow::field("id", arrow::int64()))), - DataField::ConvertDataFieldToArrowField( - DataField(1, arrow::field("value", arrow::utf8())))}); + return RealtimePrimaryKeyLayout::CreateSchema( + {FieldWithId("id", arrow::int64(), 0), FieldWithId("value", arrow::utf8(), 1)}); } std::shared_ptr NestedTransportSchema() { - return arrow::schema( - {DataField::ConvertDataFieldToArrowField(SpecialFields::ValueKind())->WithNullable(false), - DataField::ConvertDataFieldToArrowField(SpecialFields::SequenceNumber()) - ->WithNullable(false), - DataField::ConvertDataFieldToArrowField(SpecialFields::RealtimeOffset()), - DataField::ConvertDataFieldToArrowField(DataField(0, arrow::field("id", arrow::int64()))), - DataField::ConvertDataFieldToArrowField(DataField( - 1, - arrow::field("value", - arrow::struct_({arrow::field("name", arrow::utf8()), - arrow::field("items", arrow::list(arrow::int32()))}))))}); + return RealtimePrimaryKeyLayout::CreateSchema( + {FieldWithId("id", arrow::int64(), 0), + FieldWithId("value", + arrow::struct_({arrow::field("name", arrow::utf8()), + arrow::field("items", arrow::list(arrow::int32()))}), + 1)}); } std::unique_ptr MakeBatch(const std::string& json) { @@ -125,36 +113,6 @@ Result ReadJson(const std::vector>& re return result->ToString(); } -class TestingMemoryPool final : public MemoryPool { - public: - void* Malloc(uint64_t size, uint64_t alignment) override { - return delegate_->Malloc(size, alignment); - } - - void* Realloc(void* pointer, size_t old_size, size_t new_size, uint64_t alignment) override { - return delegate_->Realloc(pointer, old_size, new_size, alignment); - } - - void Free(void* pointer, uint64_t size) override { - delegate_->Free(pointer, size); - } - - void Free(void* pointer, uint64_t size, uint64_t alignment) override { - delegate_->Free(pointer, size, alignment); - } - - uint64_t CurrentUsage() const override { - return delegate_->CurrentUsage(); - } - - uint64_t MaxMemoryUsage() const override { - return delegate_->MaxMemoryUsage(); - } - - private: - std::unique_ptr delegate_ = GetMemoryPool(); -}; - TEST(PrimaryKeyRealtimeStoreTest, TestWriteAndSealValidation) { ASSERT_OK_AND_ASSIGN(std::shared_ptr store, PrimaryKeyRealtimeStore::Create(TransportSchema(), GetDefaultPool())); @@ -337,8 +295,8 @@ TEST(PrimaryKeyRealtimeStoreTest, TestQueryReaderPerStoredBatch) { TEST(PrimaryKeyRealtimeStoreTest, TestQueryPoolOutlivesStoreReaderAndExport) { const std::shared_ptr stored_schema = TransportSchema(); - std::shared_ptr pool = std::make_shared(); - std::weak_ptr pool_lifetime = pool; + std::shared_ptr pool = GetMemoryPool(); + std::weak_ptr pool_lifetime = pool; auto write_schema = std::make_unique(); ASSERT_TRUE(arrow::ExportSchema(*stored_schema, write_schema.get()).ok()); ArrowRealtimeStoreFactory factory; @@ -381,16 +339,13 @@ TEST(PrimaryKeyRealtimeStoreTest, TestQueryReaderProjectsNestedFields) { const std::shared_ptr stored_b = FieldWithId("b", arrow::int32(), 11); const std::shared_ptr stored_x = FieldWithId("x", arrow::int32(), 20); const std::shared_ptr stored_y = FieldWithId("y", arrow::int32(), 21); - arrow::FieldVector stored_fields = { - DataField::ConvertDataFieldToArrowField(SpecialFields::ValueKind())->WithNullable(false), - DataField::ConvertDataFieldToArrowField(SpecialFields::SequenceNumber()) - ->WithNullable(false), - DataField::ConvertDataFieldToArrowField(SpecialFields::RealtimeOffset()), + arrow::FieldVector stored_value_fields = { FieldWithId("id", arrow::int64(), 0), FieldWithId("profile", arrow::struct_({stored_profile_a}), 1), FieldWithId("items", arrow::list(arrow::struct_({stored_a, stored_b})), 2), FieldWithId("attrs", arrow::map(arrow::utf8(), arrow::struct_({stored_x, stored_y})), 3)}; - std::shared_ptr stored_schema = arrow::schema(std::move(stored_fields)); + std::shared_ptr stored_schema = + RealtimePrimaryKeyLayout::CreateSchema(stored_value_fields); ASSERT_OK_AND_ASSIGN(std::shared_ptr store, PrimaryKeyRealtimeStore::Create(stored_schema, GetDefaultPool())); ASSERT_OK(store->Write(RealtimeWriteBatch{ @@ -401,14 +356,14 @@ TEST(PrimaryKeyRealtimeStoreTest, TestQueryReaderProjectsNestedFields) { OffsetRange(0, 1)})); ASSERT_OK_AND_ASSIGN(std::shared_ptr view, store->AcquireReadView()); - arrow::FieldVector requested_fields(stored_schema->fields().begin(), - stored_schema->fields().begin() + 3); - requested_fields.push_back(FieldWithId("profile", arrow::struct_({stored_profile_a}), 1)); - requested_fields.push_back( + arrow::FieldVector requested_value_fields; + requested_value_fields.push_back(FieldWithId("profile", arrow::struct_({stored_profile_a}), 1)); + requested_value_fields.push_back( FieldWithId("items", arrow::list(arrow::struct_({stored_b, stored_a})), 2)); - requested_fields.push_back( + requested_value_fields.push_back( FieldWithId("attrs", arrow::map(arrow::utf8(), arrow::struct_({stored_y, stored_x})), 3)); - std::shared_ptr requested_schema = arrow::schema(std::move(requested_fields)); + std::shared_ptr requested_schema = + RealtimePrimaryKeyLayout::CreateSchema(requested_value_fields); auto c_schema = std::make_unique(); ASSERT_TRUE(arrow::ExportSchema(*requested_schema, c_schema.get()).ok()); RealtimeQueryContext context{c_schema.get(), /*predicate=*/nullptr, diff --git a/src/paimon/core/realtime/realtime_primary_key_reader.cpp b/src/paimon/core/realtime/realtime_primary_key_reader.cpp new file mode 100644 index 000000000..19ba0a985 --- /dev/null +++ b/src/paimon/core/realtime/realtime_primary_key_reader.cpp @@ -0,0 +1,510 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "paimon/core/realtime/realtime_primary_key_reader.h" + +#include +#include +#include +#include +#include +#include + +#include "arrow/array/array_base.h" +#include "arrow/array/array_primitive.h" +#include "arrow/c/bridge.h" +#include "arrow/type.h" +#include "fmt/format.h" +#include "paimon/common/data/columnar/columnar_batch_context.h" +#include "paimon/common/data/columnar/columnar_row_ref.h" +#include "paimon/common/table/special_fields.h" +#include "paimon/common/types/data_field.h" +#include "paimon/common/types/row_kind.h" +#include "paimon/common/utils/arrow/arrow_utils.h" +#include "paimon/common/utils/arrow/status_utils.h" +#include "paimon/common/utils/checked_cast.h" +#include "paimon/common/utils/scope_guard.h" +#include "paimon/core/utils/nested_projection_utils.h" +#include "paimon/macros.h" +#include "paimon/reader/batch_reader.h" +#include "paimon/status.h" +#include "paimon/utils/roaring_bitmap64.h" + +namespace paimon { + +namespace { + +template +void CloseReaders(const std::vector>& readers) { + for (const std::unique_ptr& reader : readers) { + if (reader) { + reader->Close(); + } + } +} + +class RealtimeOffsetCoverage { + public: + static Result> Create(const OffsetRange& offsets, + size_t reader_count, + bool allow_committed_prefix) { + if (offsets.begin < 0 || offsets.end < offsets.begin) { + return Status::Invalid("PK real-time store returned an invalid offset range"); + } + return std::shared_ptr( + new RealtimeOffsetCoverage(offsets, reader_count, allow_committed_prefix)); + } + + Status Add(const arrow::Int64Array& offsets) { + for (int64_t row = 0; row < offsets.length(); ++row) { + const int64_t offset = offsets.Value(row); + if (allow_committed_prefix_ && offset < 0) { + return Status::Invalid("PK real-time store reader offset must be non-negative"); + } + if (allow_committed_prefix_ && offset < offsets_.begin) { + continue; + } + if (offset < offsets_.begin || offset >= offsets_.end) { + return Status::Invalid( + allow_committed_prefix_ + ? "PK real-time store query reader offset is outside the visible range" + : "PK real-time store commit reader offset is outside the sealed range"); + } + if (!seen_offsets_.CheckedAdd(offset)) { + return CoverageError(); + } + } + return Status::OK(); + } + + Status FinishReader() { + ++finished_reader_count_; + if (finished_reader_count_ == reader_count_ && + seen_offsets_.Cardinality() != offsets_.Count()) { + return CoverageError(); + } + return Status::OK(); + } + + private: + RealtimeOffsetCoverage(const OffsetRange& offsets, size_t reader_count, + bool allow_committed_prefix) + : offsets_(offsets), + reader_count_(reader_count), + allow_committed_prefix_(allow_committed_prefix) {} + + Status CoverageError() const { + return Status::Invalid( + allow_committed_prefix_ + ? "PK real-time store query readers did not cover the visible range" + : "PK real-time store commit readers did not cover the sealed range"); + } + + OffsetRange offsets_; + size_t reader_count_; + bool allow_committed_prefix_; + RoaringBitmap64 seen_offsets_; + size_t finished_reader_count_ = 0; +}; + +Status CheckTransportField(const std::shared_ptr& schema, int32_t field_idx, + const DataField& expected_field) { + if (schema->num_fields() <= field_idx) { + return Status::Invalid( + fmt::format("realtime primary-key transport schema is missing field {} at index {}", + expected_field.Name(), field_idx)); + } + const std::shared_ptr& field = schema->field(field_idx); + PAIMON_ASSIGN_OR_RAISE(int32_t field_id, NestedProjectionUtils::GetPaimonFieldId(field)); + if (field->name() != expected_field.Name() || !field->type()->Equals(*expected_field.Type()) || + field->nullable() || field_id != expected_field.Id()) { + return Status::Invalid(fmt::format( + "realtime primary-key transport schema field {} must be non-null {}:{} with field id " + "{}, got {}:{} nullable={} field id {}", + field_idx, expected_field.Name(), expected_field.Type()->ToString(), + expected_field.Id(), field->name(), field->type()->ToString(), field->nullable(), + field_id)); + } + return Status::OK(); +} + +Result> ResolveFieldIndexes( + const std::shared_ptr& transport_schema, + const std::unordered_map& field_indexes, + const std::shared_ptr& row_schema) { + std::vector result; + result.reserve(row_schema->num_fields()); + for (const std::shared_ptr& row_field : row_schema->fields()) { + PAIMON_ASSIGN_OR_RAISE(int32_t field_id, + NestedProjectionUtils::GetPaimonFieldId(row_field)); + auto field_index = field_indexes.find(field_id); + if (field_index == field_indexes.end()) { + return Status::Invalid(fmt::format( + "cannot find field id {} in realtime primary-key transport schema", field_id)); + } + const std::shared_ptr& transport_field = + transport_schema->field(field_index->second); + if (!transport_field->type()->Equals(row_field->type())) { + return Status::Invalid(fmt::format( + "realtime primary-key transport field id {} type {} does not match row type {}", + field_id, transport_field->type()->ToString(), row_field->type()->ToString())); + } + result.push_back(field_index->second); + } + return result; +} + +class RealtimePrimaryKeyReaderPlan { + public: + static Result> Create( + const std::shared_ptr& transport_schema, + const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema) { + std::unordered_map field_indexes; + field_indexes.reserve(transport_schema->num_fields() - + RealtimePrimaryKeyLayout::kValueStartIndex); + for (int32_t i = RealtimePrimaryKeyLayout::kValueStartIndex; + i < transport_schema->num_fields(); ++i) { + PAIMON_ASSIGN_OR_RAISE(int32_t field_id, NestedProjectionUtils::GetPaimonFieldId( + transport_schema->field(i))); + if (!field_indexes.emplace(field_id, i).second) { + return Status::Invalid(fmt::format( + "duplicate field id {} in realtime primary-key transport schema", field_id)); + } + } + PAIMON_ASSIGN_OR_RAISE(std::vector key_field_indexes, + ResolveFieldIndexes(transport_schema, field_indexes, key_schema)); + PAIMON_ASSIGN_OR_RAISE(std::vector value_field_indexes, + ResolveFieldIndexes(transport_schema, field_indexes, value_schema)); + return std::shared_ptr(new RealtimePrimaryKeyReaderPlan( + transport_schema, std::move(key_field_indexes), std::move(value_field_indexes))); + } + + const std::shared_ptr& TransportSchema() const { + return transport_schema_; + } + + const std::vector& KeyFieldIndexes() const { + return key_field_indexes_; + } + + const std::vector& ValueFieldIndexes() const { + return value_field_indexes_; + } + + private: + RealtimePrimaryKeyReaderPlan(const std::shared_ptr& schema, + std::vector&& key_indexes, + std::vector&& value_indexes) + : transport_schema_(schema), + key_field_indexes_(std::move(key_indexes)), + value_field_indexes_(std::move(value_indexes)) {} + + const std::shared_ptr transport_schema_; + const std::vector key_field_indexes_; + const std::vector value_field_indexes_; +}; + +class RealtimePrimaryKeyReader final : public KeyValueRecordReader { + public: + RealtimePrimaryKeyReader(std::unique_ptr&& reader, + const std::shared_ptr& plan, + const std::optional& visible_offsets, + const std::shared_ptr& pool, + const std::shared_ptr& offset_coverage) + : reader_(std::move(reader)), + plan_(plan), + visible_offsets_(visible_offsets), + pool_(pool), + offset_coverage_(offset_coverage) {} + + class Iterator final : public KeyValueRecordReader::Iterator { + public: + explicit Iterator(RealtimePrimaryKeyReader* reader) : reader_(reader) {} + + Result HasNext() const override { + return cursor_ < reader_->RowCount(); + } + + Result Next() override { + if (cursor_ >= reader_->RowCount()) { + return Status::Invalid("No more realtime primary-key values in current iterator"); + } + const int64_t row = reader_->RowAt(cursor_); + std::shared_ptr key = + std::make_shared(reader_->key_ctx_, row); + auto value = std::make_unique(reader_->value_ctx_, row); + PAIMON_ASSIGN_OR_RAISE(const RowKind* row_kind, + RowKind::FromByteValue(reader_->row_kind_array_->Value(row))); + int64_t sequence_number = reader_->sequence_number_array_->Value(row); + ++cursor_; + return KeyValue(row_kind, sequence_number, KeyValue::UNKNOWN_LEVEL, std::move(key), + std::move(value)); + } + + private: + RealtimePrimaryKeyReader* reader_; + int64_t cursor_ = 0; + }; + + Result> NextBatch() override { + return NextBatchImpl(); + } + + std::shared_ptr GetReaderMetrics() const override { + return reader_->GetReaderMetrics(); + } + + void Close() override { + ResetBatchState(); + reader_->Close(); + } + + private: + Result> NextBatchImpl() { + while (true) { + ResetBatchState(); + PAIMON_ASSIGN_OR_RAISE(BatchReader::ReadBatchWithBitmap batch_with_bitmap, + reader_->NextBatchWithBitmap()); + if (BatchReader::IsEofBatch(batch_with_bitmap)) { + if (offset_coverage_ && !offset_coverage_finished_) { + offset_coverage_finished_ = true; + PAIMON_RETURN_NOT_OK(offset_coverage_->FinishReader()); + } + return std::unique_ptr(); + } + auto& [batch, selection] = batch_with_bitmap; + auto& [c_array, c_schema] = batch; + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr arrow_array, + arrow::ImportArray(c_array.get(), c_schema.get())); + if (!arrow_array || arrow_array->type_id() != arrow::Type::STRUCT) { + return Status::Invalid( + "cannot cast realtime primary-key transport batch to StructArray"); + } + std::shared_ptr data_batch = + checked_pointer_cast(arrow_array); + PAIMON_RETURN_NOT_OK(ValidateTransportBatch(data_batch)); + + std::shared_ptr> offset_array = + checked_pointer_cast>( + data_batch->field(RealtimePrimaryKeyLayout::kRealtimeOffsetIndex)); + if (offset_coverage_) { + PAIMON_RETURN_NOT_OK(offset_coverage_->Add(*offset_array)); + } + + row_kind_array_ = checked_pointer_cast>( + data_batch->field(RealtimePrimaryKeyLayout::kValueKindIndex)); + sequence_number_array_ = checked_pointer_cast>( + data_batch->field(RealtimePrimaryKeyLayout::kSequenceNumberIndex)); + arrow::ArrayVector key_fields; + key_fields.reserve(plan_->KeyFieldIndexes().size()); + for (int32_t index : plan_->KeyFieldIndexes()) { + key_fields.push_back(data_batch->field(index)); + } + arrow::ArrayVector value_fields; + value_fields.reserve(plan_->ValueFieldIndexes().size()); + for (int32_t index : plan_->ValueFieldIndexes()) { + value_fields.push_back(data_batch->field(index)); + } + key_ctx_ = std::make_shared(key_fields, pool_); + value_ctx_ = std::make_shared(value_fields, pool_); + PAIMON_ASSIGN_OR_RAISE(bool has_selected_rows, + SelectRows(*offset_array, std::move(selection))); + if (!has_selected_rows) { + continue; + } + ArrowUtils::TraverseArray(data_batch); + return std::make_unique(this); + } + } + + Status ValidateTransportBatch(const std::shared_ptr& data_batch) const { + if (data_batch->num_fields() != plan_->TransportSchema()->num_fields()) { + return Status::Invalid(fmt::format( + "realtime primary-key transport batch field count {} does not match schema field " + "count {}", + data_batch->num_fields(), plan_->TransportSchema()->num_fields())); + } + const arrow::FieldVector& batch_fields = data_batch->type()->fields(); + for (int32_t i = 0; i < data_batch->num_fields(); ++i) { + if (!batch_fields[i]->Equals(plan_->TransportSchema()->field(i), true)) { + return Status::Invalid(fmt::format( + "realtime primary-key transport batch field {} does not match declared schema", + i)); + } + } + if (data_batch->field(RealtimePrimaryKeyLayout::kValueKindIndex)->null_count() != 0 || + data_batch->field(RealtimePrimaryKeyLayout::kSequenceNumberIndex)->null_count() != 0 || + data_batch->field(RealtimePrimaryKeyLayout::kRealtimeOffsetIndex)->null_count() != 0) { + return Status::Invalid("realtime primary-key transport columns must not contain nulls"); + } + return Status::OK(); + } + + Result SelectRows(const arrow::Int64Array& offsets, RoaringBitmap32&& selection) { + for (auto iter = selection.Begin(); iter != selection.End(); ++iter) { + const uint32_t row = *iter; + if (static_cast(row) >= offsets.length()) { + return Status::Invalid( + fmt::format("selected row id {} is out of bounds for realtime primary-key " + "transport batch length {}", + row, offsets.length())); + } + } + if (selection.Cardinality() != offsets.length()) { + return Status::Invalid( + "PK real-time store reader bitmap must cover every raw " + "transport row"); + } + selected_rows_.reserve(offsets.length()); + for (int64_t row = 0; row < offsets.length(); ++row) { + if (!visible_offsets_.has_value() || (offsets.Value(row) >= visible_offsets_->begin && + offsets.Value(row) < visible_offsets_->end)) { + selected_rows_.push_back(row); + } + } + return !selected_rows_.empty(); + } + + int64_t RowCount() const { + return static_cast(selected_rows_.size()); + } + + int64_t RowAt(int64_t ordinal) const { + return selected_rows_[ordinal]; + } + + void ResetBatchState() { + key_ctx_.reset(); + value_ctx_.reset(); + row_kind_array_.reset(); + sequence_number_array_.reset(); + selected_rows_.clear(); + } + + private: + std::unique_ptr reader_; + std::shared_ptr plan_; + std::optional visible_offsets_; + std::shared_ptr pool_; + std::shared_ptr offset_coverage_; + bool offset_coverage_finished_ = false; + std::shared_ptr key_ctx_; + std::shared_ptr value_ctx_; + std::shared_ptr> row_kind_array_; + std::shared_ptr> sequence_number_array_; + std::vector selected_rows_; +}; + +} // namespace + +std::shared_ptr RealtimePrimaryKeyLayout::CreateSchema( + const std::vector>& value_fields) { + arrow::FieldVector fields = { + DataField::ConvertDataFieldToArrowField(SpecialFields::ValueKind())->WithNullable(false), + DataField::ConvertDataFieldToArrowField(SpecialFields::SequenceNumber()) + ->WithNullable(false), + DataField::ConvertDataFieldToArrowField(SpecialFields::RealtimeOffset())}; + fields.insert(fields.end(), value_fields.begin(), value_fields.end()); + return arrow::schema(std::move(fields)); +} + +Status RealtimePrimaryKeyLayout::ValidateSchema( + const std::shared_ptr& transport_schema) { + if (!transport_schema || transport_schema->num_fields() < kValueStartIndex) { + return Status::Invalid( + "realtime primary-key transport schema must contain transport fields"); + } + PAIMON_RETURN_NOT_OK( + CheckTransportField(transport_schema, kValueKindIndex, SpecialFields::ValueKind())); + PAIMON_RETURN_NOT_OK(CheckTransportField(transport_schema, kSequenceNumberIndex, + SpecialFields::SequenceNumber())); + PAIMON_RETURN_NOT_OK(CheckTransportField(transport_schema, kRealtimeOffsetIndex, + SpecialFields::RealtimeOffset())); + return Status::OK(); +} + +Result>> +RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::vector>&& readers, + const std::shared_ptr& transport_schema, const OffsetRange& visible_offsets, + const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema, + const std::shared_ptr& memory_pool) { + std::vector> adapted_readers; + ScopeGuard remaining_raw_readers_guard([&readers]() { CloseReaders(readers); }); + if (readers.empty() && visible_offsets.begin < visible_offsets.end) { + return Status::Invalid( + "PK real-time store returned no query readers for a non-empty visible range"); + } + for (const std::unique_ptr& reader : readers) { + if (!reader) { + return Status::Invalid("PK real-time store returned a null query reader"); + } + } + PAIMON_RETURN_NOT_OK(RealtimePrimaryKeyLayout::ValidateSchema(transport_schema)); + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr plan, + RealtimePrimaryKeyReaderPlan::Create(transport_schema, key_schema, value_schema)); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr offset_coverage, + RealtimeOffsetCoverage::Create(visible_offsets, readers.size(), + /*allow_committed_prefix=*/true)); + adapted_readers.reserve(readers.size()); + for (std::unique_ptr& reader : readers) { + adapted_readers.push_back(std::make_unique( + std::move(reader), plan, visible_offsets, memory_pool, offset_coverage)); + } + remaining_raw_readers_guard.Release(); + return adapted_readers; +} + +Result>> +RealtimePrimaryKeyReaderFactory::CreateForCommit( + std::vector>&& readers, + const std::shared_ptr& transport_schema, const OffsetRange& sealed_offsets, + const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema, + const std::shared_ptr& memory_pool) { + std::vector> adapted_readers; + ScopeGuard remaining_raw_readers_guard([&readers]() { CloseReaders(readers); }); + if (readers.empty()) { + return Status::Invalid( + "PK real-time store returned no commit readers for a sealed segment"); + } + for (const std::unique_ptr& reader : readers) { + if (!reader) { + return Status::Invalid("PK real-time store returned a null commit reader"); + } + } + PAIMON_RETURN_NOT_OK(RealtimePrimaryKeyLayout::ValidateSchema(transport_schema)); + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr plan, + RealtimePrimaryKeyReaderPlan::Create(transport_schema, key_schema, value_schema)); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr offset_coverage, + RealtimeOffsetCoverage::Create(sealed_offsets, readers.size(), + /*allow_committed_prefix=*/false)); + adapted_readers.reserve(readers.size()); + for (std::unique_ptr& reader : readers) { + adapted_readers.push_back(std::make_unique( + std::move(reader), plan, std::nullopt, memory_pool, offset_coverage)); + } + remaining_raw_readers_guard.Release(); + return adapted_readers; +} + +} // namespace paimon diff --git a/src/paimon/core/realtime/realtime_primary_key_reader.h b/src/paimon/core/realtime/realtime_primary_key_reader.h new file mode 100644 index 000000000..d175c3b63 --- /dev/null +++ b/src/paimon/core/realtime/realtime_primary_key_reader.h @@ -0,0 +1,71 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#pragma once + +#include +#include +#include + +#include "arrow/type_fwd.h" +#include "paimon/core/io/key_value_record_reader.h" +#include "paimon/realtime/offset_range.h" +#include "paimon/result.h" + +namespace paimon { +class BatchReader; +class MemoryPool; + +/// Defines the Arrow field layout for PK realtime transport batches. +class RealtimePrimaryKeyLayout { + public: + RealtimePrimaryKeyLayout() = delete; + ~RealtimePrimaryKeyLayout() = delete; + + static constexpr int32_t kValueKindIndex = 0; + static constexpr int32_t kSequenceNumberIndex = 1; + static constexpr int32_t kRealtimeOffsetIndex = 2; + static constexpr int32_t kValueStartIndex = 3; + + static std::shared_ptr CreateSchema( + const std::vector>& value_fields); + + static Status ValidateSchema(const std::shared_ptr& transport_schema); +}; + +class RealtimePrimaryKeyReaderFactory { + public: + RealtimePrimaryKeyReaderFactory() = delete; + ~RealtimePrimaryKeyReaderFactory() = delete; + + static Result>> CreateForQuery( + std::vector>&& readers, + const std::shared_ptr& transport_schema, const OffsetRange& visible_offsets, + const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema, + const std::shared_ptr& memory_pool); + + static Result>> CreateForCommit( + std::vector>&& readers, + const std::shared_ptr& transport_schema, const OffsetRange& sealed_offsets, + const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema, + const std::shared_ptr& memory_pool); +}; + +} // namespace paimon diff --git a/src/paimon/core/realtime/realtime_primary_key_reader_test.cpp b/src/paimon/core/realtime/realtime_primary_key_reader_test.cpp new file mode 100644 index 000000000..89a39cd63 --- /dev/null +++ b/src/paimon/core/realtime/realtime_primary_key_reader_test.cpp @@ -0,0 +1,688 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "paimon/core/realtime/realtime_primary_key_reader.h" + +#include +#include +#include +#include +#include + +#include "arrow/api.h" +#include "arrow/array/array_nested.h" +#include "arrow/ipc/json_simple.h" +#include "gtest/gtest.h" +#include "paimon/common/table/special_fields.h" +#include "paimon/common/types/data_field.h" +#include "paimon/memory/memory_pool.h" +#include "paimon/realtime/offset_range.h" +#include "paimon/testing/mock/mock_file_batch_reader.h" +#include "paimon/testing/utils/key_value_checker.h" +#include "paimon/testing/utils/read_result_collector.h" +#include "paimon/testing/utils/testharness.h" + +namespace paimon::test { + +namespace { + +std::shared_ptr MakeField(const std::string& name, + const std::shared_ptr& type, + int32_t field_id, bool nullable = true) { + return DataField::ConvertDataFieldToArrowField( + DataField(field_id, arrow::field(name, type, nullable))); +} + +std::shared_ptr MakeTransportSchema(const arrow::FieldVector& value_fields) { + return RealtimePrimaryKeyLayout::CreateSchema(value_fields); +} + +Result> CreateRealtimePrimaryKeyQueryReaderForTest( + std::unique_ptr&& reader, const std::shared_ptr& transport_schema, + const OffsetRange& visible_offsets, const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema, + const std::shared_ptr& memory_pool) { + std::vector> readers; + readers.push_back(std::move(reader)); + PAIMON_ASSIGN_OR_RAISE(std::vector> adapted_readers, + RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::move(readers), transport_schema, visible_offsets, key_schema, + value_schema, memory_pool)); + return std::move(adapted_readers[0]); +} + +Result> CreateRealtimePrimaryKeyCommitReaderForTest( + std::unique_ptr&& reader, const std::shared_ptr& transport_schema, + const OffsetRange& sealed_offsets, const std::shared_ptr& key_schema, + const std::shared_ptr& value_schema, + const std::shared_ptr& memory_pool) { + std::vector> readers; + readers.push_back(std::move(reader)); + PAIMON_ASSIGN_OR_RAISE(std::vector> adapted_readers, + RealtimePrimaryKeyReaderFactory::CreateForCommit( + std::move(readers), transport_schema, sealed_offsets, key_schema, + value_schema, memory_pool)); + return std::move(adapted_readers[0]); +} + +class TrackingBatchReader : public BatchReader { + public: + TrackingBatchReader(std::unique_ptr&& delegate, int32_t* close_count) + : delegate_(std::move(delegate)), close_count_(close_count) {} + + Result NextBatch() override { + return delegate_->NextBatch(); + } + + std::shared_ptr GetReaderMetrics() const override { + return delegate_->GetReaderMetrics(); + } + + void Close() override { + ++(*close_count_); + delegate_->Close(); + } + + private: + std::unique_ptr delegate_; + int32_t* close_count_; +}; + +class MalformedBitmapBatchReader : public BatchReader { + public: + MalformedBitmapBatchReader(std::unique_ptr&& delegate, int32_t row_id) + : delegate_(std::move(delegate)), row_id_(row_id) {} + + Result NextBatch() override { + return delegate_->NextBatch(); + } + + Result NextBatchWithBitmap() override { + PAIMON_ASSIGN_OR_RAISE(ReadBatchWithBitmap batch, delegate_->NextBatchWithBitmap()); + if (!IsEofBatch(batch)) { + batch.second.Add(row_id_); + } + return batch; + } + + std::shared_ptr GetReaderMetrics() const override { + return delegate_->GetReaderMetrics(); + } + + void Close() override { + delegate_->Close(); + } + + private: + std::unique_ptr delegate_; + int32_t row_id_; +}; + +} // namespace + +class RealtimePrimaryKeyReaderTest : public testing::Test { + protected: + std::shared_ptr pool_ = GetDefaultPool(); +}; + +TEST_F(RealtimePrimaryKeyReaderTest, TestTransportSchemaLayout) { + arrow::FieldVector value_fields = {arrow::field("key", arrow::int64(), false), + arrow::field("value", arrow::utf8())}; + std::shared_ptr schema = MakeTransportSchema(value_fields); + + ASSERT_EQ(RealtimePrimaryKeyLayout::kValueKindIndex, 0); + ASSERT_EQ(RealtimePrimaryKeyLayout::kSequenceNumberIndex, 1); + ASSERT_EQ(RealtimePrimaryKeyLayout::kRealtimeOffsetIndex, 2); + ASSERT_EQ(RealtimePrimaryKeyLayout::kValueStartIndex, 3); + ASSERT_EQ(schema->field(0)->name(), "_VALUE_KIND"); + ASSERT_EQ(schema->field(1)->name(), "_SEQUENCE_NUMBER"); + ASSERT_EQ(schema->field(2)->name(), "_REALTIME_OFFSET"); + ASSERT_EQ(schema->field(3)->name(), "key"); + ASSERT_EQ(schema->field(4)->name(), "value"); + ASSERT_FALSE(schema->field(0)->nullable()); + ASSERT_FALSE(schema->field(1)->nullable()); + ASSERT_EQ(schema->field(2)->nullable(), SpecialFields::RealtimeOffset().Nullable()); + ASSERT_FALSE(schema->field(3)->nullable()); + ASSERT_TRUE(schema->field(4)->nullable()); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestTransportSchemaValidation) { + const std::shared_ptr valid = MakeTransportSchema({}); + std::vector invalid_fields; + + arrow::FieldVector wrong_type = valid->fields(); + wrong_type[0] = DataField::ConvertDataFieldToArrowField( + DataField(SpecialFields::ValueKind().Id(), + arrow::field("_VALUE_KIND", arrow::int32(), false))) + ->WithNullable(false); + invalid_fields.push_back(std::move(wrong_type)); + + arrow::FieldVector nullable_sequence = valid->fields(); + nullable_sequence[1] = nullable_sequence[1]->WithNullable(true); + invalid_fields.push_back(std::move(nullable_sequence)); + + arrow::FieldVector wrong_offset_id = valid->fields(); + wrong_offset_id[2] = DataField::ConvertDataFieldToArrowField( + DataField(99, arrow::field("_REALTIME_OFFSET", arrow::int64(), false))) + ->WithNullable(false); + invalid_fields.push_back(std::move(wrong_offset_id)); + + for (const arrow::FieldVector& fields : invalid_fields) { + ASSERT_NOK_WITH_MSG(RealtimePrimaryKeyLayout::ValidateSchema(arrow::schema(fields)), + "transport schema field"); + } +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryAllowsCommittedPrefix) { + std::vector value_fields = {DataField(0, arrow::field("k0", arrow::int32())), + DataField(1, arrow::field("v0", arrow::int32()))}; + std::shared_ptr value_schema = + DataField::ConvertDataFieldsToArrowSchema(value_fields); + std::shared_ptr key_schema = arrow::schema({value_schema->field(0)}); + std::shared_ptr transport_schema = MakeTransportSchema(value_schema->fields()); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + auto transport_array = std::dynamic_pointer_cast( + arrow::ipc::internal::json::ArrayFromJSON(transport_type, R"([ + [0, 100, 0, 1, 10], + [0, 101, 1, 2, 20], + [0, 102, 2, 4, 40], + [0, 103, 3, 6, 60] + ])") + .ValueOrDie()); + + std::vector> batch_readers; + batch_readers.push_back( + std::make_unique(transport_array, transport_type, 2)); + ASSERT_OK_AND_ASSIGN(std::vector> readers, + RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::move(batch_readers), transport_schema, OffsetRange(2, 4), + key_schema, value_schema, pool_)); + ASSERT_EQ(1, readers.size()); + ASSERT_OK_AND_ASSIGN( + std::vector results, + (ReadResultCollector::CollectKeyValueResult< + KeyValueRecordReader, KeyValueRecordReader::Iterator>(readers[0].get()))); + + std::vector row_kinds = {const_cast(RowKind::Insert()), + const_cast(RowKind::Insert())}; + std::vector levels = {KeyValue::UNKNOWN_LEVEL, KeyValue::UNKNOWN_LEVEL}; + std::vector expected = KeyValueChecker::GenerateKeyValues( + row_kinds, {102, 103}, levels, {{4}, {6}}, {{4, 40}, {6, 60}}, pool_); + KeyValueChecker::CheckResult(expected, results, 1, 2); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryRejectsNegativeOffset) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, R"([[0, 10, -1, 1]])") + .ValueOrDie(); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::make_unique(transport_array, transport_type, + /*read_batch_size=*/1), + transport_schema, OffsetRange(1, 2), value_schema, value_schema, pool_)); + + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "reader offset must be non-negative"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryOffsetCoverageAcrossReadersAndBatches) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr first_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 10, 2, 1], [0, 11, 0, 2]])") + .ValueOrDie(); + std::shared_ptr second_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 12, 3, 3], [0, 13, 1, 4]])") + .ValueOrDie(); + std::vector> batch_readers; + batch_readers.push_back( + std::make_unique(first_array, transport_type, /*read_batch_size=*/1)); + batch_readers.push_back( + std::make_unique(second_array, transport_type, /*read_batch_size=*/1)); + + ASSERT_OK_AND_ASSIGN(std::vector> readers, + RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::move(batch_readers), transport_schema, OffsetRange(0, 4), + value_schema, value_schema, pool_)); + int64_t row_count = 0; + for (const std::unique_ptr& reader : readers) { + ASSERT_OK_AND_ASSIGN( + std::vector rows, + (ReadResultCollector::CollectKeyValueResult< + KeyValueRecordReader, KeyValueRecordReader::Iterator>(reader.get()))); + row_count += static_cast(rows.size()); + } + ASSERT_EQ(4, row_count); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryRejectsMissingVisibleOffset) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 10, 0, 1], [0, 11, 2, 2]])") + .ValueOrDie(); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::make_unique(transport_array, transport_type, + /*read_batch_size=*/1), + transport_schema, OffsetRange(0, 3), value_schema, value_schema, pool_)); + + ASSERT_NOK_WITH_MSG( + (ReadResultCollector::CollectKeyValueResult(reader.get())), + "query readers did not cover the visible range"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryRejectsDuplicateVisibleOffset) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr first_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, R"([[0, 10, 0, 1]])") + .ValueOrDie(); + std::shared_ptr second_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 11, 1, 2], [0, 12, 1, 3]])") + .ValueOrDie(); + std::vector> batch_readers; + batch_readers.push_back( + std::make_unique(first_array, transport_type, /*read_batch_size=*/1)); + batch_readers.push_back( + std::make_unique(second_array, transport_type, /*read_batch_size=*/1)); + ASSERT_OK_AND_ASSIGN(std::vector> readers, + RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::move(batch_readers), transport_schema, OffsetRange(0, 2), + value_schema, value_schema, pool_)); + ASSERT_OK_AND_ASSIGN( + std::vector first_rows, + (ReadResultCollector::CollectKeyValueResult< + KeyValueRecordReader, KeyValueRecordReader::Iterator>(readers[0].get()))); + ASSERT_EQ(1, first_rows.size()); + ASSERT_NOK_WITH_MSG((ReadResultCollector::CollectKeyValueResult( + readers[1].get())), + "query readers did not cover the visible range"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryRejectsEmptyEofForVisibleRange) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, R"([])").ValueOrDie(); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::make_unique(transport_array, transport_type, + /*read_batch_size=*/1), + transport_schema, OffsetRange(0, 1), value_schema, value_schema, pool_)); + + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "query readers did not cover the visible range"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryRejectsEmptyReadersForVisibleRange) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::vector> batch_readers; + + ASSERT_NOK_WITH_MSG( + RealtimePrimaryKeyReaderFactory::CreateForQuery(std::move(batch_readers), transport_schema, + OffsetRange(0, 1), value_schema, + value_schema, pool_), + "PK real-time store returned no query readers for a non-empty visible range"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryAllowsEmptyReadersForEmptyVisibleRange) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::vector> batch_readers; + + ASSERT_OK_AND_ASSIGN(std::vector> readers, + RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::move(batch_readers), transport_schema, OffsetRange(1, 1), + value_schema, value_schema, pool_)); + ASSERT_TRUE(readers.empty()); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryBitmapBounds) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, R"([[0, 10, 0, 1]])") + .ValueOrDie(); + auto batch_reader = std::make_unique( + std::make_unique(transport_array, transport_type, /*batch_size=*/1), + /*row_id=*/1); + + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::move(batch_reader), transport_schema, OffsetRange(0, 1), + value_schema, value_schema, pool_)); + Result> result = + ReadResultCollector::CollectKeyValueResult(reader.get()); + ASSERT_TRUE(result.status().IsInvalid()); + ASSERT_NOK_WITH_MSG(result, + "selected row id 1 is out of bounds for realtime primary-key transport " + "batch length 1"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryRejectsPartialBitmap) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 10, 0, 1], [0, 11, 1, 2]])") + .ValueOrDie(); + RoaringBitmap32 partial_bitmap; + partial_bitmap.Add(0); + auto batch_reader = std::make_unique( + transport_array, transport_type, partial_bitmap, /*read_batch_size=*/2); + batch_reader->EnableRandomizeBatchSize(false); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::move(batch_reader), transport_schema, OffsetRange(0, 2), + value_schema, value_schema, pool_)); + + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "must cover every raw transport row"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestCommitRejectsPartialBitmap) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 10, 0, 1], [0, 11, 1, 2]])") + .ValueOrDie(); + RoaringBitmap32 partial_bitmap; + partial_bitmap.Add(0); + auto batch_reader = std::make_unique( + transport_array, transport_type, partial_bitmap, /*read_batch_size=*/2); + batch_reader->EnableRandomizeBatchSize(false); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateRealtimePrimaryKeyCommitReaderForTest( + std::move(batch_reader), transport_schema, OffsetRange(0, 2), + value_schema, value_schema, pool_)); + + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "must cover every raw transport row"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestQueryProjection) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr extra = MakeField("extra", arrow::int32(), 1); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key, extra}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + auto transport_array = std::dynamic_pointer_cast( + arrow::ipc::internal::json::ArrayFromJSON(transport_type, R"([[0, 10, 0, 1, 2]])") + .ValueOrDie()); + + auto query_batch_reader = + std::make_unique(transport_array, transport_type, 1); + ASSERT_OK_AND_ASSIGN(std::unique_ptr query_reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::move(query_batch_reader), transport_schema, OffsetRange(0, 1), + value_schema, value_schema, pool_)); + ASSERT_OK_AND_ASSIGN( + std::vector query_results, + (ReadResultCollector::CollectKeyValueResult< + KeyValueRecordReader, KeyValueRecordReader::Iterator>(query_reader.get()))); + ASSERT_EQ(query_results.size(), 1); + ASSERT_EQ(query_results[0].value->GetFieldCount(), 1); + ASSERT_EQ(query_results[0].value->GetInt(0), 1); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestCommitOffsetCoverage) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr first_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 10, 2, 1], [0, 11, 0, 3]])") + .ValueOrDie(); + std::shared_ptr second_array = + arrow::ipc::internal::json::ArrayFromJSON(transport_type, + R"([[0, 12, 1, 2], [0, 13, 3, 4]])") + .ValueOrDie(); + std::vector> batch_readers; + batch_readers.push_back( + std::make_unique(first_array, transport_type, /*read_batch_size=*/1)); + batch_readers.push_back( + std::make_unique(second_array, transport_type, /*read_batch_size=*/1)); + + ASSERT_OK_AND_ASSIGN(std::vector> readers, + RealtimePrimaryKeyReaderFactory::CreateForCommit( + std::move(batch_readers), transport_schema, OffsetRange(0, 4), + value_schema, value_schema, pool_)); + int64_t row_count = 0; + for (const std::unique_ptr& reader : readers) { + ASSERT_OK_AND_ASSIGN( + std::vector rows, + (ReadResultCollector::CollectKeyValueResult< + KeyValueRecordReader, KeyValueRecordReader::Iterator>(reader.get()))); + row_count += static_cast(rows.size()); + } + ASSERT_EQ(4, row_count); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestCommitRejectsEmptyReaders) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::vector> batch_readers; + + ASSERT_NOK_WITH_MSG(RealtimePrimaryKeyReaderFactory::CreateForCommit( + std::move(batch_readers), transport_schema, OffsetRange(0, 1), + value_schema, value_schema, pool_), + "PK real-time store returned no commit readers for a sealed segment"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestRejectsDuplicateCommitOffset) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON( + transport_type, R"([[0, 10, 0, 1], [0, 11, 0, 2], [0, 12, 2, 3]])") + .ValueOrDie(); + std::vector> batch_readers; + batch_readers.push_back(std::make_unique(transport_array, transport_type, + /*read_batch_size=*/1)); + + ASSERT_OK_AND_ASSIGN(std::vector> readers, + RealtimePrimaryKeyReaderFactory::CreateForCommit( + std::move(batch_readers), transport_schema, OffsetRange(0, 3), + value_schema, value_schema, pool_)); + ASSERT_NOK_WITH_MSG((ReadResultCollector::CollectKeyValueResult( + readers[0].get())), + "did not cover the sealed range"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestBadCommitBatch) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value = MakeField("value", arrow::int32(), 1); + std::shared_ptr value_schema = arrow::schema({key, value}); + std::shared_ptr transport_schema = MakeTransportSchema({key, value}); + std::shared_ptr actual_schema = MakeTransportSchema({key}); + std::shared_ptr actual_type = arrow::struct_(actual_schema->fields()); + std::shared_ptr actual = + arrow::ipc::internal::json::ArrayFromJSON(actual_type, R"([[0, 10, 0, 1]])").ValueOrDie(); + + auto batch_reader = std::make_unique(actual, actual_type, 1); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateRealtimePrimaryKeyCommitReaderForTest( + std::move(batch_reader), transport_schema, OffsetRange(0, 1), + arrow::schema({key}), value_schema, pool_)); + ASSERT_NOK_WITH_MSG(reader->NextBatch(), "field count"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestSafeDecode) { + std::shared_ptr key = MakeField("key", arrow::int32(), 0); + std::shared_ptr value_schema = arrow::schema({key}); + std::shared_ptr transport_schema = MakeTransportSchema({key}); + + arrow::FieldVector invalid_fields = transport_schema->fields(); + invalid_fields[0] = invalid_fields[0]->WithName("wrong_value_kind"); + invalid_fields[3] = MakeField("wrong_key", arrow::int32(), 99); + std::shared_ptr invalid_type = arrow::struct_(invalid_fields); + auto invalid_array = std::dynamic_pointer_cast( + arrow::ipc::internal::json::ArrayFromJSON(invalid_type, R"([[0, 10, 0, 1]])").ValueOrDie()); + + auto batch_reader = std::make_unique(invalid_array, invalid_type, 1); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::move(batch_reader), transport_schema, OffsetRange(0, 1), + value_schema, value_schema, pool_)); + ASSERT_NOK_WITH_MSG( + (ReadResultCollector::CollectKeyValueResult(reader.get())), + "transport batch field"); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestNestedValues) { + std::shared_ptr id = MakeField("id", arrow::int32(), 0); + std::shared_ptr key_schema = arrow::schema({id}); + std::shared_ptr query_item_b = MakeField("renamed_b", arrow::int32(), 11); + std::shared_ptr query_item_a = MakeField("renamed_a", arrow::int32(), 10); + std::shared_ptr query_items = MakeField( + "items_renamed", + arrow::list(arrow::field("element", arrow::struct_({query_item_b, query_item_a}))), 2); + std::shared_ptr query_attr_y = MakeField("renamed_y", arrow::int32(), 21); + std::shared_ptr query_attr_x = MakeField("renamed_x", arrow::int32(), 20); + std::shared_ptr query_attrs = + MakeField("attrs_renamed", + arrow::map(arrow::utf8(), arrow::struct_({query_attr_y, query_attr_x})), 3); + std::shared_ptr query_key_right = MakeField("renamed_right", arrow::int32(), 31); + std::shared_ptr query_key_left = MakeField("renamed_left", arrow::int32(), 30); + std::shared_ptr query_keyed_values = + MakeField("keyed_values_renamed", + arrow::map(arrow::struct_({query_key_right, query_key_left}), arrow::int32()), 4); + std::shared_ptr query_value_schema = + arrow::schema({id, query_items, query_attrs, query_keyed_values}); + std::shared_ptr transport_schema = + MakeTransportSchema(query_value_schema->fields()); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + std::shared_ptr transport_array = + arrow::ipc::internal::json::ArrayFromJSON( + transport_type, + R"([[0, 10, 0, 1, [[200, 100], [400, 300]], [["k1", [8, 7]], ["k2", [10, 9]]], [[[12, 11], 13], [[22, 21], 23]]]])") + .ValueOrDie(); + + auto batch_reader = std::make_unique(transport_array, transport_type, 1); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateRealtimePrimaryKeyQueryReaderForTest( + std::move(batch_reader), transport_schema, OffsetRange(0, 1), + key_schema, query_value_schema, pool_)); + ASSERT_OK_AND_ASSIGN( + std::vector results, + (ReadResultCollector::CollectKeyValueResult(reader.get()))); + + ASSERT_EQ(results.size(), 1); + ASSERT_EQ(results[0].key->GetInt(0), 1); + ASSERT_EQ(results[0].value->GetFieldCount(), 4); + ASSERT_EQ(results[0].value->GetInt(0), 1); + + std::shared_ptr item_array = results[0].value->GetArray(1); + ASSERT_EQ(item_array->Size(), 2); + std::shared_ptr first_item = item_array->GetRow(0, 2); + ASSERT_EQ(first_item->GetInt(0), 200); + ASSERT_EQ(first_item->GetInt(1), 100); + std::shared_ptr second_item = item_array->GetRow(1, 2); + ASSERT_EQ(second_item->GetInt(0), 400); + ASSERT_EQ(second_item->GetInt(1), 300); + + std::shared_ptr attr_map = results[0].value->GetMap(2); + ASSERT_EQ(attr_map->Size(), 2); + std::shared_ptr key_array = attr_map->KeyArray(); + ASSERT_EQ(std::string(key_array->GetStringView(0)), "k1"); + ASSERT_EQ(std::string(key_array->GetStringView(1)), "k2"); + std::shared_ptr value_array = attr_map->ValueArray(); + std::shared_ptr first_attr = value_array->GetRow(0, 2); + ASSERT_EQ(first_attr->GetInt(0), 8); + ASSERT_EQ(first_attr->GetInt(1), 7); + std::shared_ptr second_attr = value_array->GetRow(1, 2); + ASSERT_EQ(second_attr->GetInt(0), 10); + ASSERT_EQ(second_attr->GetInt(1), 9); + + std::shared_ptr keyed_value_map = results[0].value->GetMap(3); + ASSERT_EQ(keyed_value_map->Size(), 2); + std::shared_ptr struct_keys = keyed_value_map->KeyArray(); + std::shared_ptr first_key = struct_keys->GetRow(0, 2); + ASSERT_EQ(first_key->GetInt(0), 12); + ASSERT_EQ(first_key->GetInt(1), 11); + std::shared_ptr second_key = struct_keys->GetRow(1, 2); + ASSERT_EQ(second_key->GetInt(0), 22); + ASSERT_EQ(second_key->GetInt(1), 21); + ASSERT_EQ(keyed_value_map->ValueArray()->GetInt(0), 13); + ASSERT_EQ(keyed_value_map->ValueArray()->GetInt(1), 23); +} + +TEST_F(RealtimePrimaryKeyReaderTest, TestFactoryFailureClosesReaders) { + std::vector value_fields = {DataField(0, arrow::field("k0", arrow::int32())), + DataField(1, arrow::field("v0", arrow::int32()))}; + std::shared_ptr value_schema = + DataField::ConvertDataFieldsToArrowSchema(value_fields); + std::shared_ptr key_schema = arrow::schema({value_schema->field(0)}); + std::shared_ptr transport_schema = MakeTransportSchema(value_schema->fields()); + std::shared_ptr transport_type = arrow::struct_(transport_schema->fields()); + auto transport_array = std::dynamic_pointer_cast( + arrow::ipc::internal::json::ArrayFromJSON(transport_type, R"([ + [0, 10, 0, 1, 100] + ])") + .ValueOrDie()); + + int32_t factory_failure_close_count = 0; + std::vector> batch_readers; + batch_readers.push_back(std::make_unique( + std::make_unique(transport_array, transport_type, 1), + &factory_failure_close_count)); + batch_readers.push_back(nullptr); + ASSERT_NOK_WITH_MSG(RealtimePrimaryKeyReaderFactory::CreateForQuery( + std::move(batch_readers), transport_schema, OffsetRange(0, 1), + key_schema, value_schema, pool_), + "PK real-time store returned a null query reader"); + ASSERT_EQ(factory_failure_close_count, 1); +} + +} // namespace paimon::test