Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
2dot5
ClickHouse
提交
8266715c
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,发现更多精彩内容 >>
提交
8266715c
编写于
5月 31, 2020
作者:
K
kssenii
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Fix build & fix style
上级
0362bb2d
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
11 addition
and
7 deletion
+11
-7
src/Storages/RabbitMQ/RabbitMQBlockInputStream.cpp
src/Storages/RabbitMQ/RabbitMQBlockInputStream.cpp
+1
-2
src/Storages/RabbitMQ/RabbitMQHandler.cpp
src/Storages/RabbitMQ/RabbitMQHandler.cpp
+1
-1
src/Storages/RabbitMQ/ReadBufferFromRabbitMQConsumer.cpp
src/Storages/RabbitMQ/ReadBufferFromRabbitMQConsumer.cpp
+3
-0
src/Storages/RabbitMQ/StorageRabbitMQ.cpp
src/Storages/RabbitMQ/StorageRabbitMQ.cpp
+6
-4
未找到文件。
src/Storages/RabbitMQ/RabbitMQBlockInputStream.cpp
浏览文件 @
8266715c
...
...
@@ -64,8 +64,7 @@ Block RabbitMQBlockInputStream::readImpl()
MutableColumns
result_columns
=
non_virtual_header
.
cloneEmptyColumns
();
MutableColumns
virtual_columns
=
virtual_header
.
cloneEmptyColumns
();
auto
input_format
=
FormatFactory
::
instance
().
getInputFormat
(
storage
.
getFormatName
(),
*
buffer
,
non_virtual_header
,
context
,
1
);
auto
input_format
=
FormatFactory
::
instance
().
getInputFormat
(
storage
.
getFormatName
(),
*
buffer
,
non_virtual_header
,
context
,
1
);
InputPort
port
(
input_format
->
getPort
().
getHeader
(),
input_format
.
get
());
connect
(
input_format
->
getPort
(),
port
);
...
...
src/Storages/RabbitMQ/RabbitMQHandler.cpp
浏览文件 @
8266715c
...
...
@@ -12,7 +12,7 @@ RabbitMQHandler::RabbitMQHandler(event_base * evbase_, Poco::Logger * log_) :
}
void
RabbitMQHandler
::
onError
(
AMQP
::
TcpConnection
*
,
const
char
*
message
)
void
RabbitMQHandler
::
onError
(
AMQP
::
TcpConnection
*
/* connection */
,
const
char
*
message
)
{
LOG_ERROR
(
log
,
"Library error report: {}"
,
message
);
stop
();
...
...
src/Storages/RabbitMQ/ReadBufferFromRabbitMQConsumer.cpp
浏览文件 @
8266715c
#include <utility>
#include <chrono>
#include <thread>
#include <mutex>
#include <atomic>
#include <memory>
#include <Storages/RabbitMQ/ReadBufferFromRabbitMQConsumer.h>
#include <Storages/RabbitMQ/RabbitMQHandler.h>
#include <common/logger_useful.h>
...
...
src/Storages/RabbitMQ/StorageRabbitMQ.cpp
浏览文件 @
8266715c
...
...
@@ -45,10 +45,8 @@ namespace DB
namespace
ErrorCodes
{
extern
const
int
NOT_IMPLEMENTED
;
extern
const
int
LOGICAL_ERROR
;
extern
const
int
BAD_ARGUMENTS
;
extern
const
int
NUMBER_OF_ARGUMENTS_DOESNT_MATCH
;
}
...
...
@@ -157,6 +155,7 @@ void StorageRabbitMQ::shutdown()
popReadBuffer
();
}
connection
.
close
();
task
->
deactivate
();
}
...
...
@@ -201,8 +200,10 @@ ConsumerBufferPtr StorageRabbitMQ::createReadBuffer()
next_channel_id
+=
num_queues
;
update_channel_id
=
true
;
ChannelPtr
consumer_channel
=
std
::
make_shared
<
AMQP
::
TcpChannel
>
(
&
connection
);
return
std
::
make_shared
<
ReadBufferFromRabbitMQConsumer
>
(
std
::
make_shared
<
AMQP
::
TcpChannel
>
(
&
connection
)
,
eventHandler
,
exchange_name
,
routing_key
,
next_channel_id
,
consumer_channel
,
eventHandler
,
exchange_name
,
routing_key
,
next_channel_id
,
log
,
row_delimiter
,
bind_by_id
,
hash_exchange
,
num_queues
,
stream_cancelled
);
}
...
...
@@ -460,7 +461,8 @@ void registerStorageRabbitMQ(StorageFactory & factory)
}
}
return
StorageRabbitMQ
::
create
(
args
.
table_id
,
args
.
context
,
args
.
columns
,
host_port
,
routing_key
,
exchange
,
return
StorageRabbitMQ
::
create
(
args
.
table_id
,
args
.
context
,
args
.
columns
,
host_port
,
routing_key
,
exchange
,
format
,
row_delimiter
,
num_consumers
,
num_queues
,
hash_exchange
);
};
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录