This is an automated email from the ASF dual-hosted git repository.
liuneng pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-gluten.git
The following commit(s) were added to refs/heads/main by this push:
new 40f2b26df [CH] Speed up mergetree metadata reloading
40f2b26df is described below
commit 40f2b26df0cd5630fd00ef79679a1d6c19051616
Author: LiuNeng <[email protected]>
AuthorDate: Tue Apr 23 17:12:48 2024 +0800
[CH] Speed up mergetree metadata reloading
What changes were proposed in this pull request?
Speed up mergetree metadata reloading. The principle is to increase the
writing speed of metadata through the thread pool.
The default thread pool size is minimum 10, which can be specified through
spark.gluten.sql.columnar.backend.ch.runtimesettings.max_threads
2000+parts
before:
runTPCHQueryBySQL cost time: 136886 ms
after:
runTPCHQueryBySQL cost time: 12017 ms
How was this patch tested?
manual tests
---
.../Storages/Mergetree/MetaDataHelper.cpp | 56 +++++++++++++++-------
1 file changed, 39 insertions(+), 17 deletions(-)
diff --git a/cpp-ch/local-engine/Storages/Mergetree/MetaDataHelper.cpp
b/cpp-ch/local-engine/Storages/Mergetree/MetaDataHelper.cpp
index 688ca8be8..21c0fc968 100644
--- a/cpp-ch/local-engine/Storages/Mergetree/MetaDataHelper.cpp
+++ b/cpp-ch/local-engine/Storages/Mergetree/MetaDataHelper.cpp
@@ -19,6 +19,13 @@
#include <Disks/ObjectStorages/MetadataStorageFromDisk.h>
+namespace CurrentMetrics
+{
+extern const Metric LocalThread;
+extern const Metric LocalThreadActive;
+extern const Metric LocalThreadScheduled;
+}
+
using namespace DB;
namespace local_engine
@@ -68,6 +75,17 @@ void restoreMetaData(CustomStorageMergeTreePtr & storage,
const MergeTreeTable &
if (not_exists_part.empty())
return;
+ // Increase the speed of metadata recovery
+ auto max_concurrency = std::max(10UL,
SerializedPlanParser::global_context->getSettings().max_threads.value);
+ auto max_threads = std::min(max_concurrency, not_exists_part.size());
+ FreeThreadPool thread_pool(
+ CurrentMetrics::LocalThread,
+ CurrentMetrics::LocalThreadActive,
+ CurrentMetrics::LocalThreadScheduled,
+ max_threads,
+ max_threads,
+ not_exists_part.size()
+ );
auto s3 = data_disk->getObjectStorage();
if (!metadata_disk->exists(table_path))
@@ -75,25 +93,29 @@ void restoreMetaData(CustomStorageMergeTreePtr & storage,
const MergeTreeTable &
for (const auto & part : not_exists_part)
{
- auto part_path = table_path / part;
- auto metadata_file_path = part_path / "metadata.gluten";
+ auto job = [&]() {
+ auto part_path = table_path / part;
+ auto metadata_file_path = part_path / "metadata.gluten";
- if (metadata_disk->exists(part_path))
- continue;
- else
- metadata_disk->createDirectories(part_path);
- auto key =
s3->generateObjectKeyForPath(metadata_file_path.generic_string());
- StoredObject metadata_object(key.serialize());
- auto part_metadata =
extractPartMetaData(*s3->readObject(metadata_object));
- for (const auto & item : part_metadata)
- {
- auto item_path = part_path / item.first;
- auto out = metadata_disk->writeFile(item_path);
- out->write(item.second.data(), item.second.size());
- out->finalize();
- out->sync();
- }
+ if (metadata_disk->exists(part_path))
+ return;
+ else
+ metadata_disk->createDirectories(part_path);
+ auto key =
s3->generateObjectKeyForPath(metadata_file_path.generic_string());
+ StoredObject metadata_object(key.serialize());
+ auto part_metadata =
extractPartMetaData(*s3->readObject(metadata_object));
+ for (const auto & item : part_metadata)
+ {
+ auto item_path = part_path / item.first;
+ auto out = metadata_disk->writeFile(item_path);
+ out->write(item.second.data(), item.second.size());
+ out->finalize();
+ out->sync();
+ }
+ };
+ thread_pool.scheduleOrThrow(job);
}
+ thread_pool.wait();
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]