Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
2dot5
ClickHouse
提交
f5631939
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,发现更多精彩内容 >>
提交
f5631939
编写于
6月 16, 2020
作者:
N
Nikolai Kochetov
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Add MergeSortingStep.
上级
815ac038
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
92 addition
and
10 deletion
+92
-10
src/Interpreters/InterpreterSelectQuery.cpp
src/Interpreters/InterpreterSelectQuery.cpp
+9
-10
src/Processors/QueryPlan/MergeSortingStep.cpp
src/Processors/QueryPlan/MergeSortingStep.cpp
+43
-0
src/Processors/QueryPlan/MergeSortingStep.h
src/Processors/QueryPlan/MergeSortingStep.h
+39
-0
src/Processors/ya.make
src/Processors/ya.make
+1
-0
未找到文件。
src/Interpreters/InterpreterSelectQuery.cpp
浏览文件 @
f5631939
...
...
@@ -80,6 +80,7 @@
#include <Processors/QueryPlan/ReadNothingStep.h>
#include <Processors/QueryPlan/ReadFromPreparedSource.h>
#include <Processors/QueryPlan/PartialSortingStep.h>
#include <Processors/QueryPlan/MergeSortingStep.h>
namespace
DB
...
...
@@ -1704,17 +1705,15 @@ void InterpreterSelectQuery::executeOrder(QueryPipeline & pipeline, InputOrderIn
partial_sorting
.
transformPipeline
(
pipeline
);
/// Merge the sorted blocks.
pipeline
.
addSimpleTransform
([
&
](
const
Block
&
header
,
QueryPipeline
::
StreamType
stream_type
)
->
ProcessorPtr
{
if
(
stream_type
==
QueryPipeline
::
StreamType
::
Totals
)
return
nullptr
;
MergeSortingStep
merge_sorting_step
(
DataStream
{.
header
=
pipeline
.
getHeader
()},
output_order_descr
,
settings
.
max_block_size
,
limit
,
settings
.
max_bytes_before_remerge_sort
/
pipeline
.
getNumStreams
(),
settings
.
max_bytes_before_external_sort
,
context
->
getTemporaryVolume
(),
settings
.
min_free_disk_space_for_temporary_data
);
return
std
::
make_shared
<
MergeSortingTransform
>
(
header
,
output_order_descr
,
settings
.
max_block_size
,
limit
,
settings
.
max_bytes_before_remerge_sort
/
pipeline
.
getNumStreams
(),
settings
.
max_bytes_before_external_sort
,
context
->
getTemporaryVolume
(),
settings
.
min_free_disk_space_for_temporary_data
);
});
merge_sorting_step
.
setStepDescription
(
"Merge sorted blocks before ORDER BY"
);
merge_sorting_step
.
transformPipeline
(
pipeline
);
/// If there are several streams, we merge them into one
executeMergeSorted
(
pipeline
,
output_order_descr
,
limit
);
...
...
src/Processors/QueryPlan/MergeSortingStep.cpp
0 → 100644
浏览文件 @
f5631939
#include <Processors/QueryPlan/MergeSortingStep.h>
#include <Processors/QueryPipeline.h>
#include <Processors/Transforms/MergeSortingTransform.h>
namespace
DB
{
MergeSortingStep
::
MergeSortingStep
(
const
DataStream
&
input_stream
,
const
SortDescription
&
description_
,
size_t
max_merged_block_size_
,
UInt64
limit_
,
size_t
max_bytes_before_remerge_
,
size_t
max_bytes_before_external_sort_
,
VolumePtr
tmp_volume_
,
size_t
min_free_disk_space_
)
:
ITransformingStep
(
input_stream
,
input_stream
)
,
description
(
description_
)
,
max_merged_block_size
(
max_merged_block_size_
)
,
limit
(
limit_
)
,
max_bytes_before_remerge
(
max_bytes_before_remerge_
)
,
max_bytes_before_external_sort
(
max_bytes_before_external_sort_
),
tmp_volume
(
tmp_volume_
)
,
min_free_disk_space
(
min_free_disk_space_
)
{
}
void
MergeSortingStep
::
transformPipeline
(
QueryPipeline
&
pipeline
)
{
pipeline
.
addSimpleTransform
([
&
](
const
Block
&
header
,
QueryPipeline
::
StreamType
stream_type
)
->
ProcessorPtr
{
if
(
stream_type
==
QueryPipeline
::
StreamType
::
Totals
)
return
nullptr
;
return
std
::
make_shared
<
MergeSortingTransform
>
(
header
,
description
,
max_merged_block_size
,
limit
,
max_bytes_before_remerge
,
max_bytes_before_external_sort
,
tmp_volume
,
min_free_disk_space
);
});
}
}
src/Processors/QueryPlan/MergeSortingStep.h
0 → 100644
浏览文件 @
f5631939
#pragma once
#include <Processors/QueryPlan/ITransformingStep.h>
#include <Core/SortDescription.h>
#include <DataStreams/SizeLimits.h>
#include <Disks/IVolume.h>
namespace
DB
{
class
MergeSortingStep
:
public
ITransformingStep
{
public:
explicit
MergeSortingStep
(
const
DataStream
&
input_stream
,
const
SortDescription
&
description_
,
size_t
max_merged_block_size_
,
UInt64
limit_
,
size_t
max_bytes_before_remerge_
,
size_t
max_bytes_before_external_sort_
,
VolumePtr
tmp_volume_
,
size_t
min_free_disk_space_
);
String
getName
()
const
override
{
return
"MergeSorting"
;
}
void
transformPipeline
(
QueryPipeline
&
pipeline
)
override
;
private:
SortDescription
description
;
size_t
max_merged_block_size
;
UInt64
limit
;
size_t
max_bytes_before_remerge
;
size_t
max_bytes_before_external_sort
;
VolumePtr
tmp_volume
;
size_t
min_free_disk_space
;
};
}
src/Processors/ya.make
浏览文件 @
f5631939
...
...
@@ -142,6 +142,7 @@ SRCS(
QueryPlan/ISourceStep.cpp
QueryPlan/ITransformingStep.cpp
QueryPlan/IQueryPlanStep.cpp
QueryPlan/MergeSortingStep.cpp
QueryPlan/PartialSortingStep.cpp
QueryPlan/ReadFromStorageStep.cpp
QueryPlan/ReadNothingStep.cpp
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录