Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
2dot5
ClickHouse
提交
2e15ce67
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,发现更多精彩内容 >>
未验证
提交
2e15ce67
编写于
3月 29, 2018
作者:
A
alexey-milovidov
提交者:
GitHub
3月 29, 2018
浏览文件
操作
浏览文件
下载
差异文件
Merge pull request #2132 from yandex/small-enhancements
Small enhancements
上级
a2bc0462
de3791db
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
67 addition
and
28 deletion
+67
-28
dbms/src/Client/Connection.cpp
dbms/src/Client/Connection.cpp
+1
-26
dbms/src/Client/TimeoutSetter.h
dbms/src/Client/TimeoutSetter.h
+46
-0
dbms/src/Server/TCPHandler.cpp
dbms/src/Server/TCPHandler.cpp
+10
-1
dbms/src/Server/TCPHandler.h
dbms/src/Server/TCPHandler.h
+4
-0
dbms/src/Storages/MergeTree/ReplicatedMergeTreeQueue.cpp
dbms/src/Storages/MergeTree/ReplicatedMergeTreeQueue.cpp
+6
-1
未找到文件。
dbms/src/Client/Connection.cpp
浏览文件 @
2e15ce67
...
...
@@ -12,6 +12,7 @@
#include <DataStreams/NativeBlockInputStream.h>
#include <DataStreams/NativeBlockOutputStream.h>
#include <Client/Connection.h>
#include <Client/TimeoutSetter.h>
#include <Common/ClickHouseRevision.h>
#include <Common/Exception.h>
#include <Common/NetException.h>
...
...
@@ -228,32 +229,6 @@ void Connection::forceConnected()
}
}
struct
TimeoutSetter
{
TimeoutSetter
(
Poco
::
Net
::
StreamSocket
&
socket_
,
const
Poco
::
Timespan
&
timeout_
)
:
socket
(
socket_
),
timeout
(
timeout_
)
{
old_send_timeout
=
socket
.
getSendTimeout
();
old_receive_timeout
=
socket
.
getReceiveTimeout
();
if
(
old_send_timeout
>
timeout
)
socket
.
setSendTimeout
(
timeout
);
if
(
old_receive_timeout
>
timeout
)
socket
.
setReceiveTimeout
(
timeout
);
}
~
TimeoutSetter
()
{
socket
.
setSendTimeout
(
old_send_timeout
);
socket
.
setReceiveTimeout
(
old_receive_timeout
);
}
Poco
::
Net
::
StreamSocket
&
socket
;
Poco
::
Timespan
timeout
;
Poco
::
Timespan
old_send_timeout
;
Poco
::
Timespan
old_receive_timeout
;
};
bool
Connection
::
ping
()
{
// LOG_TRACE(log_wrapper.get(), "Ping");
...
...
dbms/src/Client/TimeoutSetter.h
0 → 100644
浏览文件 @
2e15ce67
#pragma once
#include <Poco/Timespan.h>
#include <Poco/Net/StreamSocket.h>
namespace
DB
{
/// Temporarily overrides socket send/recieve timeouts and reset them back into destructor
/// Timeouts could be only decreased
struct
TimeoutSetter
{
TimeoutSetter
(
Poco
::
Net
::
StreamSocket
&
socket_
,
const
Poco
::
Timespan
&
send_timeout_
,
const
Poco
::
Timespan
&
recieve_timeout_
)
:
socket
(
socket_
),
send_timeout
(
send_timeout_
),
recieve_timeout
(
recieve_timeout_
)
{
old_send_timeout
=
socket
.
getSendTimeout
();
old_receive_timeout
=
socket
.
getReceiveTimeout
();
if
(
old_send_timeout
>
send_timeout
)
socket
.
setSendTimeout
(
send_timeout
);
if
(
old_receive_timeout
>
recieve_timeout
)
socket
.
setReceiveTimeout
(
recieve_timeout
);
}
TimeoutSetter
(
Poco
::
Net
::
StreamSocket
&
socket_
,
const
Poco
::
Timespan
&
timeout_
)
:
TimeoutSetter
(
socket_
,
timeout_
,
timeout_
)
{}
~
TimeoutSetter
()
{
socket
.
setSendTimeout
(
old_send_timeout
);
socket
.
setReceiveTimeout
(
old_receive_timeout
);
}
Poco
::
Net
::
StreamSocket
&
socket
;
Poco
::
Timespan
send_timeout
;
Poco
::
Timespan
recieve_timeout
;
Poco
::
Timespan
old_send_timeout
;
Poco
::
Timespan
old_receive_timeout
;
};
}
dbms/src/Server/TCPHandler.cpp
浏览文件 @
2e15ce67
...
...
@@ -31,6 +31,7 @@
#include "TCPHandler.h"
#include <Common/NetException.h>
#include <ext/scope_guard.h>
namespace
DB
...
...
@@ -139,6 +140,9 @@ void TCPHandler::runImpl()
/// Restore context of request.
query_context
=
connection_context
;
/// If a user passed query-local timeouts, reset socket to initial state at the end of the query
SCOPE_EXIT
({
state
.
timeout_setter
.
reset
();});
/** If Query - process it. If Ping or Cancel - go back to the beginning.
* There may come settings for a separate query that modify `query_context`.
*/
...
...
@@ -600,7 +604,12 @@ void TCPHandler::receiveQuery()
}
/// Per query settings.
query_context
.
getSettingsRef
().
deserialize
(
*
in
);
Settings
&
settings
=
query_context
.
getSettingsRef
();
settings
.
deserialize
(
*
in
);
/// Sync timeouts on client and server during current query to avoid dangling queries on server
/// NOTE: these settings are applied only for current connection (not for distributed tables' connections)
state
.
timeout_setter
=
std
::
make_unique
<
TimeoutSetter
>
(
socket
(),
settings
.
send_timeout
,
settings
.
receive_timeout
);
readVarUInt
(
stage
,
*
in
);
state
.
stage
=
QueryProcessingStage
::
Enum
(
stage
);
...
...
dbms/src/Server/TCPHandler.h
浏览文件 @
2e15ce67
...
...
@@ -11,6 +11,7 @@
#include <DataStreams/BlockIO.h>
#include <IO/ReadHelpers.h>
#include <IO/WriteHelpers.h>
#include <Client/TimeoutSetter.h>
#include "IServer.h"
...
...
@@ -59,6 +60,9 @@ struct QueryState
/// To output progress, the difference after the previous sending of progress.
Progress
progress
;
/// Timeouts setter for current query
std
::
unique_ptr
<
TimeoutSetter
>
timeout_setter
;
void
reset
()
{
...
...
dbms/src/Storages/MergeTree/ReplicatedMergeTreeQueue.cpp
浏览文件 @
2e15ce67
...
...
@@ -94,7 +94,12 @@ void ReplicatedMergeTreeQueue::initialize(
void
ReplicatedMergeTreeQueue
::
insertUnlocked
(
LogEntryPtr
&
entry
,
std
::
optional
<
time_t
>
&
min_unprocessed_insert_time_changed
,
std
::
lock_guard
<
std
::
mutex
>
&
)
{
virtual_parts
.
add
(
entry
->
new_part_name
);
queue
.
push_back
(
entry
);
/// Put 'DROP PARTITION' entries at the beginning of the queue not to make superfluous fetches of parts that will be eventually deleted
if
(
entry
->
type
!=
LogEntry
::
DROP_RANGE
)
queue
.
push_back
(
entry
);
else
queue
.
push_front
(
entry
);
if
(
entry
->
type
==
LogEntry
::
GET_PART
)
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录