lszskye commented on code in PR #245: URL: https://github.com/apache/paimon-cpp/pull/245#discussion_r4004039248
########## src/paimon/core/index/pksorted/pk_sorted_index_builder.cpp: ########## @@ -0,0 +1,270 @@ +/* + * 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/pksorted/pk_sorted_index_builder.h" + +#include <algorithm> +#include <map> +#include <string> +#include <utility> + +#include "arrow/api.h" +#include "arrow/array/concatenate.h" +#include "arrow/c/bridge.h" +#include "fmt/format.h" +#include "paimon/common/table/special_fields.h" +#include "paimon/common/utils/arrow/mem_utils.h" +#include "paimon/common/utils/arrow/status_utils.h" +#include "paimon/common/utils/checked_cast.h" +#include "paimon/common/utils/fields_comparator.h" +#include "paimon/common/utils/scope_guard.h" +#include "paimon/core/casting/casting_utils.h" +#include "paimon/core/global_index/global_index_file_manager.h" +#include "paimon/core/index/pk/primary_key_index_source_file.h" +#include "paimon/core/index/pk/primary_key_index_source_policy.h" +#include "paimon/core/index/pksorted/pk_sorted_data_file_reader.h" +#include "paimon/core/index/pksorted/pk_sorted_index_file.h" +#include "paimon/core/io/data_file_meta.h" +#include "paimon/core/mergetree/compact/sort_merge_reader_with_min_heap.h" +#include "paimon/core/mergetree/external_sort_buffer.h" +#include "paimon/core/mergetree/in_memory_sort_buffer.h" +#include "paimon/core/mergetree/sort_buffer.h" +#include "paimon/core/schema/table_schema.h" +#include "paimon/core/utils/file_store_path_factory.h" +#include "paimon/fs/file_system.h" +#include "paimon/global_index/io/global_index_file_writer.h" +#include "paimon/record_batch.h" + +namespace paimon { +namespace { + +class TrackingGlobalIndexFileWriter : public GlobalIndexFileWriter { + public: + explicit TrackingGlobalIndexFileWriter(const std::shared_ptr<GlobalIndexFileManager>& delegate) + : delegate_(delegate) {} + + Result<std::string> NewFileName(const std::string& prefix) const override { + PAIMON_ASSIGN_OR_RAISE(std::string file_name, delegate_->NewFileName(prefix)); + created_file_names_.push_back(file_name); + return file_name; + } + + Result<std::unique_ptr<OutputStream>> NewOutputStream( + const std::string& file_name) const override { + return delegate_->NewOutputStream(file_name); + } + + Result<int64_t> GetFileSize(const std::string& file_name) const override { + return delegate_->GetFileSize(file_name); + } + + std::string ToPath(const std::string& file_name) const override { + return delegate_->ToPath(file_name); + } + + void Cleanup(const std::shared_ptr<FileSystem>& fs) const { + for (const std::string& file_name : created_file_names_) { + [[maybe_unused]] Status status = fs->Delete(delegate_->ToPath(file_name)); + } + } + + private: + std::shared_ptr<GlobalIndexFileManager> delegate_; + mutable std::vector<std::string> created_file_names_; +}; + +} // namespace + +Result<std::unique_ptr<PkSortedIndexBuilder>> PkSortedIndexBuilder::Create( + const std::string& root_path, const std::string& branch, const BinaryRow& partition, + int32_t bucket, const std::shared_ptr<TableSchema>& table_schema, + const PrimaryKeyIndexDefinition& definition, + const std::shared_ptr<FileStorePathFactory>& path_factory, 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) { + if (definition.GetFamily() != PrimaryKeyIndexDefinition::Family::BTREE) { + return Status::Invalid("PkSortedIndexBuilder only supports BTree definitions."); + } + PAIMON_ASSIGN_OR_RAISE(DataField field, table_schema->GetField(definition.FieldId())); + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr<PkSortedDataFileReader> data_file_reader, + PkSortedDataFileReader::Create(root_path, table_schema, definition.FieldId(), path_factory, + branch, options, executor, pool)); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<IndexPathFactory> index_path_factory, + path_factory->CreateIndexFileFactory(partition, bucket)); + return std::unique_ptr<PkSortedIndexBuilder>(new PkSortedIndexBuilder( + partition, bucket, std::move(field), definition, + std::shared_ptr<PkSortedDataFileReader>(std::move(data_file_reader)), + options.GetFileSystem(), std::shared_ptr<IndexPathFactory>(std::move(index_path_factory)), + options, io_manager, enable_multi_thread_spill, pool)); +} + +Result<std::shared_ptr<IndexFileMeta>> PkSortedIndexBuilder::Build( Review Comment: The function is too long. Consider splitting it into smaller, focused functions to improve readability and maintainability. -- 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]
