lxy-9602 commented on code in PR #382:
URL: https://github.com/apache/paimon-cpp/pull/382#discussion_r4079686764
##########
include/paimon/global_index/global_index_write_task.h:
##########
@@ -46,6 +47,11 @@ class PAIMON_EXPORT GlobalIndexWriteTask {
/// by the given `indexed_split`.
/// @param options Index-specific configuration (e.g., false positive
rate for bloom
/// filters).
+ /// @param task_id When checkpoints are enabled, the caller must
provide a non-empty task
+ /// identifier that uniquely identifies an index build task. Reuse it when
retrying the same
+ /// build. If the source data, build configuration, or build source code
changes, the caller
+ /// must use a new identifier; otherwise, the index build may fail. Pass
nullopt when
+ /// checkpoints are disabled. Index types without checkpoint support
ignore this value.
/// @param pool Memory pool for temporary allocations during index
construction.
Review Comment:
The comment is misaligned; please align it with the parameter comments above
and below.
##########
src/paimon/core/global_index/global_index_file_manager.h:
##########
@@ -19,22 +19,37 @@
#pragma once
+#include <algorithm>
+#include <cstdint>
#include <memory>
+#include <optional>
#include <string>
+#include <utility>
+#include <vector>
+#include "paimon/common/utils/path_util.h"
#include "paimon/common/utils/uuid.h"
+#include "paimon/core/index/index_checkpoint_path_factory.h"
#include "paimon/core/index/index_path_factory.h"
#include "paimon/fs/file_system.h"
+#include "paimon/global_index/io/global_index_checkpoint_file_manager.h"
#include "paimon/global_index/io/global_index_file_reader.h"
#include "paimon/global_index/io/global_index_file_writer.h"
namespace paimon {
/// Helper class for managing global index files.
-class GlobalIndexFileManager : public GlobalIndexFileReader, public
GlobalIndexFileWriter {
+/// Checkpoint storage is optional and is never accessed by construction or
ordinary index I/O.
+class GlobalIndexFileManager : public GlobalIndexFileReader,
+ public GlobalIndexFileWriter,
+ public GlobalIndexCheckpointFileManager {
public:
- GlobalIndexFileManager(const std::shared_ptr<FileSystem>& fs,
- const std::shared_ptr<IndexPathFactory>&
path_factory)
- : fs_(fs), path_factory_(path_factory) {}
+ GlobalIndexFileManager(
+ const std::shared_ptr<FileSystem>& fs,
+ const std::shared_ptr<IndexPathFactory>& path_factory,
+ std::unique_ptr<IndexCheckpointPathFactory> checkpoint_path_factory =
nullptr)
+ : fs_(fs),
Review Comment:
Please avoid default parameters in production code.
##########
src/paimon/core/global_index/global_index_write_task.cpp:
##########
@@ -387,12 +381,20 @@ Result<std::shared_ptr<CommitMessage>>
GlobalIndexWriteTask::WriteIndex(
std::vector<std::string> writer_field_names =
BuildWriterFieldNames(field_name, extra_fields);
std::vector<std::string> read_field_names =
BuildReadFieldNames(field_name, extra_fields);
- // create index file manager
+ // Checkpoint capability is optional; only plugins enabling checkpoints
will use it.
PAIMON_ASSIGN_OR_RAISE(
- std::shared_ptr<GlobalIndexFileManager> index_file_manager,
- CreateGlobalIndexFileManager(table_path, table_schema, core_options,
pool));
+ std::shared_ptr<FileStorePathFactory> path_factory,
+ CreateFileStorePathFactory(table_path, table_schema, core_options,
pool));
+ std::unique_ptr<IndexCheckpointPathFactory> checkpoint_path_factory;
+ if (task_id && !task_id->empty() && indexer->SupportsCheckpoint()) {
+ PAIMON_ASSIGN_OR_RAISE(checkpoint_path_factory,
+
path_factory->CreateGlobalIndexCheckpointPathFactory(
+ index_type, field_name, range,
task_id.value()));
+ }
+ auto index_file_manager = std::make_shared<GlobalIndexFileManager>(
+ core_options.GetFileSystem(),
path_factory->CreateGlobalIndexFileFactory(),
+ std::move(checkpoint_path_factory));
- // create batch reader
PAIMON_ASSIGN_OR_RAISE(
Review Comment:
Please keep the original comment.
##########
src/paimon/core/global_index/global_index_file_manager.h:
##########
@@ -71,8 +86,105 @@ class GlobalIndexFileManager : public
GlobalIndexFileReader, public GlobalIndexF
return path_factory_->IsExternalPath();
}
+ bool SupportsCheckpoint() const override {
+ return checkpoint_path_factory_ != nullptr;
+ }
+
+ Result<std::unique_ptr<OutputStream>> CreateCheckpointOutputStream() const
override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::string file_name, NewCheckpointFileName());
+
PAIMON_RETURN_NOT_OK(fs_->Mkdirs(checkpoint_path_factory_->GetDirectoryPath()));
+ return fs_->Create(checkpoint_path_factory_->ToPath(file_name),
/*overwrite=*/false);
+ }
+
+ Result<std::unique_ptr<InputStream>> OpenCheckpointInputStream() const
override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::optional<CheckpointFile> checkpoint_file,
+ LatestCheckpointFile());
+ if (!checkpoint_file) {
+ return Status::NotExist("global index checkpoint file does not
exist");
+ }
+ return fs_->Open(checkpoint_file->path);
+ }
+
+ Result<bool> CheckpointExists() const override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::optional<CheckpointFile> checkpoint_file,
+ LatestCheckpointFile());
+ return checkpoint_file.has_value();
+ }
+
+ Status DeleteCheckpoint() const override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::vector<CheckpointFile> checkpoint_files,
ListCheckpointFiles());
+ for (const CheckpointFile& checkpoint_file : checkpoint_files) {
+ PAIMON_RETURN_NOT_OK(fs_->Delete(checkpoint_file.path,
/*recursive=*/false));
+ }
+ return Status::OK();
+ }
+
private:
+ struct CheckpointFile {
+ int64_t id;
+ std::string path;
+ };
+
+ Result<std::string> NewCheckpointFileName() const {
+ if (!checkpoint_file_id_initialized_) {
+ PAIMON_ASSIGN_OR_RAISE(std::optional<CheckpointFile>
checkpoint_file,
+ LatestCheckpointFile());
+ checkpoint_path_factory_->InitializeFileId(checkpoint_file ?
checkpoint_file->id : -1);
Review Comment:
Is this similar to `Init()` func? `checkpoint_path_factory_` and `manager`
seem a bit coupled, so I’d suggest adjusting this a bit.
##########
src/paimon/global_index/lumina/lumina_global_index.cpp:
##########
@@ -896,49 +923,110 @@ Result<std::vector<GlobalIndexIOMeta>>
LuminaIndexWriter::Finish() {
if (indexed_count_ == 0) {
return std::vector<GlobalIndexIOMeta>();
}
- ::lumina::core::MemoryResourceConfig memory_resource(pool_.get());
- PAIMON_ASSIGN_OR_RAISE_FROM_LUMINA(
- ::lumina::api::LuminaBuilder builder,
- ::lumina::api::LuminaBuilder::Create(builder_options_,
memory_resource));
- // pretrain
- LuminaDataset dataset1(indexed_count_, dimension_, array_vec_,
array_start_ids_);
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.PretrainFrom(dataset1));
-
- // insert data
- if (tag_fields_.empty()) {
- LuminaDataset dataset2(indexed_count_, dimension_, array_vec_,
array_start_ids_);
- std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.InsertFrom(dataset2));
- } else {
- ::lumina::extensions::experimental::BuildWithTagExtension
tag_extension;
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.Attach(tag_extension));
- LuminaDatasetWithTag dataset2(indexed_count_, dimension_, array_vec_,
array_start_ids_,
- tag_data_vec_);
- std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
- std::vector<std::vector<TagDimensionData>>().swap(tag_data_vec_);
-
PAIMON_RETURN_NOT_OK_FROM_LUMINA(tag_extension.InsertFromWithTag(dataset2));
+
+ bool had_checkpoint = false;
+ if (checkpoint_file_manager_) {
+ PAIMON_ASSIGN_OR_RAISE(had_checkpoint,
checkpoint_file_manager_->CheckpointExists());
}
+ auto create_build_context = [&]() ->
Result<std::unique_ptr<LuminaBuildContext>> {
+ ::lumina::core::MemoryResourceConfig memory_resource(pool_.get());
+ PAIMON_ASSIGN_OR_RAISE_FROM_LUMINA(
+ ::lumina::api::LuminaBuilder builder,
+ ::lumina::api::LuminaBuilder::Create(builder_options_,
memory_resource));
+ auto context =
std::make_unique<LuminaBuildContext>(std::move(builder));
+ if (checkpoint_file_manager_) {
+ auto checkpoint_manager =
+
std::make_unique<LuminaCheckpointManager>(checkpoint_file_manager_);
+ auto attach_checkpoint = [&](auto* extension) -> Status {
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.Attach(*extension));
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+ extension->LoadCkptManager(std::move(checkpoint_manager)));
+ return Status::OK();
+ };
+ Status checkpoint_status = Status::OK();
+ if (tag_fields_.empty()) {
+ context->checkpoint_extension = std::make_unique<
+
::lumina::extensions::experimental::BuildWithCheckpointExtension>();
+ checkpoint_status =
attach_checkpoint(context->checkpoint_extension.get());
+ } else {
Review Comment:
Is `checkpoint_status` necessary here? Could we just use
`PAIMON_RETURN_NOT_OK` directly?
##########
src/paimon/global_index/lumina/lumina_checkpoint_manager.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/global_index/lumina/lumina_checkpoint_manager.h"
+
+#include <utility>
+
+#include "paimon/global_index/lumina/lumina_file_reader.h"
+#include "paimon/global_index/lumina/lumina_file_writer.h"
+#include "paimon/global_index/lumina/lumina_utils.h"
+
+namespace paimon::lumina {
+
+std::unique_ptr<::lumina::io::FileWriter>
LuminaCheckpointManager::CreateCkptFileWriter() {
+ Result<std::unique_ptr<OutputStream>> output =
file_manager_->CreateCheckpointOutputStream();
+ if (!output.ok()) {
+ return nullptr;
+ }
+ std::shared_ptr<OutputStream> shared_output = std::move(output).value();
+ return std::make_unique<LuminaFileWriter>(shared_output);
+}
+
+::lumina::core::Result<bool> LuminaCheckpointManager::HasCkptFile() {
+ Result<bool> exists = file_manager_->CheckpointExists();
+ if (!exists.ok()) {
+ return
::lumina::core::Result<bool>::Err(PaimonToLuminaStatus(exists.status()));
Review Comment:
`PaimonToLuminaStatus`
##########
src/paimon/global_index/lumina/lumina_global_index.cpp:
##########
@@ -896,49 +923,110 @@ Result<std::vector<GlobalIndexIOMeta>>
LuminaIndexWriter::Finish() {
if (indexed_count_ == 0) {
return std::vector<GlobalIndexIOMeta>();
}
- ::lumina::core::MemoryResourceConfig memory_resource(pool_.get());
- PAIMON_ASSIGN_OR_RAISE_FROM_LUMINA(
- ::lumina::api::LuminaBuilder builder,
- ::lumina::api::LuminaBuilder::Create(builder_options_,
memory_resource));
- // pretrain
- LuminaDataset dataset1(indexed_count_, dimension_, array_vec_,
array_start_ids_);
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.PretrainFrom(dataset1));
-
- // insert data
- if (tag_fields_.empty()) {
- LuminaDataset dataset2(indexed_count_, dimension_, array_vec_,
array_start_ids_);
- std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.InsertFrom(dataset2));
- } else {
- ::lumina::extensions::experimental::BuildWithTagExtension
tag_extension;
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.Attach(tag_extension));
- LuminaDatasetWithTag dataset2(indexed_count_, dimension_, array_vec_,
array_start_ids_,
- tag_data_vec_);
- std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
- std::vector<std::vector<TagDimensionData>>().swap(tag_data_vec_);
-
PAIMON_RETURN_NOT_OK_FROM_LUMINA(tag_extension.InsertFromWithTag(dataset2));
+
+ bool had_checkpoint = false;
+ if (checkpoint_file_manager_) {
+ PAIMON_ASSIGN_OR_RAISE(had_checkpoint,
checkpoint_file_manager_->CheckpointExists());
}
+ auto create_build_context = [&]() ->
Result<std::unique_ptr<LuminaBuildContext>> {
+ ::lumina::core::MemoryResourceConfig memory_resource(pool_.get());
+ PAIMON_ASSIGN_OR_RAISE_FROM_LUMINA(
+ ::lumina::api::LuminaBuilder builder,
+ ::lumina::api::LuminaBuilder::Create(builder_options_,
memory_resource));
+ auto context =
std::make_unique<LuminaBuildContext>(std::move(builder));
+ if (checkpoint_file_manager_) {
+ auto checkpoint_manager =
+
std::make_unique<LuminaCheckpointManager>(checkpoint_file_manager_);
+ auto attach_checkpoint = [&](auto* extension) -> Status {
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.Attach(*extension));
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+ extension->LoadCkptManager(std::move(checkpoint_manager)));
+ return Status::OK();
+ };
+ Status checkpoint_status = Status::OK();
+ if (tag_fields_.empty()) {
+ context->checkpoint_extension = std::make_unique<
+
::lumina::extensions::experimental::BuildWithCheckpointExtension>();
+ checkpoint_status =
attach_checkpoint(context->checkpoint_extension.get());
+ } else {
+ context->checkpoint_tag_extension = std::make_unique<
+
::lumina::extensions::experimental::BuildWithCkptAndTagExtension>();
+ checkpoint_status =
attach_checkpoint(context->checkpoint_tag_extension.get());
+ }
+ if (!checkpoint_status.ok()) {
+ return checkpoint_status;
+ }
+ } else if (!tag_fields_.empty()) {
+ context->tag_extension =
+
std::make_unique<::lumina::extensions::experimental::BuildWithTagExtension>();
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.Attach(*context->tag_extension));
+ }
+ return context;
+ };
+
+ auto build_index = [&]() -> Result<std::unique_ptr<LuminaBuildContext>> {
+ PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<LuminaBuildContext> context,
create_build_context());
+ // pretrain
+ LuminaDataset pretrain_dataset(indexed_count_, dimension_, array_vec_,
array_start_ids_);
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.PretrainFrom(pretrain_dataset));
+ // insert data
+ if (tag_fields_.empty()) {
+ LuminaDataset insert_dataset(indexed_count_, dimension_,
array_vec_, array_start_ids_);
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.InsertFrom(insert_dataset));
+ } else {
+ LuminaDatasetWithTag insert_dataset(indexed_count_, dimension_,
array_vec_,
+ array_start_ids_,
tag_data_vec_);
+ if (context->checkpoint_tag_extension) {
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+
context->checkpoint_tag_extension->InsertFromWithTag(insert_dataset));
+ } else {
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+ context->tag_extension->InsertFromWithTag(insert_dataset));
+ }
+ }
+ return context;
+ };
+
+ Result<std::unique_ptr<LuminaBuildContext>> build_result = build_index();
+ if (!build_result.ok() && had_checkpoint) {
+ LOG(WARNING) << "Failed to build Lumina index with checkpoint, discard
it and rebuild "
+ "from scratch: "
+ << build_result.status().ToString();
+ PAIMON_RETURN_NOT_OK(checkpoint_file_manager_->DeleteCheckpoint());
+ build_result = build_index();
+ }
+ if (!build_result.ok()) {
+ return build_result.status();
+ }
+ std::unique_ptr<LuminaBuildContext> build_context =
std::move(build_result).value();
+ std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
+ std::vector<std::vector<TagDimensionData>>().swap(tag_data_vec_);
+
Review Comment:
rm swap logic.
##########
src/paimon/core/global_index/global_index_file_manager.h:
##########
@@ -71,8 +86,105 @@ class GlobalIndexFileManager : public
GlobalIndexFileReader, public GlobalIndexF
return path_factory_->IsExternalPath();
}
+ bool SupportsCheckpoint() const override {
+ return checkpoint_path_factory_ != nullptr;
+ }
+
+ Result<std::unique_ptr<OutputStream>> CreateCheckpointOutputStream() const
override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::string file_name, NewCheckpointFileName());
+
PAIMON_RETURN_NOT_OK(fs_->Mkdirs(checkpoint_path_factory_->GetDirectoryPath()));
+ return fs_->Create(checkpoint_path_factory_->ToPath(file_name),
/*overwrite=*/false);
+ }
+
+ Result<std::unique_ptr<InputStream>> OpenCheckpointInputStream() const
override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::optional<CheckpointFile> checkpoint_file,
+ LatestCheckpointFile());
+ if (!checkpoint_file) {
+ return Status::NotExist("global index checkpoint file does not
exist");
+ }
+ return fs_->Open(checkpoint_file->path);
+ }
+
+ Result<bool> CheckpointExists() const override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::optional<CheckpointFile> checkpoint_file,
+ LatestCheckpointFile());
+ return checkpoint_file.has_value();
+ }
+
+ Status DeleteCheckpoint() const override {
+ if (!SupportsCheckpoint()) {
+ return Status::Invalid("global index checkpoint storage is not
configured");
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::vector<CheckpointFile> checkpoint_files,
ListCheckpointFiles());
+ for (const CheckpointFile& checkpoint_file : checkpoint_files) {
+ PAIMON_RETURN_NOT_OK(fs_->Delete(checkpoint_file.path,
/*recursive=*/false));
+ }
+ return Status::OK();
Review Comment:
Do we want best-effort deletion here? For example, if there are three files
and the first deletion fails, should we continue trying to delete the other
two, or abort without deleting the rest? Please check the design in paimon-java
for reference.
##########
src/paimon/global_index/lumina/lumina_global_index.cpp:
##########
@@ -896,49 +923,110 @@ Result<std::vector<GlobalIndexIOMeta>>
LuminaIndexWriter::Finish() {
if (indexed_count_ == 0) {
return std::vector<GlobalIndexIOMeta>();
}
- ::lumina::core::MemoryResourceConfig memory_resource(pool_.get());
- PAIMON_ASSIGN_OR_RAISE_FROM_LUMINA(
- ::lumina::api::LuminaBuilder builder,
- ::lumina::api::LuminaBuilder::Create(builder_options_,
memory_resource));
- // pretrain
- LuminaDataset dataset1(indexed_count_, dimension_, array_vec_,
array_start_ids_);
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.PretrainFrom(dataset1));
-
- // insert data
- if (tag_fields_.empty()) {
- LuminaDataset dataset2(indexed_count_, dimension_, array_vec_,
array_start_ids_);
- std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.InsertFrom(dataset2));
- } else {
- ::lumina::extensions::experimental::BuildWithTagExtension
tag_extension;
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.Attach(tag_extension));
- LuminaDatasetWithTag dataset2(indexed_count_, dimension_, array_vec_,
array_start_ids_,
- tag_data_vec_);
- std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
- std::vector<std::vector<TagDimensionData>>().swap(tag_data_vec_);
-
PAIMON_RETURN_NOT_OK_FROM_LUMINA(tag_extension.InsertFromWithTag(dataset2));
+
+ bool had_checkpoint = false;
+ if (checkpoint_file_manager_) {
+ PAIMON_ASSIGN_OR_RAISE(had_checkpoint,
checkpoint_file_manager_->CheckpointExists());
}
+ auto create_build_context = [&]() ->
Result<std::unique_ptr<LuminaBuildContext>> {
+ ::lumina::core::MemoryResourceConfig memory_resource(pool_.get());
+ PAIMON_ASSIGN_OR_RAISE_FROM_LUMINA(
+ ::lumina::api::LuminaBuilder builder,
+ ::lumina::api::LuminaBuilder::Create(builder_options_,
memory_resource));
+ auto context =
std::make_unique<LuminaBuildContext>(std::move(builder));
+ if (checkpoint_file_manager_) {
+ auto checkpoint_manager =
+
std::make_unique<LuminaCheckpointManager>(checkpoint_file_manager_);
+ auto attach_checkpoint = [&](auto* extension) -> Status {
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.Attach(*extension));
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+ extension->LoadCkptManager(std::move(checkpoint_manager)));
+ return Status::OK();
+ };
+ Status checkpoint_status = Status::OK();
+ if (tag_fields_.empty()) {
+ context->checkpoint_extension = std::make_unique<
+
::lumina::extensions::experimental::BuildWithCheckpointExtension>();
+ checkpoint_status =
attach_checkpoint(context->checkpoint_extension.get());
+ } else {
+ context->checkpoint_tag_extension = std::make_unique<
+
::lumina::extensions::experimental::BuildWithCkptAndTagExtension>();
+ checkpoint_status =
attach_checkpoint(context->checkpoint_tag_extension.get());
+ }
+ if (!checkpoint_status.ok()) {
+ return checkpoint_status;
+ }
+ } else if (!tag_fields_.empty()) {
+ context->tag_extension =
+
std::make_unique<::lumina::extensions::experimental::BuildWithTagExtension>();
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.Attach(*context->tag_extension));
+ }
+ return context;
+ };
+
+ auto build_index = [&]() -> Result<std::unique_ptr<LuminaBuildContext>> {
+ PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<LuminaBuildContext> context,
create_build_context());
+ // pretrain
+ LuminaDataset pretrain_dataset(indexed_count_, dimension_, array_vec_,
array_start_ids_);
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.PretrainFrom(pretrain_dataset));
+ // insert data
+ if (tag_fields_.empty()) {
+ LuminaDataset insert_dataset(indexed_count_, dimension_,
array_vec_, array_start_ids_);
+
PAIMON_RETURN_NOT_OK_FROM_LUMINA(context->builder.InsertFrom(insert_dataset));
+ } else {
+ LuminaDatasetWithTag insert_dataset(indexed_count_, dimension_,
array_vec_,
+ array_start_ids_,
tag_data_vec_);
+ if (context->checkpoint_tag_extension) {
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+
context->checkpoint_tag_extension->InsertFromWithTag(insert_dataset));
+ } else {
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+ context->tag_extension->InsertFromWithTag(insert_dataset));
+ }
+ }
+ return context;
+ };
+
+ Result<std::unique_ptr<LuminaBuildContext>> build_result = build_index();
+ if (!build_result.ok() && had_checkpoint) {
+ LOG(WARNING) << "Failed to build Lumina index with checkpoint, discard
it and rebuild "
+ "from scratch: "
+ << build_result.status().ToString();
+ PAIMON_RETURN_NOT_OK(checkpoint_file_manager_->DeleteCheckpoint());
+ build_result = build_index();
+ }
+ if (!build_result.ok()) {
+ return build_result.status();
+ }
+ std::unique_ptr<LuminaBuildContext> build_context =
std::move(build_result).value();
+ std::vector<std::shared_ptr<arrow::FloatArray>>().swap(array_vec_);
+ std::vector<std::vector<TagDimensionData>>().swap(tag_data_vec_);
+
// dump index
PAIMON_ASSIGN_OR_RAISE(std::string index_file_name,
file_manager_->NewFileName(LuminaDefines::kIdentifier));
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<OutputStream> out,
file_manager_->NewOutputStream(index_file_name));
auto file_writer = std::make_unique<LuminaFileWriter>(out);
- PAIMON_RETURN_NOT_OK_FROM_LUMINA(builder.Dump(std::move(file_writer),
io_options_));
+ PAIMON_RETURN_NOT_OK_FROM_LUMINA(
+ build_context->builder.Dump(std::move(file_writer), io_options_));
+
// prepare GlobalIndexIOMeta
PAIMON_ASSIGN_OR_RAISE(int64_t file_size,
file_manager_->GetFileSize(index_file_name));
std::string options_json;
PAIMON_RETURN_NOT_OK(RapidJsonUtil::ToJsonString(lumina_options_,
&options_json));
auto meta_bytes = std::make_shared<Bytes>(options_json,
pool_->GetPaimonPool().get());
GlobalIndexIOMeta meta(file_manager_->ToPath(index_file_name), file_size,
/*metadata=*/meta_bytes);
+ if (checkpoint_file_manager_) {
+ PAIMON_RETURN_NOT_OK(checkpoint_file_manager_->DeleteCheckpoint());
+ }
return std::vector<GlobalIndexIOMeta>({meta});
Review Comment:
I assume deletion failures are meant to be silent here? In paimon-java, many
deletions are done in destructors and are also silent.
--
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]