2012-07-17 20:04:39 +00:00
|
|
|
|
#include <DB/Storages/StorageMergeTree.h>
|
2014-03-13 12:48:07 +00:00
|
|
|
|
#include <DB/Storages/MergeTree/MergeTreeBlockOutputStream.h>
|
|
|
|
|
#include <DB/Storages/MergeTree/DiskSpaceMonitor.h>
|
|
|
|
|
#include <DB/Common/escapeForFileName.h>
|
2012-07-19 20:32:10 +00:00
|
|
|
|
|
2012-07-17 20:04:39 +00:00
|
|
|
|
namespace DB
|
|
|
|
|
{
|
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
BackgroundProcessingPool StorageMergeTree::merge_pool;
|
|
|
|
|
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
2014-05-08 07:12:01 +00:00
|
|
|
|
StorageMergeTree::StorageMergeTree(const String & path_, const String & database_name_, const String & name_,
|
|
|
|
|
NamesAndTypesListPtr columns_,
|
2014-03-09 17:36:01 +00:00
|
|
|
|
const Context & context_,
|
|
|
|
|
ASTPtr & primary_expr_ast_,
|
|
|
|
|
const String & date_column_name_,
|
2014-04-08 07:58:53 +00:00
|
|
|
|
const ASTPtr & sampling_expression_, /// nullptr, если семплирование не поддерживается.
|
2014-03-09 17:36:01 +00:00
|
|
|
|
size_t index_granularity_,
|
|
|
|
|
MergeTreeData::Mode mode_,
|
|
|
|
|
const String & sign_column_,
|
|
|
|
|
const MergeTreeSettings & settings_)
|
2014-03-13 12:48:07 +00:00
|
|
|
|
: path(path_), name(name_), full_path(path + escapeForFileName(name) + '/'), increment(full_path + "increment.txt"),
|
2014-05-08 18:34:43 +00:00
|
|
|
|
data(full_path, columns_, context_, primary_expr_ast_, date_column_name_, sampling_expression_,
|
|
|
|
|
index_granularity_,mode_, sign_column_, settings_, database_name_ + "." + name),
|
2014-03-13 12:48:07 +00:00
|
|
|
|
reader(data), writer(data), merger(data),
|
2014-05-08 07:12:01 +00:00
|
|
|
|
log(&Logger::get(database_name_ + "." + name + " (StorageMergeTree)")),
|
2014-03-13 12:48:07 +00:00
|
|
|
|
shutdown_called(false)
|
|
|
|
|
{
|
|
|
|
|
increment.fixIfBroken(data.getMaxDataPartIndex());
|
2014-04-09 15:52:47 +00:00
|
|
|
|
|
|
|
|
|
data.clearOldParts();
|
2014-03-13 12:48:07 +00:00
|
|
|
|
}
|
2012-07-17 20:04:39 +00:00
|
|
|
|
|
2013-02-06 11:26:35 +00:00
|
|
|
|
StoragePtr StorageMergeTree::create(
|
2014-05-08 07:12:01 +00:00
|
|
|
|
const String & path_, const String & database_name_, const String & name_,
|
|
|
|
|
NamesAndTypesListPtr columns_,
|
2013-05-05 18:02:05 +00:00
|
|
|
|
const Context & context_,
|
2013-02-06 11:26:35 +00:00
|
|
|
|
ASTPtr & primary_expr_ast_,
|
2014-03-09 17:36:01 +00:00
|
|
|
|
const String & date_column_name_,
|
|
|
|
|
const ASTPtr & sampling_expression_,
|
2013-02-06 11:26:35 +00:00
|
|
|
|
size_t index_granularity_,
|
2014-03-09 17:36:01 +00:00
|
|
|
|
MergeTreeData::Mode mode_,
|
2013-02-06 11:26:35 +00:00
|
|
|
|
const String & sign_column_,
|
2014-03-09 17:36:01 +00:00
|
|
|
|
const MergeTreeSettings & settings_)
|
2013-02-06 11:26:35 +00:00
|
|
|
|
{
|
2014-04-11 13:05:17 +00:00
|
|
|
|
StorageMergeTree * res = new StorageMergeTree(
|
2014-05-08 07:12:01 +00:00
|
|
|
|
path_, database_name_, name_, columns_, context_, primary_expr_ast_, date_column_name_,
|
2014-04-11 13:05:17 +00:00
|
|
|
|
sampling_expression_, index_granularity_, mode_, sign_column_, settings_);
|
|
|
|
|
StoragePtr res_ptr = res->thisPtr();
|
|
|
|
|
|
|
|
|
|
merge_pool.setNumberOfThreads(res->data.settings.merging_threads);
|
|
|
|
|
merge_pool.setSleepTime(5);
|
|
|
|
|
res->merge_task_handle = merge_pool.addTask(std::bind(&StorageMergeTree::mergeTask, res, std::placeholders::_1));
|
|
|
|
|
|
|
|
|
|
return res_ptr;
|
2013-02-06 11:26:35 +00:00
|
|
|
|
}
|
|
|
|
|
|
2013-09-30 01:29:19 +00:00
|
|
|
|
void StorageMergeTree::shutdown()
|
2012-07-30 20:32:36 +00:00
|
|
|
|
{
|
2014-03-13 12:48:07 +00:00
|
|
|
|
if (shutdown_called)
|
|
|
|
|
return;
|
|
|
|
|
shutdown_called = true;
|
|
|
|
|
merger.cancelAll();
|
2014-04-11 13:05:17 +00:00
|
|
|
|
merge_pool.removeTask(merge_task_handle);
|
2012-07-18 19:44:04 +00:00
|
|
|
|
}
|
|
|
|
|
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
|
|
|
|
StorageMergeTree::~StorageMergeTree()
|
|
|
|
|
{
|
|
|
|
|
shutdown();
|
|
|
|
|
}
|
2012-07-18 19:44:04 +00:00
|
|
|
|
|
2012-07-21 05:07:14 +00:00
|
|
|
|
BlockInputStreams StorageMergeTree::read(
|
2014-03-09 17:36:01 +00:00
|
|
|
|
const Names & column_names,
|
2012-07-21 05:07:14 +00:00
|
|
|
|
ASTPtr query,
|
2013-02-01 19:02:04 +00:00
|
|
|
|
const Settings & settings,
|
2012-07-21 05:07:14 +00:00
|
|
|
|
QueryProcessingStage::Enum & processed_stage,
|
|
|
|
|
size_t max_block_size,
|
|
|
|
|
unsigned threads)
|
|
|
|
|
{
|
2014-03-13 12:48:07 +00:00
|
|
|
|
return reader.read(column_names, query, settings, processed_stage, max_block_size, threads);
|
2012-12-06 09:45:09 +00:00
|
|
|
|
}
|
|
|
|
|
|
2014-03-09 17:36:01 +00:00
|
|
|
|
BlockOutputStreamPtr StorageMergeTree::write(ASTPtr query)
|
2013-01-23 11:16:32 +00:00
|
|
|
|
{
|
2014-03-19 10:45:13 +00:00
|
|
|
|
return new MergeTreeBlockOutputStream(*this);
|
2013-01-23 11:16:32 +00:00
|
|
|
|
}
|
|
|
|
|
|
2014-03-20 13:28:49 +00:00
|
|
|
|
void StorageMergeTree::drop()
|
2012-08-16 18:17:01 +00:00
|
|
|
|
{
|
2014-04-11 15:53:32 +00:00
|
|
|
|
shutdown();
|
2014-03-13 12:48:07 +00:00
|
|
|
|
data.dropAllData();
|
2013-08-07 13:07:42 +00:00
|
|
|
|
}
|
|
|
|
|
|
2014-03-09 17:36:01 +00:00
|
|
|
|
void StorageMergeTree::rename(const String & new_path_to_db, const String & new_name)
|
2014-03-04 11:30:50 +00:00
|
|
|
|
{
|
2014-03-13 12:48:07 +00:00
|
|
|
|
std::string new_full_path = new_path_to_db + escapeForFileName(new_name) + '/';
|
2014-03-13 19:14:25 +00:00
|
|
|
|
|
|
|
|
|
data.setPath(new_full_path);
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
|
|
|
|
path = new_path_to_db;
|
|
|
|
|
name = new_name;
|
|
|
|
|
full_path = new_full_path;
|
|
|
|
|
|
|
|
|
|
increment.setPath(full_path + "increment.txt");
|
2014-05-08 07:12:01 +00:00
|
|
|
|
|
|
|
|
|
/// TODO: Можно обновить названия логгеров у this, data, reader, writer, merger.
|
2014-03-04 11:30:50 +00:00
|
|
|
|
}
|
|
|
|
|
|
2013-08-09 00:12:59 +00:00
|
|
|
|
void StorageMergeTree::alter(const ASTAlterQuery::Parameters & params)
|
2013-08-07 13:07:42 +00:00
|
|
|
|
{
|
2014-03-09 17:36:01 +00:00
|
|
|
|
data.alter(params);
|
2013-10-03 12:46:17 +00:00
|
|
|
|
}
|
|
|
|
|
|
2014-03-20 13:00:42 +00:00
|
|
|
|
void StorageMergeTree::prepareAlterModify(const ASTAlterQuery::Parameters & params)
|
|
|
|
|
{
|
|
|
|
|
data.prepareAlterModify(params);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void StorageMergeTree::commitAlterModify(const ASTAlterQuery::Parameters & params)
|
|
|
|
|
{
|
|
|
|
|
data.commitAlterModify(params);
|
|
|
|
|
}
|
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
bool StorageMergeTree::merge(bool aggressive, BackgroundProcessingPool::Context * pool_context)
|
2014-03-13 12:48:07 +00:00
|
|
|
|
{
|
2014-04-11 13:05:17 +00:00
|
|
|
|
auto structure_lock = lockStructure(false);
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
/// Удаляем старые куски.
|
|
|
|
|
data.clearOldParts();
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
size_t disk_space = DiskSpaceMonitor::getUnreservedFreeSpace(full_path);
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
/// Нужно вызывать деструктор под незалоченным currently_merging_mutex.
|
|
|
|
|
CurrentlyMergingPartsTaggerPtr merging_tagger;
|
|
|
|
|
String merged_name;
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
|
|
|
|
{
|
2014-04-11 13:05:17 +00:00
|
|
|
|
Poco::ScopedLock<Poco::FastMutex> lock(currently_merging_mutex);
|
2014-03-27 11:30:54 +00:00
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
MergeTreeData::DataPartsVector parts;
|
|
|
|
|
auto can_merge = std::bind(&StorageMergeTree::canMergeParts, this, std::placeholders::_1, std::placeholders::_2);
|
|
|
|
|
/// Если слияние запущено из пула потоков, и хотя бы половина потоков сливает большие куски,
|
|
|
|
|
/// не будем сливать большие куски.
|
|
|
|
|
int big_merges = merge_pool.getCounter("big merges");
|
|
|
|
|
bool only_small = pool_context && big_merges * 2 >= merge_pool.getNumberOfThreads();
|
|
|
|
|
|
|
|
|
|
if (!merger.selectPartsToMerge(parts, merged_name, disk_space, false, aggressive, only_small, can_merge) &&
|
|
|
|
|
!merger.selectPartsToMerge(parts, merged_name, disk_space, true, aggressive, only_small, can_merge))
|
|
|
|
|
{
|
|
|
|
|
return false;
|
|
|
|
|
}
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
merging_tagger = new CurrentlyMergingPartsTagger(parts, merger.estimateDiskSpaceForMerge(parts), *this);
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
/// Если собираемся сливать большие куски, увеличим счетчик потоков, сливающих большие куски.
|
|
|
|
|
if (pool_context)
|
|
|
|
|
{
|
|
|
|
|
for (const auto & part : parts)
|
2014-03-13 12:48:07 +00:00
|
|
|
|
{
|
2014-04-11 13:05:17 +00:00
|
|
|
|
if (part->size * data.index_granularity > 25 * 1024 * 1024)
|
2014-03-13 19:07:17 +00:00
|
|
|
|
{
|
2014-04-11 13:05:17 +00:00
|
|
|
|
pool_context->incrementCounter("big merges");
|
|
|
|
|
break;
|
2014-03-13 19:07:17 +00:00
|
|
|
|
}
|
2014-03-13 12:48:07 +00:00
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
merger.mergeParts(merging_tagger->parts, merged_name);
|
2014-03-13 12:48:07 +00:00
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
return true;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
bool StorageMergeTree::mergeTask(BackgroundProcessingPool::Context & context)
|
2014-03-13 12:48:07 +00:00
|
|
|
|
{
|
2014-04-11 13:05:17 +00:00
|
|
|
|
if (shutdown_called)
|
|
|
|
|
return false;
|
2014-04-11 18:04:21 +00:00
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
return merge(false, &context);
|
|
|
|
|
}
|
|
|
|
|
catch (Exception & e)
|
|
|
|
|
{
|
|
|
|
|
if (e.code() == ErrorCodes::ABORTED)
|
|
|
|
|
{
|
|
|
|
|
LOG_INFO(log, "Merge cancelled");
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
throw;
|
|
|
|
|
}
|
2014-03-13 12:48:07 +00:00
|
|
|
|
}
|
|
|
|
|
|
2014-04-11 13:05:17 +00:00
|
|
|
|
|
2014-03-13 17:44:00 +00:00
|
|
|
|
bool StorageMergeTree::canMergeParts(const MergeTreeData::DataPartPtr & left, const MergeTreeData::DataPartPtr & right)
|
|
|
|
|
{
|
|
|
|
|
return !currently_merging.count(left) && !currently_merging.count(right);
|
|
|
|
|
}
|
|
|
|
|
|
2012-07-17 20:04:39 +00:00
|
|
|
|
}
|