Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
DiDi
kafka-manager
提交
a5fa9de5
K
kafka-manager
项目概览
DiDi
/
kafka-manager
9 个月 前同步成功
通知
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 搜索 >>
提交
a5fa9de5
编写于
9月 28, 2022
作者:
Z
zengqiao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
修复Group指标防重复不生效问题
上级
1e256ae1
变更
1
隐藏空白更改
内联
并排
Showing
1 changed file
with
17 addition
and
9 deletion
+17
-9
km-core/src/main/java/com/xiaojukeji/know/streaming/km/core/service/group/impl/GroupMetricServiceImpl.java
...ng/km/core/service/group/impl/GroupMetricServiceImpl.java
+17
-9
未找到文件。
km-core/src/main/java/com/xiaojukeji/know/streaming/km/core/service/group/impl/GroupMetricServiceImpl.java
浏览文件 @
a5fa9de5
...
...
@@ -90,23 +90,31 @@ public class GroupMetricServiceImpl extends BaseMetricService implements GroupMe
@Override
public
Result
<
List
<
GroupMetrics
>>
collectGroupMetricsFromKafka
(
Long
clusterId
,
String
groupName
,
List
<
String
>
metrics
)
{
List
<
GroupMetrics
>
allGroupMetrics
=
new
ArrayList
<>();
Map
<
String
,
GroupMetrics
>
topicPartitionGroupMap
=
new
HashMap
<>();
List
<
GroupMetrics
>
allGroupMetrics
=
new
ArrayList
<>();
Map
<
String
,
GroupMetrics
>
topicPartitionGroupMap
=
new
HashMap
<>();
GroupMetrics
groupMetrics
=
new
GroupMetrics
(
clusterId
,
groupName
,
true
);
for
(
String
metric
:
metrics
){
if
(
null
!=
groupMetrics
.
getMetrics
().
get
(
metric
)){
continue
;}
Set
<
String
>
existMetricSet
=
new
HashSet
<>();
for
(
String
metric
:
metrics
)
{
if
(
existMetricSet
.
contains
(
metric
))
{
continue
;
}
Result
<
List
<
GroupMetrics
>>
ret
=
collectGroupMetricsFromKafka
(
clusterId
,
groupName
,
metric
);
if
(
null
!=
ret
&&
ret
.
successful
())
{
if
(
null
!=
ret
&&
ret
.
successful
())
{
List
<
GroupMetrics
>
groupMetricsList
=
ret
.
getData
();
for
(
GroupMetrics
gm
:
groupMetricsList
){
if
(
gm
.
isBGroupMetric
()){
for
(
GroupMetrics
gm
:
groupMetricsList
)
{
//记录已存在的指标
existMetricSet
.
addAll
(
gm
.
getMetrics
().
keySet
());
if
(
gm
.
isBGroupMetric
())
{
groupMetrics
.
getMetrics
().
putAll
(
gm
.
getMetrics
());
}
else
{
}
else
{
GroupMetrics
topicGroupMetric
=
topicPartitionGroupMap
.
getOrDefault
(
gm
.
getTopic
()
+
gm
.
getPartitionId
(),
new
GroupMetrics
(
clusterId
,
groupName
,
false
));
new
GroupMetrics
(
clusterId
,
gm
.
getPartitionId
(),
gm
.
getTopic
()
,
groupName
,
false
));
topicGroupMetric
.
getMetrics
().
putAll
(
gm
.
getMetrics
());
topicPartitionGroupMap
.
put
(
gm
.
getTopic
()
+
gm
.
getPartitionId
(),
topicGroupMetric
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录