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
2 changes: 2 additions & 0 deletions src/paimon/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions src/paimon/core/io/merged_key_value_record_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@

#include <memory>
#include <utility>
#include <vector>

#include "arrow/api.h"
#include "arrow/array/array_nested.h"
Expand Down
85 changes: 20 additions & 65 deletions src/paimon/core/realtime/primary_key_realtime_store_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@

#include "paimon/core/realtime/primary_key_realtime_store.h"

#include <cstddef>
#include <cstdint>
#include <memory>
#include <optional>
Expand All @@ -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"
Expand All @@ -50,28 +49,17 @@ std::shared_ptr<arrow::Field> FieldWithId(const std::string& name,
}

std::shared_ptr<arrow::Schema> 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<arrow::Schema> 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<RecordBatch> MakeBatch(const std::string& json) {
Expand Down Expand Up @@ -125,36 +113,6 @@ Result<std::string> ReadJson(const std::vector<std::unique_ptr<BatchReader>>& 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<MemoryPool> delegate_ = GetMemoryPool();
};

TEST(PrimaryKeyRealtimeStoreTest, TestWriteAndSealValidation) {
ASSERT_OK_AND_ASSIGN(std::shared_ptr<PrimaryKeyRealtimeStore> store,
PrimaryKeyRealtimeStore::Create(TransportSchema(), GetDefaultPool()));
Expand Down Expand Up @@ -337,8 +295,8 @@ TEST(PrimaryKeyRealtimeStoreTest, TestQueryReaderPerStoredBatch) {

TEST(PrimaryKeyRealtimeStoreTest, TestQueryPoolOutlivesStoreReaderAndExport) {
const std::shared_ptr<arrow::Schema> stored_schema = TransportSchema();
std::shared_ptr<TestingMemoryPool> pool = std::make_shared<TestingMemoryPool>();
std::weak_ptr<TestingMemoryPool> pool_lifetime = pool;
std::shared_ptr<MemoryPool> pool = GetMemoryPool();
std::weak_ptr<MemoryPool> pool_lifetime = pool;
auto write_schema = std::make_unique<ArrowSchema>();
ASSERT_TRUE(arrow::ExportSchema(*stored_schema, write_schema.get()).ok());
ArrowRealtimeStoreFactory factory;
Expand Down Expand Up @@ -381,16 +339,13 @@ TEST(PrimaryKeyRealtimeStoreTest, TestQueryReaderProjectsNestedFields) {
const std::shared_ptr<arrow::Field> stored_b = FieldWithId("b", arrow::int32(), 11);
const std::shared_ptr<arrow::Field> stored_x = FieldWithId("x", arrow::int32(), 20);
const std::shared_ptr<arrow::Field> 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<arrow::Schema> stored_schema = arrow::schema(std::move(stored_fields));
std::shared_ptr<arrow::Schema> stored_schema =
RealtimePrimaryKeyLayout::CreateSchema(stored_value_fields);
ASSERT_OK_AND_ASSIGN(std::shared_ptr<PrimaryKeyRealtimeStore> store,
PrimaryKeyRealtimeStore::Create(stored_schema, GetDefaultPool()));
ASSERT_OK(store->Write(RealtimeWriteBatch{
Expand All @@ -401,14 +356,14 @@ TEST(PrimaryKeyRealtimeStoreTest, TestQueryReaderProjectsNestedFields) {
OffsetRange(0, 1)}));
ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeReadView> 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<arrow::Schema> requested_schema = arrow::schema(std::move(requested_fields));
std::shared_ptr<arrow::Schema> requested_schema =
RealtimePrimaryKeyLayout::CreateSchema(requested_value_fields);
auto c_schema = std::make_unique<ArrowSchema>();
ASSERT_TRUE(arrow::ExportSchema(*requested_schema, c_schema.get()).ok());
RealtimeQueryContext context{c_schema.get(), /*predicate=*/nullptr,
Expand Down
Loading
Loading