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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -583,7 +583,7 @@ if(PAIMON_ENABLE_LANCE)
add_subdirectory(src/paimon/format/lance)
endif()
if(PAIMON_ENABLE_LUMINA)
add_subdirectory(src/paimon/global_index/lumina)
add_subdirectory(src/paimon/indexer/lumina)
endif()
add_subdirectory(src/paimon/global_index/lucene)
if(PAIMON_ENABLE_TANTIVY)
Expand Down
1,024 changes: 0 additions & 1,024 deletions src/paimon/global_index/lumina/lumina_global_index.cpp

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,14 @@
# limitations under the License.

if(PAIMON_ENABLE_LUMINA)
set(PAIMON_LUMINA_INDEX lumina_global_index.cpp lumina_global_index_factory.cpp)
set(PAIMON_LUMINA_INDEX
lumina_dataset.cpp
lumina_global_index.cpp
lumina_global_index_factory.cpp
lumina_index_accumulator.cpp
lumina_index_options.cpp
lumina_index_searcher.cpp
lumina_tag_utils.cpp)

add_paimon_lib(paimon_lumina_index
SOURCES
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,9 @@
#include "lumina/core/Types.h"
#include "lumina/extensions/SearchWithFilterExtension.h"
#include "paimon/fs/local/local_file_system.h"
#include "paimon/global_index/lumina/lumina_file_reader.h"
#include "paimon/global_index/lumina/lumina_file_writer.h"
#include "paimon/global_index/lumina/lumina_memory_pool.h"
#include "paimon/indexer/lumina/lumina_file_reader.h"
#include "paimon/indexer/lumina/lumina_file_writer.h"
#include "paimon/indexer/lumina/lumina_memory_pool.h"
#include "paimon/testing/utils/testharness.h"
namespace paimon::lumina::test {
class LuminaInterfaceTest : public ::testing::Test {
Expand Down
105 changes: 105 additions & 0 deletions src/paimon/indexer/lumina/lumina_dataset.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
/*
* 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/indexer/lumina/lumina_dataset.h"

#include <cstring>
#include <numeric>
#include <utility>

namespace paimon::lumina {

LuminaDataset::LuminaDataset(int64_t element_count, uint32_t dimension,
const std::vector<std::shared_ptr<arrow::FloatArray>>& arrays,
const std::vector<int64_t>& start_ids)
: element_count_(element_count),
dimension_(dimension),
arrays_(arrays),
start_ids_(start_ids) {}

uint32_t LuminaDataset::Dim() const noexcept {
return dimension_;
}

uint64_t LuminaDataset::TotalSize() const noexcept {
return static_cast<uint64_t>(element_count_);
}

::lumina::core::Result<uint64_t> LuminaDataset::GetNextBatch(
std::vector<float>& vector_buffer,
std::vector<::lumina::core::vector_id_t>& id_buffer) noexcept {
if (cursor_ >= arrays_.size()) {
return ::lumina::core::Result<uint64_t>::Ok(0);
}
std::shared_ptr<arrow::FloatArray>& values = arrays_[cursor_];
int64_t value_count = values->length();
int64_t vector_count = value_count / dimension_;
vector_buffer.resize(static_cast<size_t>(value_count));
std::memcpy(vector_buffer.data(), values->raw_values(),
sizeof(float) * static_cast<size_t>(value_count));
id_buffer.resize(static_cast<size_t>(vector_count));
std::iota(id_buffer.begin(), id_buffer.end(),
static_cast<::lumina::core::vector_id_t>(start_ids_[cursor_]));

// release the array when copy to vector_buffer
values.reset();
++cursor_;
return ::lumina::core::Result<uint64_t>::Ok(static_cast<uint64_t>(vector_count));
}

LuminaDatasetWithTag::LuminaDatasetWithTag(
int64_t element_count, uint32_t dimension,
const std::vector<std::shared_ptr<arrow::FloatArray>>& arrays,
const std::vector<int64_t>& start_ids,
const std::vector<std::vector<TagDimensionData>>& tag_data)
: element_count_(element_count),
dimension_(dimension),
arrays_(arrays),
start_ids_(start_ids),
tag_data_(tag_data) {}

uint32_t LuminaDatasetWithTag::Dim() const noexcept {
return dimension_;
}

uint64_t LuminaDatasetWithTag::TotalSize() const noexcept {
return static_cast<uint64_t>(element_count_);
}

::lumina::core::Result<uint64_t> LuminaDatasetWithTag::GetNextBatch(
std::vector<float>& vector_buffer, std::vector<::lumina::core::vector_id_t>& id_buffer,
std::vector<TagDimensionData>& tag_dimensions_data) noexcept {
if (cursor_ >= arrays_.size()) {
return ::lumina::core::Result<uint64_t>::Ok(0);
}
std::shared_ptr<arrow::FloatArray>& values = arrays_[cursor_];
int64_t value_count = values->length();
int64_t vector_count = value_count / dimension_;
vector_buffer.resize(static_cast<size_t>(value_count));
std::memcpy(vector_buffer.data(), values->raw_values(),
sizeof(float) * static_cast<size_t>(value_count));
id_buffer.resize(static_cast<size_t>(vector_count));
std::iota(id_buffer.begin(), id_buffer.end(),
static_cast<::lumina::core::vector_id_t>(start_ids_[cursor_]));
tag_dimensions_data = std::move(tag_data_[cursor_]);
values.reset();
++cursor_;
return ::lumina::core::Result<uint64_t>::Ok(static_cast<uint64_t>(vector_count));
}

} // namespace paimon::lumina
80 changes: 80 additions & 0 deletions src/paimon/indexer/lumina/lumina_dataset.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
/*
* 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 <cstddef>
#include <cstdint>
#include <memory>
#include <vector>

#include "arrow/array.h"
#include "lumina/api/Dataset.h"
#include "lumina/extensions/experimental/DatasetWithTag.h"

namespace paimon::lumina {

class LuminaDataset final : public ::lumina::api::Dataset {
public:
LuminaDataset(int64_t element_count, uint32_t dimension,
const std::vector<std::shared_ptr<arrow::FloatArray>>& arrays,
const std::vector<int64_t>& start_ids);

uint32_t Dim() const noexcept override;

uint64_t TotalSize() const noexcept override;

::lumina::core::Result<uint64_t> GetNextBatch(
std::vector<float>& vector_buffer,
std::vector<::lumina::core::vector_id_t>& id_buffer) noexcept override;

private:
int64_t element_count_;
uint32_t dimension_;
std::vector<std::shared_ptr<arrow::FloatArray>> arrays_;
std::vector<int64_t> start_ids_;
size_t cursor_ = 0;
};

class LuminaDatasetWithTag final : public ::lumina::extensions::experimental::DatasetWithTag {
public:
using TagDimensionData = ::lumina::extensions::experimental::TagDimensionData;

LuminaDatasetWithTag(int64_t element_count, uint32_t dimension,
const std::vector<std::shared_ptr<arrow::FloatArray>>& arrays,
const std::vector<int64_t>& start_ids,
const std::vector<std::vector<TagDimensionData>>& tag_data);

uint32_t Dim() const noexcept override;

uint64_t TotalSize() const noexcept override;

::lumina::core::Result<uint64_t> GetNextBatch(
std::vector<float>& vector_buffer, std::vector<::lumina::core::vector_id_t>& id_buffer,
std::vector<TagDimensionData>& tag_dimensions_data) noexcept override;

private:
int64_t element_count_;
uint32_t dimension_;
std::vector<std::shared_ptr<arrow::FloatArray>> arrays_;
std::vector<int64_t> start_ids_;
std::vector<std::vector<TagDimensionData>> tag_data_;
size_t cursor_ = 0;
};

} // namespace paimon::lumina
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,8 @@
*/
#include <future>

#include "paimon/global_index/lumina/lumina_file_reader.h"
#include "paimon/global_index/lumina/lumina_file_writer.h"
#include "paimon/indexer/lumina/lumina_file_reader.h"
#include "paimon/indexer/lumina/lumina_file_writer.h"
#include "paimon/testing/utils/testharness.h"
namespace paimon::lumina::test {
class LuminaFileIOTest : public ::testing::Test {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
#include "lumina/io/FileReader.h"
#include "paimon/common/utils/math.h"
#include "paimon/fs/file_system.h"
#include "paimon/global_index/lumina/lumina_utils.h"
#include "paimon/indexer/lumina/lumina_utils.h"
namespace paimon::lumina {
class LuminaFileReader : public ::lumina::io::FileReader {
public:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
#include "lumina/io/FileWriter.h"
#include "paimon/common/utils/math.h"
#include "paimon/fs/file_system.h"
#include "paimon/global_index/lumina/lumina_utils.h"
#include "paimon/indexer/lumina/lumina_utils.h"
namespace paimon::lumina {
class LuminaFileWriter : public ::lumina::io::FileWriter {
public:
Expand Down
Loading
Loading