lszskye commented on code in PR #382:
URL: https://github.com/apache/paimon-cpp/pull/382#discussion_r4081173317
##########
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:
make sense
--
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]