切换导航条
此项目
正在载入...
登录
mmm-go
/
partnermg
·
提交
转到一个项目
GitLab
转到群组
项目
活动
文件
提交
管道
0
构建
0
图表
里程碑
问题
0
合并请求
0
成员
标记
维基
派生
网络
创建新的问题
下载为
邮件补丁
差异文件
浏览文件
作者
唐旭辉
4 years ago
提交
8668ebb84e178eba241ab88d11c4377e556da300
1 个父辈
ccc127fb
调试
隐藏空白字符变更
内嵌
并排对比
正在显示
2 个修改的文件
包含
3 行增加
和
3 行删除
pkg/port/consumer/consumer.go
pkg/port/consumer/produce/produce.go
pkg/port/consumer/consumer.go
查看文件 @
8668ebb
...
...
@@ -95,9 +95,9 @@ func NewRuner() *Runer {
func
(
r
*
Runer
)
InitConsumer
()
error
{
config
:=
sarama
.
NewConfig
()
config
.
Consumer
.
Group
.
Rebalance
.
Strategy
=
sarama
.
BalanceStrategyRoundRobin
//
config.Consumer.Group.Rebalance.Strategy = sarama.BalanceStrategyRoundRobin
config
.
Consumer
.
Offsets
.
Initial
=
sarama
.
OffsetNewest
config
.
Version
=
sarama
.
V0_10_2_
1
config
.
Version
=
sarama
.
V0_10_2_
0
consumerGroup
,
err
:=
sarama
.
NewConsumerGroup
(
r
.
msgConsumer
.
kafkaHosts
,
r
.
msgConsumer
.
groupId
,
config
)
if
err
!=
nil
{
return
err
...
...
pkg/port/consumer/produce/produce.go
查看文件 @
8668ebb
...
...
@@ -21,7 +21,7 @@ func init() {
var
err
error
mqConfig
:=
sarama
.
NewConfig
()
mqConfig
.
Producer
.
Return
.
Successes
=
true
mqConfig
.
Version
=
sarama
.
V0_10_2_
1
mqConfig
.
Version
=
sarama
.
V0_10_2_
0
if
err
=
mqConfig
.
Validate
();
err
!=
nil
{
msg
:=
fmt
.
Sprintf
(
"Kafka producer config invalidate. config: %v. err: %v"
,
configs
.
Cfg
,
err
)
logs
.
Info
(
msg
)
...
...
请
注册
或
登录
后发表评论