lxy-9602 commented on code in PR #372: URL: https://github.com/apache/paimon-cpp/pull/372#discussion_r4228393334
########## src/paimon/common/file_index/scored_file_index_result_test.cpp: ########## @@ -0,0 +1,55 @@ +/* + * 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/file_index/scored_file_index_result.h" + +#include <vector> + +#include "gtest/gtest.h" +#include "paimon/testing/utils/testharness.h" + +namespace paimon::test { + +TEST(ScoredFileIndexResultTest, TestCreate) { + ASSERT_OK_AND_ASSIGN( + std::shared_ptr<ScoredFileIndexResult> result, + ScoredFileIndexResult::Create(RoaringBitmap32::From({2, 5}), {0.25f, 0.75f})); + EXPECT_FALSE(result->IsEmpty()); + EXPECT_EQ(RoaringBitmap32::From({2, 5}), result->GetRowPositions()); + EXPECT_EQ(std::vector<float>({0.25f, 0.75f}), result->GetScores()); + ASSERT_OK_AND_ASSIGN(bool remain, result->IsRemain()); + EXPECT_TRUE(remain); + EXPECT_EQ("row positions: {2,5}, scores: {0.25,0.75}", result->ToString()); + + std::shared_ptr<FileIndexResult> file_index_result = result; + ASSERT_OK_AND_ASSIGN(remain, file_index_result->IsRemain()); Review Comment: why need `file_index_result`? ########## src/paimon/common/reader/complete_index_score_batch_reader_test.cpp: ########## @@ -256,4 +257,78 @@ TEST_F(CompleteIndexScoreBatchReaderTest, TestGlobalScoresRequireReturnedRowId) "requires an int64 _ROW_ID field"); } +TEST_F(CompleteIndexScoreBatchReaderTest, TestFileReaderForwardsOperationsAndResetsScores) { + arrow::FieldVector fields = {arrow::field("f0", arrow::utf8()), + arrow::field("_INDEX_SCORE", arrow::float32())}; + auto data = arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), R"([ + ["Alice", null], + ["Bob", null] + ])") + .ValueOrDie(); + auto inner_reader = std::make_unique<MockFileBatchReader>(data, data->type(), /*batch_size=*/1); + MockFileBatchReader* inner = inner_reader.get(); + auto reader = std::make_unique<CompleteIndexScoreFileBatchReader>( + std::move(inner_reader), std::vector<float>{1.25f, 2.5f}, GetArrowPool(GetDefaultPool())); + + ASSERT_OK_AND_ASSIGN(uint64_t row_count, reader->GetNumberOfRows()); + EXPECT_EQ(2, row_count); + EXPECT_FALSE(reader->SupportPreciseBitmapSelection()); + reader->Warmup(); + EXPECT_EQ(1, inner->GetWarmupCount()); + + ASSERT_OK_AND_ASSIGN(BatchReader::ReadBatchWithBitmap first, reader->NextBatchWithBitmap()); + ASSERT_OK_AND_ASSIGN(uint64_t file_row_id, reader->GetPreviousBatchFileRowId(0)); + EXPECT_EQ(0, file_row_id); + auto first_array = + arrow::ImportArray(first.first.first.get(), first.first.second.get()).ValueOrDie(); Review Comment: simplify `first.first.first`. ########## src/paimon/core/io/file_index_evaluator_test.cpp: ########## @@ -416,8 +421,9 @@ TEST_F(FileIndexEvaluatorTest, TestInvalidEvaluate) { auto predicate = PredicateBuilder::IsNull(/*field_index=*/2, /*field_name=*/"f2", FieldType::INT); ASSERT_NOK_WITH_MSG( - FileIndexEvaluator::Evaluate(data_schema_, predicate, /*data_file_path_factory=*/nullptr, - data_file_meta, /*file_system=*/nullptr, pool_), + FileIndexEvaluator::Evaluate(data_schema_, core_options_, predicate, + /*data_file_path_factory=*/nullptr, data_file_meta, + /*file_system=*/nullptr, pool_), "read process for FileIndexEvaluator must have data_file_path_factory and file_system"); } Review Comment: add case for `EvaluateVectorSearch` and `EvaluateFTS` ########## src/paimon/indexer/lumina/lumina_file_index_factory.cpp: ########## @@ -0,0 +1,46 @@ +/* + * 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 <map> +#include <memory> +#include <string> +#include <utility> + +#include "paimon/factories/factory.h" +#include "paimon/file_index/file_indexer_factory.h" +#include "paimon/indexer/lumina/lumina_file_index.h" + +namespace paimon::lumina { +namespace { + +class LuminaFileIndexFactory final : public FileIndexerFactory { + public: + const char* Identifier() const override { + return "lumina"; + } + + Result<std::unique_ptr<FileIndexer>> Create( + const std::map<std::string, std::string>& options) const override { + return std::make_unique<LuminaFileIndexer>(options); + } +}; + +REGISTER_PAIMON_FACTORY(LuminaFileIndexFactory); + Review Comment: add a `.h`? ########## src/paimon/CMakeLists.txt: ########## @@ -138,6 +138,8 @@ set(PAIMON_COMMON_SRCS common/predicate/predicate_utils.cpp common/predicate/starts_with.cpp Review Comment: add an inte write and read case. ########## src/paimon/format/blob/blob_format_writer.h: ########## @@ -134,10 +134,12 @@ class BlobFormatWriter : public FormatWriter { int64_t offset, int32_t count); /// The input stream of a blob value, as Java's BlobCopySource. A `reused` stream is a view on - /// the kept source stream and must not be closed; any other stream is the value's own. + /// the kept source stream and must not be closed. For a dynamic-length descriptor, the view + /// shares an independently opened file stream, which must be closed instead of the view. struct BlobCopySource { std::unique_ptr<InputStream> stream; bool reused = false; + std::shared_ptr<InputStream> owned_stream; Review Comment: please double check ########## src/paimon/indexer/lumina/lumina_file_index_test.cpp: ########## @@ -0,0 +1,241 @@ +/* + * 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_file_index.h" + +#include <map> +#include <memory> +#include <optional> +#include <string> +#include <vector> + +#include "arrow/api.h" +#include "arrow/c/bridge.h" +#include "gtest/gtest.h" +#include "paimon/io/byte_array_input_stream.h" +#include "paimon/memory/memory_pool.h" +#include "paimon/predicate/predicate_builder.h" +#include "paimon/predicate/vector_search.h" +#include "paimon/testing/utils/testharness.h" +#include "paimon/utils/roaring_bitmap32.h" + +namespace paimon::lumina::test { +namespace { + +std::shared_ptr<arrow::StructArray> CreateVectors() { + std::shared_ptr<arrow::FloatBuilder> values = + std::make_shared<arrow::FloatBuilder>(arrow::default_memory_pool()); + arrow::ListBuilder vectors(arrow::default_memory_pool(), values); + EXPECT_TRUE(vectors.Append().ok()); + EXPECT_TRUE(values->AppendValues({0.0f, 0.0f, 0.0f, 0.0f}).ok()); + EXPECT_TRUE(vectors.AppendNull().ok()); + EXPECT_TRUE(vectors.Append().ok()); + EXPECT_TRUE(values->AppendValues({1.0f, 1.0f, 1.0f, 1.0f}).ok()); + std::shared_ptr<arrow::Array> vector_array; + EXPECT_TRUE(vectors.Finish(&vector_array).ok()); + return arrow::StructArray::Make({vector_array}, + {arrow::field("embedding", vector_array->type())}) + .ValueOrDie(); Review Comment: Please use json ########## src/paimon/core/operation/raw_file_split_read.h: ########## Review Comment: add `CompleteIndexScoreFileBatchReader` ########## src/paimon/core/io/file_index_options.cpp: ########## @@ -36,13 +37,28 @@ constexpr char kFileIndexPrefix[] = "file-index."; constexpr char kColumnsSuffix[] = ".columns"; constexpr size_t kFileIndexPrefixLength = sizeof(kFileIndexPrefix) - 1; constexpr size_t kColumnsSuffixLength = sizeof(kColumnsSuffix) - 1; +constexpr int64_t kDefaultInManifestThreshold = 500; } // namespace Result<FileIndexOptions> FileIndexOptions::FromCoreOptions(const CoreOptions& options) { + return Parse(options.ToMap(), options.FileIndexInManifestThreshold()); +} + +Result<FileIndexOptions> FileIndexOptions::FromMap( + const std::map<std::string, std::string>& raw_options) { + int64_t in_manifest_threshold = kDefaultInManifestThreshold; + auto iter = raw_options.find(Options::FILE_INDEX_IN_MANIFEST_THRESHOLD); + if (iter != raw_options.end()) { + PAIMON_ASSIGN_OR_RAISE(in_manifest_threshold, MemorySize::ParseBytes(iter->second)); + } + return Parse(raw_options, in_manifest_threshold); +} Review Comment: Why re-define `kDefaultInManifestThreshold`? ########## src/paimon/indexer/lumina/lumina_file_index.cpp: ########## @@ -0,0 +1,203 @@ +/* + * 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_file_index.h" + +#include <cstring> +#include <map> +#include <utility> + +#include "arrow/c/bridge.h" +#include "fmt/format.h" +#include "lumina/api/LuminaBuilder.h" +#include "lumina/core/Types.h" +#include "paimon/common/io/byte_array_output_stream.h" +#include "paimon/common/io/memory_segment_output_stream.h" +#include "paimon/common/io/offset_input_stream.h" +#include "paimon/common/utils/arrow/status_utils.h" +#include "paimon/file_index/scored_file_index_result.h" +#include "paimon/indexer/lumina/lumina_file_writer.h" +#include "paimon/indexer/lumina/lumina_utils.h" +#include "paimon/memory/bytes.h" +#include "paimon/status.h" + +namespace paimon::lumina { +namespace { + +Result<std::shared_ptr<arrow::Schema>> ImportVectorSchema(::ArrowSchema* c_schema, + const std::string& owner) { + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Schema> schema, + arrow::ImportSchema(c_schema)); + if (schema->num_fields() == 0) { + return Status::Invalid(fmt::format("{} requires at least one field", owner)); + } + std::shared_ptr<arrow::ListType> list_type = + std::dynamic_pointer_cast<arrow::ListType>(schema->field(0)->type()); + if (!list_type || list_type->value_type()->id() != arrow::Type::FLOAT) { + return Status::Invalid(fmt::format("{} field type must be list[float]", owner)); + } + return schema; +} + +} // namespace + +LuminaFileIndexWriter::LuminaFileIndexWriter(std::string field_name, + std::shared_ptr<arrow::DataType> arrow_type, + const LuminaIndexInfo& index_info, + ::lumina::api::BuilderOptions&& builder_options, + std::vector<LuminaTagField>&& tag_fields, + std::shared_ptr<LuminaMemoryPool> pool) + : field_name_(std::move(field_name)), + arrow_type_(std::move(arrow_type)), + index_info_(index_info), + builder_options_(std::move(builder_options)), + tag_fields_(std::move(tag_fields)), + pool_(std::move(pool)) {} + +Status LuminaFileIndexWriter::AddBatch(::ArrowArray* batch) { + if (serialized_) { + return Status::Invalid("Cannot add data after serializing a Lumina File Index"); + } + if (!batch || !batch->release) { + return Status::Invalid("Lumina File Index batch cannot be null or released"); + } + if (batch->length < 0 || + row_count_ > static_cast<int64_t>(RoaringBitmap32::MAX_VALUE) - batch->length) { + return Status::Invalid("Lumina File Index row count exceeds the bitmap32 limit"); + } + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Array> array, + arrow::ImportArray(batch, arrow_type_)); + if (array->null_count() != 0) { + return Status::Invalid("Lumina File Index struct array must not contain null rows"); + } + std::shared_ptr<arrow::StructArray> struct_array = + std::dynamic_pointer_cast<arrow::StructArray>(array); + if (!struct_array) { + return Status::Invalid("Lumina File Index input must be a struct array"); + } + std::shared_ptr<arrow::ListArray> vectors = + std::dynamic_pointer_cast<arrow::ListArray>(struct_array->GetFieldByName(field_name_)); + if (!vectors) { + return Status::Invalid("Lumina File Index field must be a list array"); + } Review Comment: `check_pointer_cast`? -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
