2012-07-18 19:16:16 +00:00
# include <boost/bind.hpp>
2012-11-28 08:52:15 +00:00
# include <numeric>
2012-07-17 20:04:39 +00:00
2012-07-19 20:32:10 +00:00
# include <Poco/DirectoryIterator.h>
# include <Poco/NumberParser.h>
2013-04-08 19:47:25 +00:00
# include <Poco/Ext/ScopedTry.h>
2012-07-19 20:32:10 +00:00
# include <Yandex/time2str.h>
2012-07-17 20:04:39 +00:00
# include <DB/Common/escapeForFileName.h>
# include <DB/IO/WriteBufferFromString.h>
# include <DB/IO/WriteBufferFromFile.h>
# include <DB/IO/CompressedWriteBuffer.h>
2012-07-21 03:45:48 +00:00
# include <DB/IO/ReadBufferFromString.h>
2012-07-19 20:32:10 +00:00
# include <DB/IO/ReadBufferFromFile.h>
# include <DB/IO/CompressedReadBuffer.h>
2012-07-17 20:04:39 +00:00
# include <DB/Columns/ColumnsNumber.h>
2012-08-30 17:43:31 +00:00
# include <DB/Columns/ColumnArray.h>
2012-07-17 20:04:39 +00:00
2012-07-21 03:45:48 +00:00
# include <DB/DataTypes/DataTypesNumberFixed.h>
2012-08-30 17:43:31 +00:00
# include <DB/DataTypes/DataTypeArray.h>
2012-07-21 03:45:48 +00:00
2012-07-19 20:32:10 +00:00
# include <DB/DataStreams/IProfilingBlockInputStream.h>
2012-07-30 20:32:36 +00:00
# include <DB/DataStreams/MergingSortedBlockInputStream.h>
2012-08-16 17:27:40 +00:00
# include <DB/DataStreams/CollapsingSortedBlockInputStream.h>
2012-07-31 17:22:40 +00:00
# include <DB/DataStreams/ExpressionBlockInputStream.h>
2012-11-28 08:52:15 +00:00
# include <DB/DataStreams/ConcatBlockInputStream.h>
2012-07-21 06:47:17 +00:00
# include <DB/DataStreams/narrowBlockInputStreams.h>
2012-07-30 20:32:36 +00:00
# include <DB/DataStreams/copyData.h>
2012-12-12 14:25:55 +00:00
# include <DB/DataStreams/FilterBlockInputStream.h>
2012-07-19 20:32:10 +00:00
2012-07-21 03:45:48 +00:00
# include <DB/Parsers/ASTExpressionList.h>
# include <DB/Parsers/ASTSelectQuery.h>
# include <DB/Parsers/ASTFunction.h>
# include <DB/Parsers/ASTLiteral.h>
2012-12-12 14:25:55 +00:00
# include <DB/Parsers/ASTIdentifier.h>
2012-07-21 03:45:48 +00:00
2012-07-17 20:04:39 +00:00
# include <DB/Interpreters/sortBlock.h>
# include <DB/Storages/StorageMergeTree.h>
2013-04-24 10:31:32 +00:00
# include <DB/Storages/MergeTree/PKCondition.h>
# include <DB/Storages/MergeTree/MergeTreeBlockOutputStream.h>
# include <DB/Storages/MergeTree/MergeTreeBlockInputStream.h>
# include <DB/Storages/MergeTree/MergedBlockOutputStream.h>
2012-07-19 20:32:10 +00:00
2012-07-17 20:04:39 +00:00
namespace DB
{
StorageMergeTree : : StorageMergeTree (
const String & path_ , const String & name_ , NamesAndTypesListPtr columns_ ,
Context & context_ ,
2012-12-12 14:25:55 +00:00
ASTPtr & primary_expr_ast_ ,
2012-12-12 15:45:08 +00:00
const String & date_column_name_ , const ASTPtr & sampling_expression_ ,
2012-08-16 17:27:40 +00:00
size_t index_granularity_ ,
2012-08-20 05:32:50 +00:00
const String & sign_column_ ,
2012-08-29 20:23:19 +00:00
const StorageMergeTreeSettings & settings_ )
2012-07-17 20:04:39 +00:00
: path ( path_ ) , name ( name_ ) , full_path ( path + escapeForFileName ( name ) + ' / ' ) , columns ( columns_ ) ,
context ( context_ ) , primary_expr_ast ( primary_expr_ast_ - > clone ( ) ) ,
2012-12-12 15:45:08 +00:00
date_column_name ( date_column_name_ ) , sampling_expression ( sampling_expression_ ) ,
2012-12-12 14:25:55 +00:00
index_granularity ( index_granularity_ ) ,
2012-08-20 05:32:50 +00:00
sign_column ( sign_column_ ) ,
2012-08-29 20:23:19 +00:00
settings ( settings_ ) ,
2012-07-19 20:32:10 +00:00
increment ( full_path + " increment.txt " ) , log ( & Logger : : get ( " StorageMergeTree: " + name ) )
2012-07-17 20:04:39 +00:00
{
2012-12-06 13:07:29 +00:00
min_marks_for_seek = ( settings . min_rows_for_seek + index_granularity - 1 ) / index_granularity ;
min_marks_for_concurrent_read = ( settings . min_rows_for_concurrent_read + index_granularity - 1 ) / index_granularity ;
2012-07-17 20:04:39 +00:00
/// создаём директорию, если её нет
Poco : : File ( full_path ) . createDirectories ( ) ;
/// инициализируем описание сортировки
sort_descr . reserve ( primary_expr_ast - > children . size ( ) ) ;
for ( ASTs : : iterator it = primary_expr_ast - > children . begin ( ) ;
it ! = primary_expr_ast - > children . end ( ) ;
+ + it )
{
2012-07-18 20:14:41 +00:00
String name = ( * it ) - > getColumnName ( ) ;
2012-07-17 20:04:39 +00:00
sort_descr . push_back ( SortColumnDescription ( name , 1 ) ) ;
}
2012-07-18 20:14:41 +00:00
2012-08-02 17:33:31 +00:00
context . setColumns ( * columns ) ;
2012-07-18 20:14:41 +00:00
primary_expr = new Expression ( primary_expr_ast , context ) ;
2012-07-21 07:02:55 +00:00
primary_key_sample = primary_expr - > getSampleBlock ( ) ;
2012-07-19 20:32:10 +00:00
2012-11-28 08:52:15 +00:00
merge_threads = new boost : : threadpool : : pool ( settings . merging_threads ) ;
2012-07-19 20:32:10 +00:00
loadDataParts ( ) ;
2012-07-17 20:04:39 +00:00
}
2013-02-06 11:26:35 +00:00
StoragePtr StorageMergeTree : : create (
const String & path_ , const String & name_ , NamesAndTypesListPtr columns_ ,
Context & context_ ,
ASTPtr & primary_expr_ast_ ,
const String & date_column_name_ , const ASTPtr & sampling_expression_ ,
size_t index_granularity_ ,
const String & sign_column_ ,
const StorageMergeTreeSettings & settings_ )
{
return ( new StorageMergeTree ( path_ , name_ , columns_ , context_ , primary_expr_ast_ , date_column_name_ , sampling_expression_ , index_granularity_ , sign_column_ , settings_ ) ) - > thisPtr ( ) ;
}
2012-07-17 20:04:39 +00:00
2012-07-30 20:32:36 +00:00
StorageMergeTree : : ~ StorageMergeTree ( )
{
2012-11-28 08:52:15 +00:00
joinMergeThreads ( ) ;
2012-07-30 20:32:36 +00:00
}
2012-07-18 19:44:04 +00:00
BlockOutputStreamPtr StorageMergeTree : : write ( ASTPtr query )
{
2013-01-23 17:38:03 +00:00
return new MergeTreeBlockOutputStream ( thisPtr ( ) ) ;
2012-07-18 19:44:04 +00:00
}
2012-07-21 05:07:14 +00:00
BlockInputStreams StorageMergeTree : : read (
2012-12-12 16:11:27 +00:00
const Names & column_names_to_return ,
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 )
{
2012-12-15 02:11:37 +00:00
check ( column_names_to_return ) ;
processed_stage = QueryProcessingStage : : FetchColumns ;
2012-12-10 10:23:10 +00:00
PKCondition key_condition ( query , context , sort_descr ) ;
PKCondition date_condition ( query , context , SortDescription ( 1 , SortColumnDescription ( date_column_name , 1 ) ) ) ;
2012-07-21 03:45:48 +00:00
2012-12-13 10:08:54 +00:00
typedef std : : vector < DataPartPtr > PartsList ;
PartsList parts ;
2013-01-09 10:19:24 +00:00
/// Выберем куски, в которых могут быть данные, удовлетворяющие date_condition.
2012-12-13 10:08:54 +00:00
{
Poco : : ScopedLock < Poco : : FastMutex > lock ( data_parts_mutex ) ;
for ( DataParts : : iterator it = data_parts . begin ( ) ; it ! = data_parts . end ( ) ; + + it )
if ( date_condition . mayBeTrueInRange ( Row ( 1 , static_cast < UInt64 > ( ( * it ) - > left_date ) ) , Row ( 1 , static_cast < UInt64 > ( ( * it ) - > right_date ) ) ) )
parts . push_back ( * it ) ;
}
2012-12-12 16:26:18 +00:00
/// Семплирование.
2012-12-12 16:11:27 +00:00
Names column_names_to_read = column_names_to_return ;
2012-12-13 09:32:08 +00:00
UInt64 sampling_column_value_limit = 0 ;
2012-12-12 16:26:18 +00:00
typedef Poco : : SharedPtr < ASTFunction > ASTFunctionPtr ;
ASTFunctionPtr filter_function ;
ExpressionPtr filter_expression ;
2012-12-12 14:25:55 +00:00
ASTSelectQuery & select = * dynamic_cast < ASTSelectQuery * > ( & * query ) ;
if ( select . sample_size )
{
2013-01-05 20:03:19 +00:00
double size = apply_visitor ( FieldVisitorConvertToNumber < double > ( ) ,
2012-12-13 10:08:54 +00:00
dynamic_cast < ASTLiteral & > ( * select . sample_size ) . value ) ;
2012-12-12 14:25:55 +00:00
if ( size < 0 )
throw Exception ( " Negative sample size " , ErrorCodes : : ARGUMENT_OUT_OF_BOUND ) ;
if ( size > 1 )
{
2013-01-05 20:03:19 +00:00
size_t requested_count = apply_visitor ( FieldVisitorConvertToNumber < UInt64 > ( ) , dynamic_cast < ASTLiteral & > ( * select . sample_size ) . value ) ;
2012-12-13 10:08:54 +00:00
/// Узнаем, сколько строк мы бы прочли без семплирования.
LOG_DEBUG ( log , " Preliminary index scan with condition: " < < key_condition . toString ( ) ) ;
2012-12-13 09:32:08 +00:00
size_t total_count = 0 ;
2012-12-13 10:08:54 +00:00
for ( size_t i = 0 ; i < parts . size ( ) ; + + i )
2012-12-13 09:32:08 +00:00
{
2012-12-13 10:08:54 +00:00
DataPartPtr & part = parts [ i ] ;
MarkRanges ranges = MergeTreeBlockInputStream : : markRangesFromPkRange ( full_path + part - > name + ' / ' ,
part - > size ,
* this ,
key_condition ) ;
for ( size_t j = 0 ; j < ranges . size ( ) ; + + j )
total_count + = ranges [ j ] . end - ranges [ j ] . begin ;
2012-12-13 09:32:08 +00:00
}
total_count * = index_granularity ;
size = std : : min ( 1. , static_cast < double > ( requested_count ) / total_count ) ;
2012-12-13 09:38:36 +00:00
2012-12-13 10:08:54 +00:00
LOG_DEBUG ( log , " Selected relative sample size: " < < size ) ;
2012-12-12 14:25:55 +00:00
}
2012-12-13 09:32:08 +00:00
UInt64 sampling_column_max = 0 ;
DataTypePtr type = Expression ( sampling_expression , context ) . getReturnTypes ( ) [ 0 ] ;
if ( type - > getName ( ) = = " UInt64 " )
sampling_column_max = std : : numeric_limits < UInt64 > : : max ( ) ;
else if ( type - > getName ( ) = = " UInt32 " )
sampling_column_max = std : : numeric_limits < UInt32 > : : max ( ) ;
else if ( type - > getName ( ) = = " UInt16 " )
sampling_column_max = std : : numeric_limits < UInt16 > : : max ( ) ;
else if ( type - > getName ( ) = = " UInt8 " )
sampling_column_max = std : : numeric_limits < UInt8 > : : max ( ) ;
2012-12-12 14:25:55 +00:00
else
2012-12-13 09:32:08 +00:00
throw Exception ( " Invalid sampling column type in storage parameters: " + type - > getName ( ) + " . Must be unsigned integer type. " , ErrorCodes : : ILLEGAL_TYPE_OF_COLUMN_FOR_FILTER ) ;
2012-12-13 10:08:54 +00:00
/// Добавим условие, чтобы отсечь еще что-нибудь при повторном просмотре индекса.
2012-12-13 09:32:08 +00:00
sampling_column_value_limit = static_cast < UInt64 > ( size * sampling_column_max ) ;
if ( ! key_condition . addCondition ( sampling_expression - > getColumnName ( ) ,
Range : : RightBounded ( sampling_column_value_limit , true ) ) )
throw Exception ( " Sampling column not in primary key " , ErrorCodes : : ILLEGAL_COLUMN ) ;
/// Выражение для фильтрации: sampling_expression <= sampling_column_value_limit
2012-12-12 16:26:18 +00:00
ASTPtr filter_function_args = new ASTExpressionList ;
filter_function_args - > children . push_back ( sampling_expression ) ;
2012-12-13 09:32:08 +00:00
filter_function_args - > children . push_back ( new ASTLiteral ( StringRange ( ) , sampling_column_value_limit ) ) ;
2012-12-12 16:26:18 +00:00
filter_function = new ASTFunction ;
filter_function - > name = " lessOrEquals " ;
filter_function - > arguments = filter_function_args ;
filter_function - > children . push_back ( filter_function - > arguments ) ;
filter_expression = new Expression ( filter_function , context ) ;
std : : vector < String > add_columns = filter_expression - > getRequiredColumns ( ) ;
column_names_to_read . insert ( column_names_to_read . end ( ) , add_columns . begin ( ) , add_columns . end ( ) ) ;
std : : sort ( column_names_to_read . begin ( ) , column_names_to_read . end ( ) ) ;
column_names_to_read . erase ( std : : unique ( column_names_to_read . begin ( ) , column_names_to_read . end ( ) ) , column_names_to_read . end ( ) ) ;
}
2012-12-12 14:25:55 +00:00
2012-12-13 10:08:54 +00:00
LOG_DEBUG ( log , " Key condition: " < < key_condition . toString ( ) ) ;
LOG_DEBUG ( log , " Date condition: " < < date_condition . toString ( ) ) ;
2012-11-28 08:52:15 +00:00
2012-12-06 11:48:41 +00:00
RangesInDataParts parts_with_ranges ;
2012-11-28 08:52:15 +00:00
/// Найдем, какой диапазон читать из каждого куска.
size_t sum_marks = 0 ;
2012-12-06 11:48:41 +00:00
size_t sum_ranges = 0 ;
2012-11-28 08:52:15 +00:00
for ( size_t i = 0 ; i < parts . size ( ) ; + + i )
{
2012-12-06 11:48:41 +00:00
DataPartPtr & part = parts [ i ] ;
RangesInDataPart ranges ( part ) ;
ranges . ranges = MergeTreeBlockInputStream : : markRangesFromPkRange ( full_path + part - > name + ' / ' ,
2012-12-13 10:08:54 +00:00
part - > size ,
* this ,
key_condition ) ;
2012-12-06 11:48:41 +00:00
if ( ! ranges . ranges . empty ( ) )
{
parts_with_ranges . push_back ( ranges ) ;
sum_ranges + = ranges . ranges . size ( ) ;
for ( size_t j = 0 ; j < ranges . ranges . size ( ) ; + + j )
{
sum_marks + = ranges . ranges [ j ] . end - ranges . ranges [ j ] . begin ;
}
}
2012-11-28 08:52:15 +00:00
}
2012-12-06 12:24:08 +00:00
LOG_DEBUG ( log , " Selected " < < parts . size ( ) < < " parts by date, " < < parts_with_ranges . size ( ) < < " parts by key, "
2012-12-06 11:48:41 +00:00
< < sum_marks < < " marks to read from " < < sum_ranges < < " ranges " ) ;
2012-11-28 08:52:15 +00:00
2013-04-23 11:08:41 +00:00
BlockInputStreams res ;
if ( select . final )
{
2013-04-24 11:39:12 +00:00
std : : vector < String > add_columns = primary_expr - > getRequiredColumns ( ) ;
column_names_to_read . insert ( column_names_to_read . end ( ) , add_columns . begin ( ) , add_columns . end ( ) ) ;
std : : sort ( column_names_to_read . begin ( ) , column_names_to_read . end ( ) ) ;
column_names_to_read . erase ( std : : unique ( column_names_to_read . begin ( ) , column_names_to_read . end ( ) ) , column_names_to_read . end ( ) ) ;
2013-04-23 11:08:41 +00:00
res = spreadMarkRangesAmongThreadsCollapsing ( parts_with_ranges , threads , column_names_to_read , max_block_size ) ;
}
else
{
res = spreadMarkRangesAmongThreads ( parts_with_ranges , threads , column_names_to_read , max_block_size ) ;
}
2012-12-12 14:25:55 +00:00
2012-12-13 09:32:08 +00:00
if ( select . sample_size )
2012-12-12 14:25:55 +00:00
{
for ( size_t i = 0 ; i < res . size ( ) ; + + i )
{
BlockInputStreamPtr original_stream = res [ i ] ;
BlockInputStreamPtr expression_stream = new ExpressionBlockInputStream ( original_stream , filter_expression ) ;
BlockInputStreamPtr filter_stream = new FilterBlockInputStream ( expression_stream , filter_function - > getColumnName ( ) ) ;
res [ i ] = filter_stream ;
}
}
return res ;
2012-12-06 09:45:09 +00:00
}
/// Примерно поровну распределить засечки между потоками.
BlockInputStreams StorageMergeTree : : spreadMarkRangesAmongThreads ( RangesInDataParts parts , size_t threads , const Names & column_names , size_t max_block_size )
{
/// Н а всякий случай перемешаем куски.
std : : random_shuffle ( parts . begin ( ) , parts . end ( ) ) ;
/// Посчитаем засечки для каждого куска.
std : : vector < size_t > sum_marks_in_parts ( parts . size ( ) ) ;
size_t sum_marks = 0 ;
for ( size_t i = 0 ; i < parts . size ( ) ; + + i )
{
/// Пусть отрезки будут перечислены справа налево, чтобы можно было выбрасывать самый левый отрезок с помощью pop_back().
std : : reverse ( parts [ i ] . ranges . begin ( ) , parts [ i ] . ranges . end ( ) ) ;
sum_marks_in_parts [ i ] = 0 ;
for ( size_t j = 0 ; j < parts [ i ] . ranges . size ( ) ; + + j )
{
MarkRange & range = parts [ i ] . ranges [ j ] ;
sum_marks_in_parts [ i ] + = range . end - range . begin ;
}
sum_marks + = sum_marks_in_parts [ i ] ;
}
2012-11-28 08:52:15 +00:00
BlockInputStreams res ;
2012-11-30 00:52:45 +00:00
2012-11-30 08:33:36 +00:00
if ( sum_marks > 0 )
2012-11-28 08:52:15 +00:00
{
2012-12-06 09:45:09 +00:00
size_t min_marks_per_thread = ( sum_marks - 1 ) / threads + 1 ;
2012-11-30 08:33:36 +00:00
2012-12-06 09:45:09 +00:00
for ( size_t i = 0 ; i < threads & & ! parts . empty ( ) ; + + i )
2012-11-28 08:52:15 +00:00
{
2012-12-06 09:45:09 +00:00
size_t need_marks = min_marks_per_thread ;
2012-11-30 08:33:36 +00:00
BlockInputStreams streams ;
2012-11-30 00:52:45 +00:00
2012-12-06 09:45:09 +00:00
/// Цикл по кускам.
while ( need_marks > 0 & & ! parts . empty ( ) )
2012-11-28 08:52:15 +00:00
{
2012-12-06 09:45:09 +00:00
RangesInDataPart & part = parts . back ( ) ;
size_t & marks_in_part = sum_marks_in_parts . back ( ) ;
/// Н е будем брать из куска слишком мало строк.
2012-12-06 13:07:29 +00:00
if ( marks_in_part > = min_marks_for_concurrent_read & &
need_marks < min_marks_for_concurrent_read )
need_marks = min_marks_for_concurrent_read ;
2012-11-30 08:33:36 +00:00
2012-12-06 09:45:09 +00:00
/// Н е будем оставлять в куске слишком мало строк.
if ( marks_in_part > need_marks & &
2012-12-06 13:07:29 +00:00
marks_in_part - need_marks < min_marks_for_concurrent_read )
2012-12-06 09:45:09 +00:00
need_marks = marks_in_part ;
2012-11-30 08:33:36 +00:00
2012-12-06 09:45:09 +00:00
/// Возьмем весь кусок, если он достаточно мал.
if ( marks_in_part < = need_marks )
2012-11-30 08:33:36 +00:00
{
2012-12-06 09:45:09 +00:00
/// Восстановим порядок отрезков.
std : : reverse ( part . ranges . begin ( ) , part . ranges . end ( ) ) ;
streams . push_back ( new MergeTreeBlockInputStream ( full_path + part . data_part - > name + ' / ' ,
max_block_size , column_names , * this ,
2013-01-23 17:38:03 +00:00
part . data_part , part . ranges , thisPtr ( ) ) ) ;
2012-12-06 09:45:09 +00:00
need_marks - = marks_in_part ;
parts . pop_back ( ) ;
sum_marks_in_parts . pop_back ( ) ;
2012-11-30 08:33:36 +00:00
continue ;
}
2012-12-06 09:45:09 +00:00
MarkRanges ranges_to_get_from_part ;
2012-11-30 08:33:36 +00:00
2012-12-06 09:45:09 +00:00
/// Цикл по отрезкам куска.
while ( need_marks > 0 )
{
if ( part . ranges . empty ( ) )
throw Exception ( " Unexpected end of ranges while spreading marks among threads " , ErrorCodes : : LOGICAL_ERROR ) ;
MarkRange & range = part . ranges . back ( ) ;
size_t marks_in_range = range . end - range . begin ;
2012-11-30 08:33:36 +00:00
2012-12-06 09:45:09 +00:00
size_t marks_to_get_from_range = std : : min ( marks_in_range , need_marks ) ;
ranges_to_get_from_part . push_back ( MarkRange ( range . begin , range . begin + marks_to_get_from_range ) ) ;
range . begin + = marks_to_get_from_range ;
marks_in_part - = marks_to_get_from_range ;
need_marks - = marks_to_get_from_range ;
if ( range . begin = = range . end )
part . ranges . pop_back ( ) ;
}
2012-11-30 08:33:36 +00:00
2012-12-06 09:45:09 +00:00
streams . push_back ( new MergeTreeBlockInputStream ( full_path + part . data_part - > name + ' / ' ,
max_block_size , column_names , * this ,
2013-01-23 17:38:03 +00:00
part . data_part , ranges_to_get_from_part , thisPtr ( ) ) ) ;
2012-11-28 08:52:15 +00:00
}
2012-11-30 00:52:45 +00:00
2012-11-30 08:33:36 +00:00
if ( streams . size ( ) = = 1 )
res . push_back ( streams [ 0 ] ) ;
2012-11-29 08:41:20 +00:00
else
2012-11-30 08:33:36 +00:00
res . push_back ( new ConcatBlockInputStream ( streams ) ) ;
2012-11-28 08:52:15 +00:00
}
2012-12-06 09:45:09 +00:00
if ( ! parts . empty ( ) )
2012-12-03 09:39:04 +00:00
throw Exception ( " Couldn't spread marks among threads " , ErrorCodes : : LOGICAL_ERROR ) ;
2012-11-28 08:52:15 +00:00
}
2012-12-06 11:10:05 +00:00
return res ;
2012-07-21 03:45:48 +00:00
}
2013-04-24 11:39:12 +00:00
/// Распределить засечки между потоками и сделать, чтобы в ответе (почти) все данные были сколлапсированы (модификатор FINAL).
2013-04-23 11:08:41 +00:00
BlockInputStreams StorageMergeTree : : spreadMarkRangesAmongThreadsCollapsing ( RangesInDataParts parts , size_t threads , const Names & column_names , size_t max_block_size )
{
BlockInputStreams streams ;
for ( size_t part_index = 0 ; part_index < parts . size ( ) ; + + part_index )
{
RangesInDataPart & part = parts [ part_index ] ;
2013-04-23 11:35:09 +00:00
streams . push_back ( new ExpressionBlockInputStream ( new MergeTreeBlockInputStream ( full_path + part . data_part - > name + ' / ' ,
2013-04-23 11:08:41 +00:00
max_block_size , column_names , * this ,
2013-04-23 11:35:09 +00:00
part . data_part , part . ranges , thisPtr ( ) ) , primary_expr ) ) ;
2013-04-23 11:08:41 +00:00
}
return BlockInputStreams ( 1 , new CollapsingSortedBlockInputStream ( streams , sort_descr , sign_column , max_block_size ) ) ;
}
2012-07-19 20:32:10 +00:00
String StorageMergeTree : : getPartName ( Yandex : : DayNum_t left_date , Yandex : : DayNum_t right_date , UInt64 left_id , UInt64 right_id , UInt64 level )
2012-07-17 20:04:39 +00:00
{
Yandex : : DateLUTSingleton & date_lut = Yandex : : DateLUTSingleton : : instance ( ) ;
2012-07-19 20:32:10 +00:00
/// Имя директории для куска иммет вид: YYYYMMDD_YYYYMMDD_N_N_L.
2012-07-17 20:04:39 +00:00
String res ;
{
2012-07-19 20:32:10 +00:00
unsigned left_date_id = Yandex : : Date2OrderedIdentifier ( date_lut . fromDayNum ( left_date ) ) ;
unsigned right_date_id = Yandex : : Date2OrderedIdentifier ( date_lut . fromDayNum ( right_date ) ) ;
2012-07-17 20:04:39 +00:00
WriteBufferFromString wb ( res ) ;
2012-07-19 20:32:10 +00:00
writeIntText ( left_date_id , wb ) ;
2012-07-17 20:04:39 +00:00
writeChar ( ' _ ' , wb ) ;
2012-07-19 20:32:10 +00:00
writeIntText ( right_date_id , wb ) ;
2012-07-17 20:04:39 +00:00
writeChar ( ' _ ' , wb ) ;
writeIntText ( left_id , wb ) ;
writeChar ( ' _ ' , wb ) ;
writeIntText ( right_id , wb ) ;
writeChar ( ' _ ' , wb ) ;
writeIntText ( level , wb ) ;
}
return res ;
}
2012-07-19 20:32:10 +00:00
void StorageMergeTree : : loadDataParts ( )
{
LOG_DEBUG ( log , " Loading data parts " ) ;
2012-08-10 20:04:34 +00:00
Poco : : ScopedLock < Poco : : FastMutex > lock ( data_parts_mutex ) ;
Poco : : ScopedLock < Poco : : FastMutex > lock_all ( all_data_parts_mutex ) ;
2012-07-19 20:32:10 +00:00
Yandex : : DateLUTSingleton & date_lut = Yandex : : DateLUTSingleton : : instance ( ) ;
2012-08-10 20:04:34 +00:00
data_parts . clear ( ) ;
2012-07-19 20:32:10 +00:00
static Poco : : RegularExpression file_name_regexp ( " ^( \\ d{8})_( \\ d{8})_( \\ d+)_( \\ d+)_( \\ d+) " ) ;
Poco : : DirectoryIterator end ;
Poco : : RegularExpression : : MatchVec matches ;
for ( Poco : : DirectoryIterator it ( full_path ) ; it ! = end ; + + it )
{
std : : string file_name = it . name ( ) ;
if ( ! ( file_name_regexp . match ( file_name , 0 , matches ) & & 6 = = matches . size ( ) ) )
continue ;
2012-07-31 20:03:53 +00:00
DataPartPtr part = new DataPart ( * this ) ;
2012-07-23 06:23:29 +00:00
part - > left_date = date_lut . toDayNum ( Yandex : : OrderedIdentifier2Date ( file_name . substr ( matches [ 1 ] . offset , matches [ 1 ] . length ) ) ) ;
part - > right_date = date_lut . toDayNum ( Yandex : : OrderedIdentifier2Date ( file_name . substr ( matches [ 2 ] . offset , matches [ 2 ] . length ) ) ) ;
part - > left = Poco : : NumberParser : : parseUnsigned64 ( file_name . substr ( matches [ 3 ] . offset , matches [ 3 ] . length ) ) ;
part - > right = Poco : : NumberParser : : parseUnsigned64 ( file_name . substr ( matches [ 4 ] . offset , matches [ 4 ] . length ) ) ;
part - > level = Poco : : NumberParser : : parseUnsigned ( file_name . substr ( matches [ 5 ] . offset , matches [ 5 ] . length ) ) ;
part - > name = file_name ;
2012-07-19 20:32:10 +00:00
/// Размер - в количестве засечек.
2012-07-23 06:23:29 +00:00
part - > size = Poco : : File ( full_path + file_name + " / " + escapeForFileName ( columns - > front ( ) . first ) + " .mrk " ) . getSize ( )
2012-07-19 20:32:10 +00:00
/ MERGE_TREE_MARK_SIZE ;
2012-07-23 06:23:29 +00:00
part - > modification_time = it - > getLastModified ( ) . epochTime ( ) ;
2012-07-19 20:32:10 +00:00
2012-08-31 20:38:05 +00:00
part - > left_month = date_lut . toFirstDayNumOfMonth ( part - > left_date ) ;
part - > right_month = date_lut . toFirstDayNumOfMonth ( part - > right_date ) ;
2012-07-19 20:32:10 +00:00
2012-08-10 20:04:34 +00:00
data_parts . insert ( part ) ;
2012-07-19 20:32:10 +00:00
}
2012-08-10 20:04:34 +00:00
all_data_parts = data_parts ;
2012-07-31 20:03:53 +00:00
2012-07-31 20:13:14 +00:00
/** Удаляем из набора актуальных кусков куски, которые содержатся в другом куске (которые были склеены),
* н о п о к а к и м - т о п р и ч и н а м о с т а л и с ь л е ж а т ь в ф а й л о в о й с и с т е м е .
* У д а л е н и е ф а й л о в б у д е т п р о и з в е д е н о п о т о м в м е т о д е clearOldParts .
*/
2012-08-10 20:04:34 +00:00
if ( data_parts . size ( ) > = 2 )
2012-07-31 20:13:14 +00:00
{
2012-08-10 20:04:34 +00:00
DataParts : : iterator prev_jt = data_parts . begin ( ) ;
2012-08-07 20:37:45 +00:00
DataParts : : iterator curr_jt = prev_jt ;
+ + curr_jt ;
2012-08-10 20:04:34 +00:00
while ( curr_jt ! = data_parts . end ( ) )
2012-07-31 20:13:14 +00:00
{
2012-08-07 20:37:45 +00:00
/// Куски данных за разные месяцы рассматривать не будем
if ( ( * curr_jt ) - > left_month ! = ( * curr_jt ) - > right_month
| | ( * curr_jt ) - > right_month ! = ( * prev_jt ) - > left_month
| | ( * prev_jt ) - > left_month ! = ( * prev_jt ) - > right_month )
{
+ + prev_jt ;
+ + curr_jt ;
continue ;
}
2012-07-31 20:13:14 +00:00
2012-08-07 20:37:45 +00:00
if ( ( * curr_jt ) - > contains ( * * prev_jt ) )
{
LOG_WARNING ( log , " Part " < < ( * curr_jt ) - > name < < " contains " < < ( * prev_jt ) - > name ) ;
2012-08-10 20:04:34 +00:00
data_parts . erase ( prev_jt ) ;
2012-08-07 20:37:45 +00:00
prev_jt = curr_jt ;
+ + curr_jt ;
}
else if ( ( * prev_jt ) - > contains ( * * curr_jt ) )
{
LOG_WARNING ( log , " Part " < < ( * prev_jt ) - > name < < " contains " < < ( * curr_jt ) - > name ) ;
2012-08-10 20:04:34 +00:00
data_parts . erase ( curr_jt + + ) ;
2012-08-07 20:37:45 +00:00
}
else
{
+ + prev_jt ;
+ + curr_jt ;
}
2012-07-31 20:13:14 +00:00
}
}
2012-07-31 20:03:53 +00:00
2012-08-10 20:04:34 +00:00
LOG_DEBUG ( log , " Loaded data parts ( " < < data_parts . size ( ) < < " items) " ) ;
2012-07-19 20:32:10 +00:00
}
2012-07-23 06:23:29 +00:00
void StorageMergeTree : : clearOldParts ( )
{
Poco : : ScopedTry < Poco : : FastMutex > lock ;
/// Если метод уже вызван из другого потока (или если all_data_parts прямо сейчас меняют), то можно ничего не делать.
if ( ! lock . lock ( & all_data_parts_mutex ) )
{
LOG_TRACE ( log , " Already clearing or modifying old parts " ) ;
return ;
}
LOG_TRACE ( log , " Clearing old parts " ) ;
for ( DataParts : : iterator it = all_data_parts . begin ( ) ; it ! = all_data_parts . end ( ) ; )
{
2012-07-31 18:25:16 +00:00
int ref_count = it - > referenceCount ( ) ;
LOG_TRACE ( log , ( * it ) - > name < < " : ref_count = " < < ref_count ) ;
2012-08-10 20:04:34 +00:00
if ( ref_count = = 1 ) /// После этого ref_count не может увеличиться.
2012-07-23 06:23:29 +00:00
{
LOG_DEBUG ( log , " Removing part " < < ( * it ) - > name ) ;
( * it ) - > remove ( ) ;
all_data_parts . erase ( it + + ) ;
}
else
+ + it ;
}
}
2012-09-10 19:05:06 +00:00
void StorageMergeTree : : merge ( size_t iterations , bool async )
2012-07-23 06:23:29 +00:00
{
2012-11-28 08:52:15 +00:00
bool while_can = false ;
if ( iterations = = 0 ) {
while_can = true ;
iterations = settings . merging_threads ;
2012-08-01 20:08:59 +00:00
}
2012-11-28 17:17:17 +00:00
for ( size_t i = 0 ; i < iterations ; + + i )
2012-11-28 08:52:15 +00:00
merge_threads - > schedule ( boost : : bind ( & StorageMergeTree : : mergeThread , this , while_can ) ) ;
if ( ! async )
joinMergeThreads ( ) ;
2012-08-13 19:13:11 +00:00
}
2012-11-28 08:52:15 +00:00
void StorageMergeTree : : mergeThread ( bool while_can )
2012-08-13 19:13:11 +00:00
{
try
2012-07-23 06:23:29 +00:00
{
2012-11-28 08:52:15 +00:00
std : : vector < DataPartPtr > parts ;
2013-02-14 11:22:56 +00:00
while ( selectPartsToMerge ( parts , false ) | |
selectPartsToMerge ( parts , true ) )
2012-08-13 19:13:11 +00:00
{
2012-11-28 08:52:15 +00:00
mergeParts ( parts ) ;
2012-08-13 19:13:11 +00:00
/// Удаляем старые куски.
2012-11-28 08:52:15 +00:00
parts . clear ( ) ;
2012-08-13 19:13:11 +00:00
clearOldParts ( ) ;
2012-11-28 08:52:15 +00:00
if ( ! while_can )
break ;
2012-08-13 19:13:11 +00:00
}
}
catch ( const Exception & e )
{
LOG_ERROR ( log , " Code: " < < e . code ( ) < < " . " < < e . displayText ( ) < < std : : endl
< < std : : endl
< < " Stack trace: " < < std : : endl
< < e . getStackTrace ( ) . toString ( ) ) ;
}
catch ( const Poco : : Exception & e )
{
LOG_ERROR ( log , " Poco::Exception: " < < e . code ( ) < < " . " < < e . displayText ( ) ) ;
}
catch ( const std : : exception & e )
{
LOG_ERROR ( log , " std::exception: " < < e . what ( ) ) ;
}
catch ( . . . )
{
LOG_ERROR ( log , " Unknown exception " ) ;
}
2012-07-23 06:23:29 +00:00
}
2012-11-28 08:52:15 +00:00
void StorageMergeTree : : joinMergeThreads ( )
{
LOG_DEBUG ( log , " Waiting for merge thread to finish. " ) ;
merge_threads - > wait ( ) ;
}
2012-11-29 10:50:17 +00:00
/// Выбираем отрезок из не более чем max_parts_to_merge_at_once кусков так, чтобы максимальный размер был меньше чем в max_size_ratio_to_merge_parts раз больше суммы остальных.
/// Это обеспечивает в худшем случае время O(n log n) на все слияния, независимо от выбора сливаемых кусков, порядка слияния и добавления.
2012-11-29 17:43:23 +00:00
/// При max_parts_to_merge_at_once >= log(max_rows_to_merge_parts/index_granularity)/log(max_size_ratio_to_merge_parts),
/// несложно доказать, что всегда будет что сливать, пока количество кусков больше
/// log(max_rows_to_merge_parts/index_granularity)/log(max_size_ratio_to_merge_parts)*(количество кусков размером больше max_rows_to_merge_parts).
2012-11-29 12:24:08 +00:00
/// Дальше эвристики.
/// Будем выбирать максимальный по включению подходящий отрезок.
/// Из всех таких выбираем отрезок с минимальным максимумом размера.
/// Из всех таких выбираем отрезок с минимальным минимумом размера.
/// Из всех таких выбираем отрезок с максимальной длиной.
2013-02-14 11:22:56 +00:00
bool StorageMergeTree : : selectPartsToMerge ( std : : vector < DataPartPtr > & parts , bool merge_anything_for_old_months )
2012-07-23 06:23:29 +00:00
{
LOG_DEBUG ( log , " Selecting parts to merge " ) ;
2012-08-10 20:04:34 +00:00
Poco : : ScopedLock < Poco : : FastMutex > lock ( data_parts_mutex ) ;
2012-07-23 06:23:29 +00:00
2013-02-14 11:22:56 +00:00
Yandex : : DateLUTSingleton & date_lut = Yandex : : DateLUTSingleton : : instance ( ) ;
2012-11-29 11:32:29 +00:00
size_t min_max = - 1U ;
size_t min_min = - 1U ;
int max_len = 0 ;
2012-11-29 10:50:17 +00:00
DataParts : : iterator best_begin ;
bool found = false ;
2013-02-14 11:22:56 +00:00
Yandex : : DayNum_t now_day = date_lut . toDayNum ( time ( 0 ) ) ;
Yandex : : DayNum_t now_month = date_lut . toFirstDayNumOfMonth ( now_day ) ;
2012-11-29 12:24:08 +00:00
/// Сколько кусков, начиная с текущего, можно включить в валидный отрезок, начинающийся левее текущего куска.
/// Нужно для определения максимальности по включению.
2012-11-29 12:26:34 +00:00
int max_count_from_left = 0 ;
2012-11-29 10:50:17 +00:00
2012-11-29 12:24:08 +00:00
/// Левый конец отрезка.
2012-11-29 10:50:17 +00:00
for ( DataParts : : iterator it = data_parts . begin ( ) ; it ! = data_parts . end ( ) ; + + it )
2012-07-23 06:23:29 +00:00
{
2012-11-29 12:24:08 +00:00
const DataPartPtr & first_part = * it ;
2012-11-29 12:26:34 +00:00
max_count_from_left = std : : max ( 0 , max_count_from_left - 1 ) ;
2012-11-29 10:50:17 +00:00
2012-11-29 12:24:08 +00:00
/// Кусок не занят и достаточно мал.
if ( first_part - > currently_merging | |
first_part - > size * index_granularity > settings . max_rows_to_merge_parts )
continue ;
/// Кусок в одном месяце.
if ( first_part - > left_month ! = first_part - > right_month )
{
LOG_WARNING ( log , " Part " < < first_part - > name < < " spans more than one month " ) ;
2012-11-29 10:50:17 +00:00
continue ;
2012-11-29 12:24:08 +00:00
}
/// Самый длинный валидный отрезок, начинающийся здесь.
2012-11-29 16:39:29 +00:00
size_t cur_longest_max = - 1U ;
size_t cur_longest_min = - 1U ;
2012-11-29 12:24:08 +00:00
int cur_longest_len = 0 ;
2012-11-29 10:50:17 +00:00
2012-11-29 12:24:08 +00:00
/// Текущий отрезок, не обязательно валидный.
size_t cur_max = first_part - > size ;
size_t cur_min = first_part - > size ;
size_t cur_sum = first_part - > size ;
2012-11-29 10:50:17 +00:00
int cur_len = 1 ;
2012-11-29 12:24:08 +00:00
Yandex : : DayNum_t month = first_part - > left_month ;
UInt64 cur_id = first_part - > right ;
2012-11-29 10:50:17 +00:00
2013-02-14 11:22:56 +00:00
/// Этот месяц кончился хотя бы день назад.
bool is_old_month = now_day - now_month > = 1 & & now_month > month ;
2012-11-29 12:24:08 +00:00
/// Правый конец отрезка.
2012-11-29 10:50:17 +00:00
DataParts : : iterator jt = it ;
2012-11-29 11:32:29 +00:00
for ( + + jt ; jt ! = data_parts . end ( ) & & cur_len < static_cast < int > ( settings . max_parts_to_merge_at_once ) ; + + jt )
2012-07-23 06:23:29 +00:00
{
2012-11-29 12:24:08 +00:00
const DataPartPtr & last_part = * jt ;
/// Кусок не занят, достаточно мал и в одном правильном месяце.
if ( last_part - > currently_merging | |
last_part - > size * index_granularity > settings . max_rows_to_merge_parts | |
last_part - > left_month ! = last_part - > right_month | |
last_part - > left_month ! = month )
2012-11-29 10:50:17 +00:00
break ;
2012-11-29 12:24:08 +00:00
/// Кусок правее предыдущего.
2012-11-30 00:52:45 +00:00
if ( last_part - > left < cur_id )
2012-11-29 12:24:08 +00:00
{
LOG_WARNING ( log , " Part " < < last_part - > name < < " intersects previous part " ) ;
break ;
}
cur_max = std : : max ( cur_max , last_part - > size ) ;
cur_min = std : : min ( cur_min , last_part - > size ) ;
cur_sum + = last_part - > size ;
2012-11-29 10:50:17 +00:00
+ + cur_len ;
2012-11-29 12:24:08 +00:00
cur_id = last_part - > right ;
2012-11-29 10:50:17 +00:00
2012-11-29 12:24:08 +00:00
/// Если отрезок валидный, то он самый длинный валидный, начинающийся тут.
2012-11-29 10:50:17 +00:00
if ( cur_len > = 2 & &
2013-02-14 11:22:56 +00:00
( static_cast < double > ( cur_max ) / ( cur_sum - cur_max ) < settings . max_size_ratio_to_merge_parts | |
( is_old_month & & merge_anything_for_old_months ) ) ) /// З а старый месяц объединяем что угодно, если разрешено.
2012-11-29 12:24:08 +00:00
{
cur_longest_max = cur_max ;
cur_longest_min = cur_min ;
cur_longest_len = cur_len ;
}
}
2012-11-29 12:26:34 +00:00
/// Это максимальный по включению валидный отрезок.
2012-11-29 12:24:08 +00:00
if ( cur_longest_len > max_count_from_left )
{
max_count_from_left = cur_longest_len ;
if ( ! found | |
std : : make_pair ( std : : make_pair ( cur_longest_max , cur_longest_min ) , - cur_longest_len ) <
std : : make_pair ( std : : make_pair ( min_max , min_min ) , - max_len ) )
2012-11-29 10:50:17 +00:00
{
found = true ;
2012-11-29 12:24:08 +00:00
min_max = cur_longest_max ;
min_min = cur_longest_min ;
max_len = cur_longest_len ;
2012-11-29 10:50:17 +00:00
best_begin = it ;
}
2012-07-23 06:23:29 +00:00
}
}
2012-11-29 10:50:17 +00:00
if ( found )
2012-07-23 06:23:29 +00:00
{
2012-11-29 10:50:17 +00:00
parts . clear ( ) ;
DataParts : : iterator it = best_begin ;
for ( int i = 0 ; i < max_len ; + + i )
2012-07-23 06:23:29 +00:00
{
2012-11-29 10:50:17 +00:00
parts . push_back ( * it ) ;
2012-11-29 11:48:27 +00:00
parts . back ( ) - > currently_merging = true ;
2012-11-29 10:50:17 +00:00
+ + it ;
2012-07-23 06:23:29 +00:00
}
2012-11-29 10:50:17 +00:00
2012-11-30 02:01:02 +00:00
LOG_DEBUG ( log , " Selected " < < parts . size ( ) < < " parts from " < < parts . front ( ) - > name < < " to " < < parts . back ( ) - > name ) ;
2012-07-23 06:23:29 +00:00
}
2012-11-29 10:50:17 +00:00
else
2012-07-23 06:23:29 +00:00
{
2012-11-29 10:50:17 +00:00
LOG_DEBUG ( log , " No parts to merge " ) ;
2012-07-23 06:23:29 +00:00
}
2012-11-29 10:50:17 +00:00
return found ;
2012-07-23 06:23:29 +00:00
}
2013-04-24 10:31:32 +00:00
/// parts должны быть отсортированы.
2012-11-28 08:52:15 +00:00
void StorageMergeTree : : mergeParts ( std : : vector < DataPartPtr > parts )
2012-07-23 06:23:29 +00:00
{
2012-11-30 02:01:02 +00:00
LOG_DEBUG ( log , " Merging " < < parts . size ( ) < < " parts: from " < < parts . front ( ) - > name < < " to " < < parts . back ( ) - > name ) ;
2012-07-30 20:32:36 +00:00
2012-08-13 19:13:11 +00:00
Names all_column_names ;
for ( NamesAndTypesList : : const_iterator it = columns - > begin ( ) ; it ! = columns - > end ( ) ; + + it )
all_column_names . push_back ( it - > first ) ;
2012-07-30 20:32:36 +00:00
2012-08-13 19:13:11 +00:00
Yandex : : DateLUTSingleton & date_lut = Yandex : : DateLUTSingleton : : instance ( ) ;
2012-07-30 20:32:36 +00:00
2012-08-13 19:13:11 +00:00
StorageMergeTree : : DataPartPtr new_data_part = new DataPart ( * this ) ;
2012-12-06 12:51:15 +00:00
new_data_part - > left_date = std : : numeric_limits < UInt16 > : : max ( ) ;
new_data_part - > right_date = std : : numeric_limits < UInt16 > : : min ( ) ;
2012-11-30 02:01:02 +00:00
new_data_part - > left = parts . front ( ) - > left ;
2012-11-28 08:52:15 +00:00
new_data_part - > right = parts . back ( ) - > right ;
new_data_part - > level = 0 ;
for ( size_t i = 0 ; i < parts . size ( ) ; + + i )
{
new_data_part - > level = std : : max ( new_data_part - > level , parts [ i ] - > level ) ;
2012-12-06 12:51:15 +00:00
new_data_part - > left_date = std : : min ( new_data_part - > left_date , parts [ i ] - > left_date ) ;
new_data_part - > right_date = std : : max ( new_data_part - > right_date , parts [ i ] - > right_date ) ;
2012-11-28 08:52:15 +00:00
}
+ + new_data_part - > level ;
2012-08-13 19:13:11 +00:00
new_data_part - > name = getPartName (
new_data_part - > left_date , new_data_part - > right_date , new_data_part - > left , new_data_part - > right , new_data_part - > level ) ;
2012-08-31 20:38:05 +00:00
new_data_part - > left_month = date_lut . toFirstDayNumOfMonth ( new_data_part - > left_date ) ;
new_data_part - > right_month = date_lut . toFirstDayNumOfMonth ( new_data_part - > right_date ) ;
2012-08-13 19:13:11 +00:00
2012-11-29 16:39:29 +00:00
/** Читаем из всех кусков, сливаем и пишем в новый.
2012-08-13 19:13:11 +00:00
* П о п у т н о в ы ч и с л я е м в ы р а ж е н и е д л я с о р т и р о в к и .
*/
BlockInputStreams src_streams ;
2012-07-30 20:32:36 +00:00
2012-11-28 08:52:15 +00:00
for ( size_t i = 0 ; i < parts . size ( ) ; + + i )
{
2012-12-06 09:45:09 +00:00
MarkRanges ranges ( 1 , MarkRange ( 0 , parts [ i ] - > size ) ) ;
2012-11-28 08:52:15 +00:00
src_streams . push_back ( new ExpressionBlockInputStream ( new MergeTreeBlockInputStream (
2013-01-23 17:38:03 +00:00
full_path + parts [ i ] - > name + ' / ' , DEFAULT_BLOCK_SIZE , all_column_names , * this , parts [ i ] , ranges , StoragePtr ( ) ) , primary_expr ) ) ;
2012-11-28 08:52:15 +00:00
}
2012-07-30 20:32:36 +00:00
2013-04-24 10:31:32 +00:00
/// Порядок потоков важен: при совпадении ключа элементы идет в порядке номера потока-источника.
/// В слитом куске строки с одинаковым ключом должны идти в порядке возрастания идентификатора исходного куска, то есть (примерного) возрастания времени вставки.
2012-11-30 00:52:45 +00:00
BlockInputStreamPtr merged_stream = sign_column . empty ( )
? new MergingSortedBlockInputStream ( src_streams , sort_descr , DEFAULT_BLOCK_SIZE )
: new CollapsingSortedBlockInputStream ( src_streams , sort_descr , sign_column , DEFAULT_BLOCK_SIZE ) ;
2012-08-13 20:16:06 +00:00
2012-12-03 08:52:58 +00:00
MergedBlockOutputStreamPtr to = new MergedBlockOutputStream ( * this ,
2012-11-28 08:52:15 +00:00
new_data_part - > left_date , new_data_part - > right_date , new_data_part - > left , new_data_part - > right , new_data_part - > level ) ;
2012-07-30 20:32:36 +00:00
2012-08-13 19:13:11 +00:00
copyData ( * merged_stream , * to ) ;
2012-07-31 20:03:53 +00:00
2012-12-03 08:52:58 +00:00
new_data_part - > size = to - > marksCount ( ) ;
2012-08-13 19:13:11 +00:00
new_data_part - > modification_time = time ( 0 ) ;
2012-07-31 17:07:20 +00:00
2012-08-13 19:13:11 +00:00
{
Poco : : ScopedLock < Poco : : FastMutex > lock ( data_parts_mutex ) ;
Poco : : ScopedLock < Poco : : FastMutex > lock_all ( all_data_parts_mutex ) ;
2012-07-31 17:07:20 +00:00
2012-08-13 19:13:11 +00:00
/// Добавляем новый кусок в набор.
2012-11-28 08:52:15 +00:00
for ( size_t i = 0 ; i < parts . size ( ) ; + + i )
{
if ( data_parts . end ( ) = = data_parts . find ( parts [ i ] ) )
throw Exception ( " Logical error: cannot find data part " + parts [ i ] - > name + " in list " , ErrorCodes : : LOGICAL_ERROR ) ;
}
2012-07-31 17:07:20 +00:00
2012-08-13 19:13:11 +00:00
data_parts . insert ( new_data_part ) ;
all_data_parts . insert ( new_data_part ) ;
2012-11-28 08:52:15 +00:00
for ( size_t i = 0 ; i < parts . size ( ) ; + + i )
{
data_parts . erase ( data_parts . find ( parts [ i ] ) ) ;
}
2012-07-31 17:07:20 +00:00
}
2012-08-13 19:13:11 +00:00
2012-12-13 20:28:10 +00:00
LOG_TRACE ( log , " Merged " < < parts . size ( ) < < " parts: from " < < parts . front ( ) - > name < < " to " < < parts . back ( ) - > name ) ;
2012-07-23 06:23:29 +00:00
}
2012-08-16 18:17:01 +00:00
2013-01-23 11:16:32 +00:00
void StorageMergeTree : : rename ( const String & new_path_to_db , const String & new_name )
{
joinMergeThreads ( ) ;
std : : string new_full_path = new_path_to_db + escapeForFileName ( new_name ) + ' / ' ;
Poco : : File ( full_path ) . renameTo ( new_full_path ) ;
path = new_path_to_db ;
full_path = new_full_path ;
name = new_name ;
increment . setPath ( full_path + " increment.txt " ) ;
}
2013-01-23 17:43:19 +00:00
void StorageMergeTree : : dropImpl ( )
2012-08-16 18:17:01 +00:00
{
2012-11-28 08:52:15 +00:00
joinMergeThreads ( ) ;
2012-08-16 18:17:01 +00:00
Poco : : ScopedLock < Poco : : FastMutex > lock ( data_parts_mutex ) ;
Poco : : ScopedLock < Poco : : FastMutex > lock_all ( all_data_parts_mutex ) ;
data_parts . clear ( ) ;
all_data_parts . clear ( ) ;
Poco : : File ( full_path ) . remove ( true ) ;
}
2012-07-17 20:04:39 +00:00
}