wangyong9999 commented on code in PR #245: URL: https://github.com/apache/paimon-cpp/pull/245#discussion_r3851760354
########## src/paimon/core/index/pk/bucketed_primary_key_index_maintainer.cpp: ########## @@ -0,0 +1,261 @@ +/* + * 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/core/index/pk/bucketed_primary_key_index_maintainer.h" + +#include <algorithm> +#include <map> +#include <set> +#include <unordered_set> +#include <utility> + +#include "fmt/format.h" +#include "paimon/core/index/global_index_meta.h" +#include "paimon/core/index/pk/primary_key_index_source_policy.h" +#include "paimon/core/index/pksorted/pk_sorted_bucket_index_state.h" +#include "paimon/core/index/pksorted/pk_sorted_index_builder.h" +#include "paimon/core/index/pksorted/pk_sorted_index_group.h" +#include "paimon/core/utils/commit_increment.h" +#include "paimon/logging.h" + +namespace paimon { +namespace { + +Logger* GetLogger() { + static std::unique_ptr<Logger> logger = Logger::GetLogger("BucketedPrimaryKeyIndexMaintainer"); + return logger.get(); +} + +void RemoveDataFiles(const std::vector<std::shared_ptr<DataFileMeta>>& files, + std::map<std::string, std::shared_ptr<DataFileMeta>>* active) { + for (const std::shared_ptr<DataFileMeta>& file : files) { + if (file != nullptr) { + active->erase(file->file_name); + } + } +} + +Status AddSourceFiles(const std::vector<std::shared_ptr<DataFileMeta>>& files, + std::map<std::string, std::shared_ptr<DataFileMeta>>* active) { + for (const std::shared_ptr<DataFileMeta>& file : files) { + if (file == nullptr) { + return Status::Invalid("Primary-key index data increment contains a null file."); + } + if (PrimaryKeyIndexSourcePolicy::ShouldRead(*file)) { + (*active)[file->file_name] = file; + } + } + return Status::OK(); +} + +Status ValidateAppendFiles(const std::vector<std::shared_ptr<DataFileMeta>>& files) { + for (const std::shared_ptr<DataFileMeta>& file : files) { + if (file == nullptr) { + return Status::Invalid("Primary-key index append increment contains a null file."); + } + if (PrimaryKeyIndexSourcePolicy::ShouldRead(*file)) { + return Status::Invalid(fmt::format( + "Append file {} must not be a primary-key sorted-index source.", file->file_name)); + } + } + return Status::OK(); +} + +std::string PayloadIdentity(const std::shared_ptr<IndexFileMeta>& payload) { + if (payload == nullptr) { + return std::string(); + } + return payload->ExternalPath().value_or(payload->FileName()); +} + +void AddUniquePayload(const std::shared_ptr<IndexFileMeta>& payload, + std::unordered_set<std::string>* identities, + std::vector<std::shared_ptr<IndexFileMeta>>* payloads) { + std::string identity = PayloadIdentity(payload); + if (!identity.empty() && identities->insert(identity).second) { + payloads->push_back(payload); + } +} + +} // namespace + +Result<std::shared_ptr<BucketedPrimaryKeyIndexMaintainer::Factory>> +BucketedPrimaryKeyIndexMaintainer::Factory::Create( + const std::string& root_path, const std::string& branch, + const std::shared_ptr<TableSchema>& table_schema, + const std::vector<PrimaryKeyIndexDefinition>& definitions, + const std::shared_ptr<FileStorePathFactory>& path_factory, + const std::shared_ptr<IndexFileHandler>& index_file_handler, const CoreOptions& options, + const std::shared_ptr<IOManager>& io_manager, bool enable_multi_thread_spill, + const std::shared_ptr<Executor>& executor, const std::shared_ptr<MemoryPool>& pool) { + std::vector<PrimaryKeyIndexDefinition> btree_definitions; + for (const PrimaryKeyIndexDefinition& definition : definitions) { + if (definition.GetFamily() == PrimaryKeyIndexDefinition::Family::BTREE) { + btree_definitions.push_back(definition); + } + } + std::sort(btree_definitions.begin(), btree_definitions.end(), + [](const PrimaryKeyIndexDefinition& left, const PrimaryKeyIndexDefinition& right) { + return left.FieldId() < right.FieldId(); + }); + return std::shared_ptr<Factory>(new Factory( + root_path, branch, table_schema, std::move(btree_definitions), path_factory, + index_file_handler, options, io_manager, enable_multi_thread_spill, executor, pool)); +} + +Result<std::shared_ptr<BucketedPrimaryKeyIndexMaintainer>> +BucketedPrimaryKeyIndexMaintainer::Factory::CreateMaintainer( + const BinaryRow& partition, int32_t bucket, + const std::vector<std::shared_ptr<DataFileMeta>>& restored_data_files, + const std::vector<std::shared_ptr<IndexFileMeta>>& restored_payloads) const { + std::map<std::string, std::shared_ptr<DataFileMeta>> active_data_files; + PAIMON_RETURN_NOT_OK(AddSourceFiles(restored_data_files, &active_data_files)); + std::vector<FieldMaintainer> fields; + fields.reserve(definitions_.size()); + for (const PrimaryKeyIndexDefinition& definition : definitions_) { + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr<PkSortedIndexBuilder> builder, + PkSortedIndexBuilder::Create(root_path_, branch_, partition, bucket, table_schema_, + definition, path_factory_, options_, io_manager_, + enable_multi_thread_spill_, executor_, pool_)); + fields.push_back( + FieldMaintainer{definition, std::shared_ptr<PkSortedIndexBuilder>(std::move(builder))}); + } + return std::shared_ptr<BucketedPrimaryKeyIndexMaintainer>(new BucketedPrimaryKeyIndexMaintainer( + std::move(fields), std::move(active_data_files), restored_payloads)); +} + +Status BucketedPrimaryKeyIndexMaintainer::PrepareCommit(CommitIncrement* increment) { + if (increment == nullptr) { + return Status::Invalid("Primary-key index commit increment is null."); + } + auto previous_data_files = active_data_files_; + const DataIncrement& data_increment = increment->GetNewFilesIncrement(); + const CompactIncrement& compact_increment = increment->GetCompactIncrement(); + PAIMON_RETURN_NOT_OK(ValidateAppendFiles(data_increment.NewFiles())); + RemoveDataFiles(compact_increment.CompactBefore(), &active_data_files_); + Status update_status = AddSourceFiles(compact_increment.CompactAfter(), &active_data_files_); + if (!update_status.ok()) { + active_data_files_ = std::move(previous_data_files); + return update_status; + } + + std::vector<std::shared_ptr<DataFileMeta>> active_data; + active_data.reserve(active_data_files_.size()); + for (const auto& file : active_data_files_) { + active_data.push_back(file.second); + } + + std::vector<std::shared_ptr<IndexFileMeta>> deleted_payloads; + std::vector<std::shared_ptr<IndexFileMeta>> new_payloads; + std::unordered_set<std::string> deleted_identities; + std::unordered_set<std::string> new_identities; + + Status build_status = Status::OK(); + std::string failed_column; + for (const FieldMaintainer& field : fields_) { + std::vector<std::shared_ptr<IndexFileMeta>> field_payloads; + for (const std::shared_ptr<IndexFileMeta>& payload : active_payloads_) { + if (payload == nullptr || payload->IndexType() != field.definition.IndexType()) { + continue; + } + const std::optional<GlobalIndexMeta>& meta = payload->GetGlobalIndexMeta(); + if (meta != std::nullopt && meta->index_field_id == field.definition.FieldId()) { + field_payloads.push_back(payload); + } + } + PkSortedBucketIndexState state = PkSortedBucketIndexState::FromActiveDataFiles( + field.definition.FieldId(), field.definition.IndexType(), active_data, field_payloads); + std::set<int32_t> current_levels; + for (const std::shared_ptr<PkSortedIndexGroup>& group : state.Groups()) { + current_levels.insert(group->DataLevel()); + } + for (const std::shared_ptr<IndexFileMeta>& rejected : state.RejectedPayloads()) { + AddUniquePayload(rejected, &deleted_identities, &deleted_payloads); + } + + std::map<int32_t, std::vector<std::shared_ptr<DataFileMeta>>> desired_by_level; + for (const std::shared_ptr<DataFileMeta>& file : active_data) { + if (file != nullptr && PrimaryKeyIndexSourcePolicy::ShouldRead(*file)) { + desired_by_level[file->level].push_back(file); + } + } + for (auto& level_files : desired_by_level) { + std::sort(level_files.second.begin(), level_files.second.end(), + [](const std::shared_ptr<DataFileMeta>& left, + const std::shared_ptr<DataFileMeta>& right) { + return left->file_name < right->file_name; + }); + if (current_levels.count(level_files.first) > 0) { + continue; + } + Result<std::shared_ptr<IndexFileMeta>> build_result = + field.builder->Build(level_files.second); + if (!build_result.ok()) { + build_status = build_result.status(); + failed_column = field.definition.Column(); + break; + } + AddUniquePayload(std::move(build_result).value(), &new_identities, &new_payloads); + } + if (!build_status.ok()) { + break; + } + } + + if (!build_status.ok()) { + for (const std::shared_ptr<IndexFileMeta>& payload : new_payloads) { + for (const FieldMaintainer& field : fields_) { + const std::optional<GlobalIndexMeta>& meta = payload->GetGlobalIndexMeta(); + if (meta != std::nullopt && meta->index_field_id == field.definition.FieldId()) { + [[maybe_unused]] Status cleanup_status = field.builder->DeletePayload(payload); + break; + } + } + } + PAIMON_LOG_WARN(GetLogger(), + "Failed to build primary-key BTree index for column %s; committing data " + "files without new index payloads. %s", + failed_column.c_str(), build_status.ToString().c_str()); + return Status::OK(); + } + + std::vector<std::shared_ptr<IndexFileMeta>> next_payloads; + next_payloads.reserve(active_payloads_.size() + new_payloads.size()); + for (const std::shared_ptr<IndexFileMeta>& payload : active_payloads_) { + if (deleted_identities.count(PayloadIdentity(payload)) == 0) { + next_payloads.push_back(payload); + } + } + next_payloads.insert(next_payloads.end(), new_payloads.begin(), new_payloads.end()); + active_payloads_ = std::move(next_payloads); + + bool has_compaction_transition = Review Comment: Two writers based on the same snapshot can each publish a different payload for the same uncovered level. The second commit does not compare source metadata, and GlobalFileNameCombiner keeps both filenames; the next restore rejects the whole level because it has two candidates and falls back to a full scan. Deduplicate or conflict-check by partition/bucket/type/field/source set during commit retry. ########## src/paimon/core/index/pk/bucketed_primary_key_index_maintainer.cpp: ########## @@ -0,0 +1,261 @@ +/* + * 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/core/index/pk/bucketed_primary_key_index_maintainer.h" + +#include <algorithm> +#include <map> +#include <set> +#include <unordered_set> +#include <utility> + +#include "fmt/format.h" +#include "paimon/core/index/global_index_meta.h" +#include "paimon/core/index/pk/primary_key_index_source_policy.h" +#include "paimon/core/index/pksorted/pk_sorted_bucket_index_state.h" +#include "paimon/core/index/pksorted/pk_sorted_index_builder.h" +#include "paimon/core/index/pksorted/pk_sorted_index_group.h" +#include "paimon/core/utils/commit_increment.h" +#include "paimon/logging.h" + +namespace paimon { +namespace { + +Logger* GetLogger() { + static std::unique_ptr<Logger> logger = Logger::GetLogger("BucketedPrimaryKeyIndexMaintainer"); + return logger.get(); +} + +void RemoveDataFiles(const std::vector<std::shared_ptr<DataFileMeta>>& files, + std::map<std::string, std::shared_ptr<DataFileMeta>>* active) { + for (const std::shared_ptr<DataFileMeta>& file : files) { + if (file != nullptr) { + active->erase(file->file_name); + } + } +} + +Status AddSourceFiles(const std::vector<std::shared_ptr<DataFileMeta>>& files, + std::map<std::string, std::shared_ptr<DataFileMeta>>* active) { + for (const std::shared_ptr<DataFileMeta>& file : files) { + if (file == nullptr) { + return Status::Invalid("Primary-key index data increment contains a null file."); + } + if (PrimaryKeyIndexSourcePolicy::ShouldRead(*file)) { + (*active)[file->file_name] = file; + } + } + return Status::OK(); +} + +Status ValidateAppendFiles(const std::vector<std::shared_ptr<DataFileMeta>>& files) { + for (const std::shared_ptr<DataFileMeta>& file : files) { + if (file == nullptr) { + return Status::Invalid("Primary-key index append increment contains a null file."); + } + if (PrimaryKeyIndexSourcePolicy::ShouldRead(*file)) { + return Status::Invalid(fmt::format( + "Append file {} must not be a primary-key sorted-index source.", file->file_name)); + } + } + return Status::OK(); +} + +std::string PayloadIdentity(const std::shared_ptr<IndexFileMeta>& payload) { + if (payload == nullptr) { + return std::string(); + } + return payload->ExternalPath().value_or(payload->FileName()); +} + +void AddUniquePayload(const std::shared_ptr<IndexFileMeta>& payload, + std::unordered_set<std::string>* identities, + std::vector<std::shared_ptr<IndexFileMeta>>* payloads) { + std::string identity = PayloadIdentity(payload); + if (!identity.empty() && identities->insert(identity).second) { + payloads->push_back(payload); + } +} + +} // namespace + +Result<std::shared_ptr<BucketedPrimaryKeyIndexMaintainer::Factory>> +BucketedPrimaryKeyIndexMaintainer::Factory::Create( + const std::string& root_path, const std::string& branch, + const std::shared_ptr<TableSchema>& table_schema, + const std::vector<PrimaryKeyIndexDefinition>& definitions, + const std::shared_ptr<FileStorePathFactory>& path_factory, + const std::shared_ptr<IndexFileHandler>& index_file_handler, const CoreOptions& options, + const std::shared_ptr<IOManager>& io_manager, bool enable_multi_thread_spill, + const std::shared_ptr<Executor>& executor, const std::shared_ptr<MemoryPool>& pool) { + std::vector<PrimaryKeyIndexDefinition> btree_definitions; + for (const PrimaryKeyIndexDefinition& definition : definitions) { + if (definition.GetFamily() == PrimaryKeyIndexDefinition::Family::BTREE) { + btree_definitions.push_back(definition); + } + } + std::sort(btree_definitions.begin(), btree_definitions.end(), + [](const PrimaryKeyIndexDefinition& left, const PrimaryKeyIndexDefinition& right) { + return left.FieldId() < right.FieldId(); + }); + return std::shared_ptr<Factory>(new Factory( + root_path, branch, table_schema, std::move(btree_definitions), path_factory, + index_file_handler, options, io_manager, enable_multi_thread_spill, executor, pool)); +} + +Result<std::shared_ptr<BucketedPrimaryKeyIndexMaintainer>> +BucketedPrimaryKeyIndexMaintainer::Factory::CreateMaintainer( + const BinaryRow& partition, int32_t bucket, + const std::vector<std::shared_ptr<DataFileMeta>>& restored_data_files, + const std::vector<std::shared_ptr<IndexFileMeta>>& restored_payloads) const { + std::map<std::string, std::shared_ptr<DataFileMeta>> active_data_files; + PAIMON_RETURN_NOT_OK(AddSourceFiles(restored_data_files, &active_data_files)); + std::vector<FieldMaintainer> fields; + fields.reserve(definitions_.size()); + for (const PrimaryKeyIndexDefinition& definition : definitions_) { + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr<PkSortedIndexBuilder> builder, + PkSortedIndexBuilder::Create(root_path_, branch_, partition, bucket, table_schema_, + definition, path_factory_, options_, io_manager_, + enable_multi_thread_spill_, executor_, pool_)); + fields.push_back( + FieldMaintainer{definition, std::shared_ptr<PkSortedIndexBuilder>(std::move(builder))}); + } + return std::shared_ptr<BucketedPrimaryKeyIndexMaintainer>(new BucketedPrimaryKeyIndexMaintainer( + std::move(fields), std::move(active_data_files), restored_payloads)); +} + +Status BucketedPrimaryKeyIndexMaintainer::PrepareCommit(CommitIncrement* increment) { + if (increment == nullptr) { + return Status::Invalid("Primary-key index commit increment is null."); + } + auto previous_data_files = active_data_files_; + const DataIncrement& data_increment = increment->GetNewFilesIncrement(); + const CompactIncrement& compact_increment = increment->GetCompactIncrement(); + PAIMON_RETURN_NOT_OK(ValidateAppendFiles(data_increment.NewFiles())); + RemoveDataFiles(compact_increment.CompactBefore(), &active_data_files_); + Status update_status = AddSourceFiles(compact_increment.CompactAfter(), &active_data_files_); + if (!update_status.ok()) { + active_data_files_ = std::move(previous_data_files); + return update_status; + } + + std::vector<std::shared_ptr<DataFileMeta>> active_data; + active_data.reserve(active_data_files_.size()); + for (const auto& file : active_data_files_) { + active_data.push_back(file.second); + } + + std::vector<std::shared_ptr<IndexFileMeta>> deleted_payloads; + std::vector<std::shared_ptr<IndexFileMeta>> new_payloads; + std::unordered_set<std::string> deleted_identities; + std::unordered_set<std::string> new_identities; + + Status build_status = Status::OK(); + std::string failed_column; + for (const FieldMaintainer& field : fields_) { + std::vector<std::shared_ptr<IndexFileMeta>> field_payloads; + for (const std::shared_ptr<IndexFileMeta>& payload : active_payloads_) { + if (payload == nullptr || payload->IndexType() != field.definition.IndexType()) { + continue; + } + const std::optional<GlobalIndexMeta>& meta = payload->GetGlobalIndexMeta(); + if (meta != std::nullopt && meta->index_field_id == field.definition.FieldId()) { + field_payloads.push_back(payload); + } + } + PkSortedBucketIndexState state = PkSortedBucketIndexState::FromActiveDataFiles( + field.definition.FieldId(), field.definition.IndexType(), active_data, field_payloads); + std::set<int32_t> current_levels; + for (const std::shared_ptr<PkSortedIndexGroup>& group : state.Groups()) { + current_levels.insert(group->DataLevel()); + } + for (const std::shared_ptr<IndexFileMeta>& rejected : state.RejectedPayloads()) { + AddUniquePayload(rejected, &deleted_identities, &deleted_payloads); + } + + std::map<int32_t, std::vector<std::shared_ptr<DataFileMeta>>> desired_by_level; + for (const std::shared_ptr<DataFileMeta>& file : active_data) { + if (file != nullptr && PrimaryKeyIndexSourcePolicy::ShouldRead(*file)) { + desired_by_level[file->level].push_back(file); + } + } + for (auto& level_files : desired_by_level) { + std::sort(level_files.second.begin(), level_files.second.end(), + [](const std::shared_ptr<DataFileMeta>& left, + const std::shared_ptr<DataFileMeta>& right) { + return left->file_name < right->file_name; + }); + if (current_levels.count(level_files.first) > 0) { + continue; + } + Result<std::shared_ptr<IndexFileMeta>> build_result = + field.builder->Build(level_files.second); + if (!build_result.ok()) { + build_status = build_result.status(); + failed_column = field.definition.Column(); + break; + } + AddUniquePayload(std::move(build_result).value(), &new_identities, &new_payloads); + } + if (!build_status.ok()) { + break; + } + } + + if (!build_status.ok()) { + for (const std::shared_ptr<IndexFileMeta>& payload : new_payloads) { + for (const FieldMaintainer& field : fields_) { + const std::optional<GlobalIndexMeta>& meta = payload->GetGlobalIndexMeta(); + if (meta != std::nullopt && meta->index_field_id == field.definition.FieldId()) { + [[maybe_unused]] Status cleanup_status = field.builder->DeletePayload(payload); + break; + } + } + } + PAIMON_LOG_WARN(GetLogger(), + "Failed to build primary-key BTree index for column %s; committing data " + "files without new index payloads. %s", + failed_column.c_str(), build_status.ToString().c_str()); + return Status::OK(); + } + + std::vector<std::shared_ptr<IndexFileMeta>> next_payloads; + next_payloads.reserve(active_payloads_.size() + new_payloads.size()); + for (const std::shared_ptr<IndexFileMeta>& payload : active_payloads_) { Review Comment: active_payloads_ can contain payloads for a field removed from pk-btree.index.columns, but only current fields_ populate deleted_identities, so this loop retains the old payload forever; removing the last definition skips creating the maintainer entirely. Either reject that schema change, or emit DELETEs for restored BTree payloads with no current owner, including the zero-definition case. ########## src/paimon/core/operation/expire_snapshots.cpp: ########## @@ -185,6 +189,7 @@ Result<int32_t> ExpireSnapshots::ExpireUntil(int64_t earliest_snapshot_id, continue; } PAIMON_ASSIGN_OR_RAISE(Snapshot snapshot, snapshot_manager_->LoadSnapshot(id)); + PAIMON_RETURN_NOT_OK(CleanUnusedIndexManifest(snapshot.IndexManifest(), &skipping_sets)); Review Comment: createBranch(branch, tag) copies the tag into the branch, but deleting the source-branch tag does not check branch ownership. The branch snapshot and copied tag still reference the same table-root payload, and main-branch expiration can delete it. Please traverse branch metadata or reject expiration while another branch exists. ########## src/paimon/core/operation/expire_snapshots.cpp: ########## @@ -290,6 +295,50 @@ Status ExpireSnapshots::CleanUnusedManifests(const std::string& manifest_list_na return Status::OK(); } +Status ExpireSnapshots::CleanUnusedIndexManifest(const std::optional<std::string>& index_manifest, + std::set<std::string>* skipping_manifest_set) { + if (!index_manifest || index_manifest->empty() || + skipping_manifest_set->count(index_manifest.value()) > 0) { + return Status::OK(); + } + if (index_manifest_file_ == nullptr) { + return Status::Invalid("index manifest file is null"); + } + + std::vector<IndexManifestEntry> entries; + Status read_status = + index_manifest_file_->ReadIfFileExist(index_manifest.value(), /*filter=*/nullptr, &entries); + if (read_status.IsNotExist()) { + return Status::OK(); + } + PAIMON_RETURN_NOT_OK(read_status); + + std::vector<std::pair<std::string, std::string>> index_files_to_delete; + std::set<std::string> planned_file_names; + for (const IndexManifestEntry& entry : entries) { + const std::string& file_name = entry.index_file->FileName(); + if (skipping_manifest_set->count(file_name) > 0 || + !planned_file_names.insert(file_name).second) { + continue; + } + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr<IndexPathFactory> index_path_factory, + path_factory_->CreateIndexFileFactory(entry.partition, entry.bucket)); + index_files_to_delete.emplace_back(file_name, index_path_factory->ToPath(entry.index_file)); + } + + for (const auto& [file_name, file_path] : index_files_to_delete) { + skipping_manifest_set->insert(file_name); + auto delete_status = fs_->Delete(file_path); + // Index payload cleanup is best effort, consistent with data file expiration. + (void)delete_status; + } + + skipping_manifest_set->insert(index_manifest.value()); + index_manifest_file_->DeleteQuietly(index_manifest.value()); Review Comment: Java parity does not restore a retry path here: this PR explicitly keeps global-index external paths out of orphan cleanup. Once this manifest is deleted, a transient payload-delete failure is no longer discoverable. Keep a durable retry record or retain the manifest until external deletes succeed or are confirmed absent. ########## src/paimon/core/index/pksorted/pk_sorted_index_file.cpp: ########## @@ -95,24 +139,95 @@ Result<std::shared_ptr<IndexFileMeta>> PkSortedIndexFile::Build( ScopeGuard array_guard([&]() { ArrowArrayRelease(&c_array); }); PAIMON_RETURN_NOT_OK(writer->AddBatch(&c_array, std::move(sorted_ordinals))); PAIMON_ASSIGN_OR_RAISE(std::vector<GlobalIndexIOMeta> io_metas, writer->Finish()); - if (io_metas.size() != 1) { - return Status::Invalid(fmt::format( - "Sorted index build must produce exactly one payload file, but produced {}.", - io_metas.size())); + return FinishIndexFile(field.Id(), index_type, source_row_count, source_meta, io_metas, + is_external_path, pool); +} + +Result<std::shared_ptr<IndexFileMeta>> PkSortedIndexFile::BuildFromSortedReader( + const DataField& field, const std::string& index_type, + const std::map<std::string, std::string>& options, int32_t data_level, + const std::vector<PrimaryKeyIndexSourceFile>& source_files, + std::unique_ptr<SortMergeReader>&& sorted_reader, + const std::shared_ptr<GlobalIndexFileWriter>& file_writer, bool is_external_path, + int32_t write_batch_size, const std::shared_ptr<MemoryPool>& pool) { + if (sorted_reader == nullptr) { + return Status::Invalid("Sorted index reader is null."); } - const GlobalIndexIOMeta& io_meta = io_metas[0]; + if (write_batch_size <= 0) { + return Status::Invalid("Sorted index write batch size must be positive."); + } + PAIMON_ASSIGN_OR_RAISE(PrimaryKeyIndexSourceMeta source_meta, + PrimaryKeyIndexSourceMeta::Create(data_level, source_files)); + PAIMON_ASSIGN_OR_RAISE(int64_t source_row_count, ValidateAndCountSourceRows(source_files)); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<GlobalIndexer> indexer, + GlobalIndexerFactory::Get(index_type, options)); + if (indexer == nullptr) { + return Status::Invalid(fmt::format("Index type {} is not registered.", index_type)); + } + auto arrow_schema = arrow::schema({DataField::ConvertDataFieldToArrowField(field)}); + ArrowSchema c_arrow_schema; + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*arrow_schema, &c_arrow_schema)); + ScopeGuard schema_guard([&]() { ArrowSchemaRelease(&c_arrow_schema); }); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<GlobalIndexWriter> writer, + indexer->CreateWriter(field.Name(), &c_arrow_schema, file_writer, pool)); - PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<Bytes> source_meta_bytes, source_meta.Serialize(pool)); - std::optional<std::string> external_path; - if (is_external_path) { - PAIMON_ASSIGN_OR_RAISE(Path path, PathUtil::ToPath(io_meta.file_path)); - external_path = path.ToString(); + auto projection_schema = SpecialFields::CompleteSequenceAndValueKindField(arrow_schema); + auto create_consumer = + [projection_schema, + pool]() -> Result<std::unique_ptr<RowToArrowArrayConverter<KeyValue, KeyValueBatch>>> { + return KeyValueMetaProjectionConsumer::Create(projection_schema, pool); + }; + auto producer = std::make_unique<AsyncKeyValueProducerAndConsumer<KeyValue, KeyValueBatch>>( + std::move(sorted_reader), create_consumer, write_batch_size, + /*projection_thread_num=*/1, pool); + ScopeGuard close_guard([&]() { producer->Close(); }); + int64_t rows_written = 0; + while (true) { + PAIMON_ASSIGN_OR_RAISE(KeyValueBatch key_value_batch, producer->NextBatch()); + if (key_value_batch.batch == nullptr) { + break; + } + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( + std::shared_ptr<arrow::RecordBatch> record_batch, + arrow::ImportRecordBatch(key_value_batch.batch.get(), projection_schema)); + if (record_batch->num_columns() != 3 || + record_batch->column(0)->type_id() != arrow::Type::INT64) { + return Status::Invalid("Sorted index projection produced an invalid batch."); + } + auto sequence_numbers = checked_pointer_cast<arrow::Int64Array>(record_batch->column(0)); + std::vector<int64_t> ordinals; + ordinals.reserve(static_cast<size_t>(sequence_numbers->length())); + for (int64_t index = 0; index < sequence_numbers->length(); ++index) { + if (sequence_numbers->IsNull(index)) { + return Status::Invalid("Sorted index row id must not be null."); + } + int64_t ordinal = sequence_numbers->Value(index); + if (ordinal < 0 || ordinal >= source_row_count) { + return Status::Invalid( + fmt::format("Row id {} is outside sorted index group row range [0, {}).", + ordinal, source_row_count)); + } + ordinals.push_back(ordinal); + } + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( + std::shared_ptr<arrow::StructArray> values, + arrow::StructArray::Make({record_batch->column(2)}, {field.Name()})); + ArrowArray c_array; + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportArray(*values, &c_array)); + ScopeGuard array_guard([&]() { ArrowArrayRelease(&c_array); }); + PAIMON_RETURN_NOT_OK(writer->AddBatch(&c_array, std::move(ordinals))); Review Comment: The shared format does not require accepting process OOM. A hard memory-budget check can return Status before current_row_ids_ and the encoding buffer allocate; the new maintainer fallback then leaves the level uncovered. Please add that bound even if chunking is deferred. -- 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]
