Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
DiDi
kafka-manager
提交
db40a5cd
K
kafka-manager
项目概览
DiDi
/
kafka-manager
8 个月 前同步成功
通知
58
Star
6372
Fork
1229
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
K
kafka-manager
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
未验证
提交
db40a5cd
编写于
8月 02, 2023
作者:
E
EricZeng
提交者:
GitHub
8月 02, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[Optimize]增加AdminClient观测信息 (#1111)
1、增加AdminClient的ClientID; 2、关闭时,增加超时时间; 3、增加关闭错误的日志;
上级
55161e43
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
12 addition
and
4 deletion
+12
-4
km-core/src/main/java/com/xiaojukeji/know/streaming/km/core/service/group/impl/GroupServiceImpl.java
...treaming/km/core/service/group/impl/GroupServiceImpl.java
+2
-0
km-persistence/src/main/java/com/xiaojukeji/know/streaming/km/persistence/kafka/KafkaAdminClient.java
...know/streaming/km/persistence/kafka/KafkaAdminClient.java
+10
-4
未找到文件。
km-core/src/main/java/com/xiaojukeji/know/streaming/km/core/service/group/impl/GroupServiceImpl.java
浏览文件 @
db40a5cd
...
...
@@ -78,6 +78,7 @@ public class GroupServiceImpl extends BaseKafkaVersionControlService implements
}
props
.
put
(
AdminClientConfig
.
BOOTSTRAP_SERVERS_CONFIG
,
clusterPhy
.
getBootstrapServers
());
props
.
put
(
AdminClientConfig
.
CLIENT_ID_CONFIG
,
String
.
format
(
"KSPartialAdminClient||clusterPhyId=%d"
,
clusterPhy
.
getId
()));
adminClient
=
KSPartialKafkaAdminClient
.
create
(
props
);
KSListGroupsResult
listConsumerGroupsResult
=
adminClient
.
listConsumerGroups
(
...
...
@@ -178,6 +179,7 @@ public class GroupServiceImpl extends BaseKafkaVersionControlService implements
}
props
.
put
(
AdminClientConfig
.
BOOTSTRAP_SERVERS_CONFIG
,
clusterPhy
.
getBootstrapServers
());
props
.
put
(
AdminClientConfig
.
CLIENT_ID_CONFIG
,
String
.
format
(
"KSPartialAdminClient||clusterPhyId=%d"
,
clusterPhy
.
getId
()));
adminClient
=
KSPartialKafkaAdminClient
.
create
(
props
);
...
...
km-persistence/src/main/java/com/xiaojukeji/know/streaming/km/persistence/kafka/KafkaAdminClient.java
浏览文件 @
db40a5cd
...
...
@@ -12,6 +12,7 @@ import org.apache.kafka.clients.admin.AdminClientConfig;
import
org.springframework.beans.factory.annotation.Value
;
import
org.springframework.stereotype.Component
;
import
java.time.Duration
;
import
java.util.ArrayList
;
import
java.util.List
;
import
java.util.Map
;
...
...
@@ -76,10 +77,12 @@ public class KafkaAdminClient extends AbstractClusterLoadedChangedHandler {
LOGGER
.
info
(
"close kafka AdminClient starting, clusterPhyId:{}"
,
clusterPhyId
);
boolean
allSuccess
=
this
.
closeAdminClientList
(
adminClientList
);
boolean
allSuccess
=
this
.
closeAdminClientList
(
clusterPhyId
,
adminClientList
);
if
(
allSuccess
)
{
LOGGER
.
info
(
"close kafka AdminClient success, clusterPhyId:{}"
,
clusterPhyId
);
}
else
{
LOGGER
.
error
(
"close kafka AdminClient exist failed and can ignore this error, clusterPhyId:{}"
,
clusterPhyId
);
}
}
catch
(
Exception
e
)
{
LOGGER
.
error
(
"close kafka AdminClient failed, clusterPhyId:{}"
,
clusterPhyId
,
e
);
...
...
@@ -116,6 +119,7 @@ public class KafkaAdminClient extends AbstractClusterLoadedChangedHandler {
adminClientList
=
new
ArrayList
<>();
for
(
int
i
=
0
;
i
<
clientCnt
;
++
i
)
{
props
.
put
(
AdminClientConfig
.
CLIENT_ID_CONFIG
,
String
.
format
(
"ApacheAdminClient||clusterPhyId=%d||Cnt=%d"
,
clusterPhyId
,
i
));
adminClientList
.
add
(
AdminClient
.
create
(
props
));
}
...
...
@@ -125,7 +129,7 @@ public class KafkaAdminClient extends AbstractClusterLoadedChangedHandler {
}
catch
(
Exception
e
)
{
LOGGER
.
error
(
"create kafka AdminClient failed, clusterPhyId:{} props:{}"
,
clusterPhyId
,
props
,
e
);
this
.
closeAdminClientList
(
adminClientList
);
this
.
closeAdminClientList
(
clusterPhyId
,
adminClientList
);
}
finally
{
modifyClientMapLock
.
unlock
();
}
...
...
@@ -133,7 +137,7 @@ public class KafkaAdminClient extends AbstractClusterLoadedChangedHandler {
return
KAFKA_ADMIN_CLIENT_MAP
.
get
(
clusterPhyId
).
get
((
int
)(
System
.
currentTimeMillis
()
%
clientCnt
));
}
private
boolean
closeAdminClientList
(
List
<
AdminClient
>
adminClientList
)
{
private
boolean
closeAdminClientList
(
L
ong
clusterPhyId
,
L
ist
<
AdminClient
>
adminClientList
)
{
if
(
adminClientList
==
null
)
{
return
true
;
}
...
...
@@ -141,9 +145,11 @@ public class KafkaAdminClient extends AbstractClusterLoadedChangedHandler {
boolean
allSuccess
=
true
;
for
(
AdminClient
adminClient:
adminClientList
)
{
try
{
adminClient
.
close
();
// 关闭客户端,超时时间为30秒
adminClient
.
close
(
Duration
.
ofSeconds
(
30
));
}
catch
(
Exception
e
)
{
// ignore
LOGGER
.
error
(
"close kafka AdminClient exist failed, clusterPhyId:{}"
,
clusterPhyId
,
e
);
allSuccess
=
false
;
}
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录