Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
2dot5
ClickHouse
提交
f0ca756f
C
ClickHouse
项目概览
2dot5
/
ClickHouse
通知
3
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
C
ClickHouse
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
提交
f0ca756f
编写于
2月 08, 2017
作者:
A
Alexey Milovidov
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
GraphiteMergeTree: fixed handling of 'version' column [#CLICKHOUSE-2804].
上级
083e9cc3
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
65 addition
and
46 deletion
+65
-46
dbms/include/DB/DataStreams/GraphiteRollupSortedBlockInputStream.h
...ude/DB/DataStreams/GraphiteRollupSortedBlockInputStream.h
+7
-7
dbms/include/DB/DataStreams/ReplacingSortedBlockInputStream.h
.../include/DB/DataStreams/ReplacingSortedBlockInputStream.h
+5
-5
dbms/src/DataStreams/GraphiteRollupSortedBlockInputStream.cpp
.../src/DataStreams/GraphiteRollupSortedBlockInputStream.cpp
+52
-33
dbms/src/DataStreams/ReplacingSortedBlockInputStream.cpp
dbms/src/DataStreams/ReplacingSortedBlockInputStream.cpp
+1
-1
未找到文件。
dbms/include/DB/DataStreams/GraphiteRollupSortedBlockInputStream.h
浏览文件 @
f0ca756f
...
...
@@ -179,15 +179,16 @@ private:
/// All data has been read.
bool
finished
=
false
;
/// This object owns a column while 'current_path' references to it
.
SharedBlockPtr
owned_current_block
;
RowRef
selected_row
;
/// Last row with maximum version for current primary key
.
UInt64
current_max_version
=
0
;
bool
is_first
=
true
;
StringRef
current_path
;
time_t
current_time
;
time_t
current_time
=
0
;
time_t
current_time_rounded
=
0
;
StringRef
next_path
;
time_t
next_time
;
UInt64
current_max_version
=
0
;
time_t
next_time
=
0
;
time_t
next_time_rounded
=
0
;
const
Graphite
::
Pattern
*
current_pattern
=
nullptr
;
std
::
vector
<
char
>
place_for_aggregate_state
;
...
...
@@ -207,8 +208,7 @@ private:
void
finishCurrentRow
(
ColumnPlainPtrs
&
merged_columns
);
/// Обновить состояние агрегатной функции новым значением value.
template
<
class
TSortCursor
>
void
accumulateRow
(
TSortCursor
&
cursor
);
void
accumulateRow
(
RowRef
&
row
);
};
}
dbms/include/DB/DataStreams/ReplacingSortedBlockInputStream.h
浏览文件 @
f0ca756f
...
...
@@ -51,15 +51,15 @@ private:
Logger
*
log
=
&
Logger
::
get
(
"ReplacingSortedBlockInputStream"
);
///
Прочитали до конца
.
///
All data has been read
.
bool
finished
=
false
;
RowRef
current_key
;
///
Текущий первичный ключ
.
RowRef
next_key
;
///
Первичный ключ следующей строки
.
RowRef
current_key
;
///
Primary key of current row
.
RowRef
next_key
;
///
Primary key of next row
.
RowRef
selected_row
;
///
Последняя строка с максимальной версией для текущего первичного ключа
.
RowRef
selected_row
;
///
Last row with maximum version for current primary key
.
UInt64
max_version
=
0
;
///
Максимальная версия для текущего первичного ключа
.
UInt64
max_version
=
0
;
///
Max version for current primary key
.
template
<
class
TSortCursor
>
void
merge
(
ColumnPlainPtrs
&
merged_columns
,
std
::
priority_queue
<
TSortCursor
>
&
queue
);
...
...
dbms/src/DataStreams/GraphiteRollupSortedBlockInputStream.cpp
浏览文件 @
f0ca756f
...
...
@@ -93,6 +93,9 @@ Block GraphiteRollupSortedBlockInputStream::readImpl()
for
(
size_t
i
=
0
;
i
<
num_columns
;
++
i
)
if
(
i
!=
time_column_num
&&
i
!=
value_column_num
&&
i
!=
version_column_num
)
unmodified_column_numbers
.
push_back
(
i
);
if
(
selected_row
.
empty
())
selected_row
.
columns
.
resize
(
num_columns
);
}
if
(
has_collation
)
...
...
@@ -123,46 +126,63 @@ void GraphiteRollupSortedBlockInputStream::merge(ColumnPlainPtrs & merged_column
bool
path_differs
=
is_first
||
next_path
!=
current_path
;
is_first
=
false
;
if
(
path_differs
)
current_pattern
=
selectPatternForPath
(
next_path
);
if
(
current_pattern
)
{
UInt32
precision
=
selectPrecision
(
current_pattern
->
retentions
,
next_time
);
next_time
=
roundTimeToPrecision
(
date_lut
,
next_time
,
precision
);
}
/// If no patterns has matched - it means that no need to do rounding.
/// Is new key before rounding.
bool
is_new_key
=
path_differs
||
next_time
!=
current_time
;
UInt64
current_version
=
current
->
all_columns
[
version_column_num
]
->
get64
(
current
->
pos
);
if
(
is_new_key
)
{
if
(
merged_rows
)
{
finishCurrentRow
(
merged_columns
);
current_max_version
=
0
;
if
(
path_differs
)
current_pattern
=
selectPatternForPath
(
next_path
);
/// if we have enough rows
if
(
merged_rows
>=
max_block_size
)
return
;
if
(
current_pattern
)
{
UInt32
precision
=
selectPrecision
(
current_pattern
->
retentions
,
next_time
);
next_time_rounded
=
roundTimeToPrecision
(
date_lut
,
next_time
,
precision
);
}
/// If no patterns has matched - it means that no need to do rounding.
startNextRow
(
merged_columns
,
current
);
/// Key will be new after rounding.
bool
will_be_new_key
=
path_differs
||
next_time_rounded
!=
current_time_rounded
;
owned_current_block
=
source_blocks
[
current
.
impl
->
order
];
current_path
=
next_path
;
current_time
=
next_time
;
current_max_version
=
0
;
if
(
will_be_new_key
)
{
if
(
merged_rows
)
{
finishCurrentRow
(
merged_columns
);
if
(
prev_pattern
)
prev_pattern
->
function
->
destroy
(
place_for_aggregate_state
.
data
());
if
(
current_pattern
)
current_pattern
->
function
->
create
(
place_for_aggregate_state
.
data
());
/// if we have enough rows
if
(
merged_rows
>=
max_block_size
)
return
;
}
++
merged_rows
;
}
startNextRow
(
merged_columns
,
current
);
accumulateRow
(
current
);
current_max_version
=
std
::
max
(
current_max_version
,
current
->
all_columns
[
version_column_num
]
->
get64
(
current
->
pos
));
current_path
=
next_path
;
current_time
=
next_time
;
current_time_rounded
=
next_time_rounded
;
if
(
prev_pattern
)
prev_pattern
->
function
->
destroy
(
place_for_aggregate_state
.
data
());
if
(
current_pattern
)
current_pattern
->
function
->
create
(
place_for_aggregate_state
.
data
());
++
merged_rows
;
}
accumulateRow
(
selected_row
);
current_max_version
=
current_version
;
}
else
if
(
current_version
>=
current_max_version
)
{
/// Within all rows with same key, we should leave only one row with maximum version;
/// and for rows with same maximum version - only last row.
current_max_version
=
current_version
;
setRowRef
(
selected_row
,
current
);
}
queue
.
pop
();
...
...
@@ -210,7 +230,7 @@ void GraphiteRollupSortedBlockInputStream::startNextRow(ColumnPlainPtrs & merged
void
GraphiteRollupSortedBlockInputStream
::
finishCurrentRow
(
ColumnPlainPtrs
&
merged_columns
)
{
/// Вставляем вычисленные значения столбцов time, value, version.
merged_columns
[
time_column_num
]
->
insert
(
UInt64
(
current_time
));
merged_columns
[
time_column_num
]
->
insert
(
UInt64
(
current_time
_rounded
));
merged_columns
[
version_column_num
]
->
insert
(
current_max_version
);
if
(
current_pattern
)
...
...
@@ -218,11 +238,10 @@ void GraphiteRollupSortedBlockInputStream::finishCurrentRow(ColumnPlainPtrs & me
}
template
<
class
TSortCursor
>
void
GraphiteRollupSortedBlockInputStream
::
accumulateRow
(
TSortCursor
&
cursor
)
void
GraphiteRollupSortedBlockInputStream
::
accumulateRow
(
RowRef
&
row
)
{
if
(
current_pattern
)
current_pattern
->
function
->
add
(
place_for_aggregate_state
.
data
(),
&
cursor
->
all_columns
[
value_column_num
],
cursor
->
pos
,
nullptr
);
current_pattern
->
function
->
add
(
place_for_aggregate_state
.
data
(),
&
row
.
columns
[
value_column_num
],
row
.
row_num
,
nullptr
);
}
}
dbms/src/DataStreams/ReplacingSortedBlockInputStream.cpp
浏览文件 @
f0ca756f
...
...
@@ -29,7 +29,7 @@ Block ReplacingSortedBlockInputStream::readImpl()
if
(
merged_columns
.
empty
())
return
Block
();
///
Дополнительная инициализация
.
///
Additional initialization
.
if
(
selected_row
.
empty
())
{
selected_row
.
columns
.
resize
(
num_columns
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录