Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
2dot5
ClickHouse
提交
e57d1c82
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,发现更多精彩内容 >>
提交
e57d1c82
编写于
8月 31, 2020
作者:
K
kssenii
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Better shutdown
上级
647cf571
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
33 addition
and
15 deletion
+33
-15
src/Storages/RabbitMQ/StorageRabbitMQ.cpp
src/Storages/RabbitMQ/StorageRabbitMQ.cpp
+31
-14
src/Storages/RabbitMQ/StorageRabbitMQ.h
src/Storages/RabbitMQ/StorageRabbitMQ.h
+2
-1
未找到文件。
src/Storages/RabbitMQ/StorageRabbitMQ.cpp
浏览文件 @
e57d1c82
...
...
@@ -210,6 +210,28 @@ void StorageRabbitMQ::loopingFunc()
}
/* Need to deactivate this way because otherwise might get a deadlock when first deactivate streaming task in shutdown and then
* inside streaming task try to deactivate any other task
*/
void
StorageRabbitMQ
::
deactivateTask
(
BackgroundSchedulePool
::
TaskHolder
&
task
,
bool
wait
,
bool
stop_loop
)
{
if
(
task_mutex
.
try_lock
())
{
if
(
stop_loop
)
event_handler
->
updateLoopState
(
Loop
::
STOP
);
task
->
deactivate
();
task_mutex
.
unlock
();
}
else
if
(
wait
)
{
/// Wait only if deactivating from shutdown
std
::
lock_guard
lock
(
task_mutex
);
task
->
deactivate
();
}
}
void
StorageRabbitMQ
::
initExchange
()
{
/* Binding scheme is the following: client's exchange -> key bindings by routing key list -> bridge exchange (fanout) ->
...
...
@@ -326,7 +348,7 @@ bool StorageRabbitMQ::restoreConnection(bool reconnecting)
if
(
reconnecting
)
{
heartbeat_task
->
deactivate
(
);
deactivateTask
(
heartbeat_task
,
0
,
0
);
connection
->
close
();
/// Connection might be unusable, but not closed
/* Connection is not closed immediately (firstly, all pending operations are completed, and then
...
...
@@ -346,7 +368,7 @@ bool StorageRabbitMQ::restoreConnection(bool reconnecting)
AMQP
::
Address
(
parsed_address
.
first
,
parsed_address
.
second
,
AMQP
::
Login
(
login_password
.
first
,
login_password
.
second
),
"/"
));
cnt_retries
=
0
;
while
(
!
connection
->
ready
()
&&
++
cnt_retries
!=
RETRIES_MAX
)
while
(
!
connection
->
ready
()
&&
!
stream_cancelled
&&
++
cnt_retries
!=
RETRIES_MAX
)
{
event_handler
->
iterateLoop
();
std
::
this_thread
::
sleep_for
(
std
::
chrono
::
milliseconds
(
CONNECT_SLEEP
));
...
...
@@ -504,11 +526,9 @@ void StorageRabbitMQ::shutdown()
stream_cancelled
=
true
;
wait_confirm
.
store
(
false
);
streaming_task
->
deactivate
();
heartbeat_task
->
deactivate
();
event_handler
->
updateLoopState
(
Loop
::
STOP
);
looping_task
->
deactivate
();
deactivateTask
(
streaming_task
,
1
,
1
);
deactivateTask
(
heartbeat_task
,
1
,
0
);
deactivateTask
(
looping_task
,
1
,
1
);
connection
->
close
();
...
...
@@ -695,14 +715,11 @@ bool StorageRabbitMQ::streamToViews()
* races inside the library, but only in case any error occurs or connection is lost while ack is being sent
*/
if
(
event_handler
->
loopRunning
())
{
event_handler
->
updateLoopState
(
Loop
::
STOP
);
looping_task
->
deactivate
();
}
deactivateTask
(
looping_task
,
0
,
1
);
if
(
!
event_handler
->
connectionRunning
())
{
if
(
restoreConnection
(
true
))
if
(
!
stream_cancelled
&&
restoreConnection
(
true
))
{
for
(
auto
&
stream
:
streams
)
stream
->
as
<
RabbitMQBlockInputStream
>
()
->
updateChannel
();
...
...
@@ -710,13 +727,13 @@ bool StorageRabbitMQ::streamToViews()
}
else
{
/// Reschedule if unable to connect to rabbitmq
/// Reschedule if unable to connect to rabbitmq
or quit if cancelled
return
false
;
}
}
else
{
heartbeat_task
->
deactivate
(
);
deactivateTask
(
heartbeat_task
,
0
,
0
);
/// Commit
for
(
auto
&
stream
:
streams
)
...
...
src/Storages/RabbitMQ/StorageRabbitMQ.h
浏览文件 @
e57d1c82
...
...
@@ -101,7 +101,7 @@ private:
size_t
num_created_consumers
=
0
;
Poco
::
Semaphore
semaphore
;
std
::
mutex
mutex
;
std
::
mutex
mutex
,
task_mutex
;
std
::
vector
<
ConsumerBufferPtr
>
buffers
;
/// available buffers for RabbitMQ consumers
String
unique_strbase
;
...
...
@@ -128,6 +128,7 @@ private:
AMQP
::
ExchangeType
defineExchangeType
(
String
exchange_type_
);
size_t
getMaxBlockSize
();
String
getTableBasedName
(
String
name
,
const
StorageID
&
table_id
);
void
deactivateTask
(
BackgroundSchedulePool
::
TaskHolder
&
task
,
bool
wait
,
bool
stop_loop
);
void
initExchange
();
void
bindExchange
();
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录