站点图标 梦呓

Kafka图形化工具kafka-ui介绍

前言

现所属公司使用了k8s来作为生产环境,使用了knative来作为服务基础构建,底层服务与服务之间的通信全由knative-eventing进行异步传输,broker底层使用的kafka作为消息队列存储。那么服务与服务之间的通信目前是处于黑盒状态,A到B的消息A是否发送成功了,kafka是否收到了,以及B是否消费了,这都是黑盒的。又因为前几天因为某家云服务器厂商因网卡硬件出现故障宕机,导致Kafka中间出现故障,服务与服务之间的通信时好时坏,重启服务、复盘代码也得不到解决,又不知道kafka消息的中间状态,就很抓瞎。之前微服务项目也使用kafka,但使用的是kafkaTool,在轻量场景下绰绰有余,但knative-eventing的重度场景下就捉襟见肘了,鉴于此,给大家介绍下Kafka-ui,支持按时间,key/value值进行搜索等好用的功能。

注:本教程Kafka是在k8s里部署的,所以相关工具也是通过k8s部署或者k8s方式使用的,kafka裸机部署及使用方法大同小异,感兴趣的同学可以自行AI。

OffsetExplore(kafkaTool)

这个工具之前名字叫kafkaTool,后来改版了叫OffsetExplore新的名字,电脑本地安装,可以支持基础的Kafka消息查询等。

安装

搜索官网找到Download下载对应的版本直接安装即可。

本地环境配置

核心原理是把k8s集群中的Kafka通过端口转发出来,本地offsetExplore进行连接的。

添加本地host

找到本机host文件,添加IP域名映射。

127.0.0.1   kafka-service-0.kafka
127.0.0.2   kafka-service-1.kafka
127.0.0.3   kafka-service-2.kafka

因为k8s集群里使用了kafka三节点部署,所以映射为三个,Kafka域名也可以根据k8s集群里实际对外暴漏的服务名进行修改。如果不知道集群里具体的域名可以通过命令kubectl get svc -n kafka进行查看:

host添加完IP域名映射后,需要执行ipconfig /flushdns进行刷新缓存,因为offsetExplore打开时会记录旧的DNS域名,所以如果是先打开offsetExplore后添加的host的IP域名映射,可能会不成功,这时候就要关闭软件重新刷新下DNS缓存即可。

端口转发

在本机powershell命令行里切换到指定kubeconfig环境,然后进行端口转发即可。

1、切换kubeconfig环境。
$env:KUBECONFIG = "D:\XXX\k8s配置文件\正式k8s配置cls-2n2-config-new"
2、可以通过此命令再次查看env是否切换成功
kubectl config get-contexts
3、进行端口转发,$ns="kafka"为k8s集群里Kafka部署的命名空间名称。
$ns="kafka"; cmd /c start /min kubectl port-forward --address 127.0.0.1 pod/kafka-0-0 9092:9092 -n $ns ; cmd /c start /min kubectl port-forward --address 127.0.0.2 pod/kafka-1-0 9092:9092 -n $ns ; cmd /c start /min kubectl port-forward --address 127.0.0.3 pod/kafka-2-0 9092:9092 -n $ns

OffsetExplore连接

打开本地OffsetExplore,填入Kafka集群地址即可连接。

功能使用

Brokers

连接到集群中的Kafka之后点击Brokers可以查看一共有几个节点。

Topics

点击topics可以查看Kafka里一共有多少个topics在使用,点开其中一个可以查看topic有几个partitions,选中某个partition可以直接查看具体partition里面的消息,或者直接点topic是直接查该topic下的所有partition的消息。

其中有个重要的就是选中该topic,在右侧特性->内容类型菜单里把钥匙、价值的内容类型改成细绳,这里是软件官方翻译错误,正确翻译应该是key、value、字符串。意思是如果不设置,等查出来的消息是十六进制的,不方便查看,这里改成字符串方便查看。

可以直接点留言->绿色播放按钮,查看按照最新或者最老排序的消息,右下角可以查看按照时间排序的多少条数据,一般都是按照最新时间排序,查询最近1000条,也可以通过筛选在列表里查询服务发送给kafka的消息是否成功,然后把消息体直接粘出来或者在下方查看发送的消息是否是正确的。

这样进行故障排除或者DEBUG调试绰绰有余了,但是在排查历史问题时,需要通过时间来查询,offsetExplore则没支持,所以给大家引入了kafka-ui来进行查询,且offsetExplore是装在本地,需要研发每个人都装一个,而kafka-ui是装在k8s集群,通过外网域名进行暴漏,只需要安装一次,大家就可通过账号密码进行登录访问了,更方便些。

Consumers

在消费者菜单里,找到具体的消费者,点击offsets可以查看当前消费者都消费了哪个topic,哪些partition,消费的offset到哪里了,以及更重要的有没有滞后,这个指标可以确认是否存在未消费信息,也能间接的证明是否发送给了下游服务。

Kafka Assistant

继OffsetExplore之后,Kafka Assistant也是一个本地工具,其优点是比OffsetExplore提供了更丰富的指标,也支持消息的时间点查询。如果不介意它需要付费使用以及需要装到本地的特点话,倒是值得一用。大家也可以到B站观看更详细的视频介绍。

同系列产品拓展

发现这个厂家同系列也有一些有意思的产品,比方说MQTT、Modbus工具在工业物联网看起来挺好用的(因为之前用过其它Modbus模拟器,是真不好用,如果再有机会一定要试下),如果大家有需要的话可以去尝试下。

核心功能

连接Kafka

查看健康指标

支持丰富数据格式

消费和发送消息

查看和更新Kafka配置

生成拓扑图

消息绘制图表

Kafka-ui

如果大家想要一个免费又好用,又不想装在本地的工具,那就非Kafka-ui莫属了。

安装

我是在k8s集群里通过helm安装的,因为安装、修改脚本过程要保留可以复用,所以先把官方helm脚本下载到本地,然后再到对应集群环境进行安装,而不是用命令直接使用官方helm脚本安装,这样脚本就可以做到保留并提交git留存。

1、helm脚本下载
helm repo add kafka-ui https://provectus.github.io/kafka-ui-charts ; 
helm repo update ; 
helm pull kafka-ui/kafka-ui --untar --destination ./kafka-ui-local
2、本地切换到指定k8s环境kubeconfig,本地直接安装
helm install kafka-ui ./kafka-ui-local/kafka-ui -f ./kafka-ui-local/kafka-ui/values-prod.yaml -n kafka-ui --create-namespace

Ingress创建

因为是安装到集群里,想能大家都通过浏览器进行访问,就需要对外(内网)暴露域名,添加ingress。

apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
  name: kafka-ui-ingress
  namespace: kafka-ui
    # annotations:
    # 根据你集群实际的 Ingress 控制器,可能需要取消注释并添加特定注解
    # 例如,强制使用 nginx 并开启一定大小的文件上传限制:
  # nginx.ingress.kubernetes.io/proxy-body-size: "50m"
spec:
  # 这里假设你的集群使用的是 Nginx Ingress Controller
  # 如果使用的是 Traefik 或其他,请将其修改为对应的 ingressClass 名称
  ingressClassName: nginx
  rules:
    - host: abc.com
      http:
        paths:
          - path: /
            pathType: Prefix
            backend:
              service:
                name: kafka-ui
                port:
                  number: 80

设置访问账号密码及集群连接配置

kafka-ui默认是没有账号密码的,通过URL直接就可以访问,这是不行的,就需要在values.yaml文件里添加账号密码,就会有一个登录页出来。

如果有多个集群安装了kafka,比如说测试集群、正式集群,建议在每个集群里都安装kafka-ui,虽然他也能跨集群连接,但是会非常麻烦,需要网络打通等,还有安全隐患。连接同集群Kafka则非常简单,填写Kafka服务名即可,不知道获取服务名的上面教程有讲到或自行AI。

replicaCount: 1


image:
  registry: docker.io
  repository: provectuslabs/kafka-ui
  pullPolicy: IfNotPresent
  # Overrides the image tag whose default is the chart appVersion.
  tag: ""

imagePullSecrets: []
nameOverride: ""
fullnameOverride: ""

serviceAccount:
  # Specifies whether a service account should be created
  create: true
  # Annotations to add to the service account
  annotations: {}
  # The name of the service account to use.
  # If not set and create is true, a name is generated using the fullname template
  name: ""

existingConfigMap: ""
yamlApplicationConfig:
  # ---------------- 1. 开启账号密码访问 ----------------
  auth:
    type: LOGIN_FORM
  spring:
    security:
      user:
        name: abc           # 你想要的登录账号
        password: acb123 # 你想要的登录密码

  # ---------------- 2. 配置多集群环境 ----------------
  kafka:
    clusters:
      # 集群 A:你原先配置的测试环境集群(保留你原来的配置)
#      - name: test-cluster
#        bootstrapServers: "之前的地址:9092"

      # 集群 B:新增你 K8s 内部 kafka 命名空间下的集群
      - name: pre-kafka-cluster
        # K8s 的 DNS 非常聪明,域名里的 .kafka 会直接被解析为去 kafka 命名空间找服务。
        # 这刚好完美匹配了我们之前排查出来的,它内部真实注册的 advertised 名字。
        bootstrapServers: "kafka-service-0.kafka:9092,kafka-service-1.kafka:9092,kafka-service-2.kafka:9092"
  # kafka:
  #   clusters:
  #     - name: yaml
  #       bootstrapServers: kafka-service:9092
  # spring:
  #   security:
  #     oauth2:
  # auth:
  #   type: disabled
  # management:
  #   health:
  #     ldap:
  #       enabled: false
yamlApplicationConfigConfigMap:
  {}
  # keyName: config.yml
  # name: configMapName
existingSecret: ""
envs:
  secret: {}
  config: {}

networkPolicy:
  enabled: false
  egressRules:
    ## Additional custom egress rules
    ## e.g:
    ## customRules:
    ##   - to:
    ##       - namespaceSelector:
    ##           matchLabels:
    ##             label: example
    customRules: []
  ingressRules:
    ## Additional custom ingress rules
    ## e.g:
    ## customRules:
    ##   - from:
    ##       - namespaceSelector:
    ##           matchLabels:
    ##             label: example
    customRules: []

podAnnotations: {}
podLabels: {}

## Annotations to be added to kafka-ui Deployment
##
annotations: {}

## Set field schema as HTTPS for readines and liveness probe
##
probes:
  useHttpsScheme: false

podSecurityContext:
  {}
  # fsGroup: 2000

securityContext:
  {}
  # capabilities:
  #   drop:
  #   - ALL
  # readOnlyRootFilesystem: true
  # runAsNonRoot: true
  # runAsUser: 1000

service:
  type: ClusterIP
  port: 80
  # In case of service type LoadBalancer, you can specify reserved static IP
  # loadBalancerIP: 10.11.12.13
  # if you want to force a specific nodePort. Must be use with service.type=NodePort
  # nodePort:

# Ingress configuration
ingress:
  # Enable ingress resource
  enabled: false

  # Annotations for the Ingress
  annotations: {}

  # ingressClassName for the Ingress
  ingressClassName: ""

  # The path for the Ingress
  path: "/"

  # The path type for the Ingress
  pathType: "Prefix"  

  # The hostname for the Ingress
  host: ""

  # configs for Ingress TLS
  tls:
    # Enable TLS termination for the Ingress
    enabled: false
    # the name of a pre-created Secret containing a TLS private key and certificate
    secretName: ""

  # HTTP paths to add to the Ingress before the default path
  precedingPaths: []

  # Http paths to add to the Ingress after the default path
  succeedingPaths: []

resources:
  {}
  # limits:
  #   cpu: 200m
  #   memory: 512Mi
  # requests:
  #   cpu: 200m
  #   memory: 256Mi

autoscaling:
  enabled: false
  minReplicas: 1
  maxReplicas: 100
  targetCPUUtilizationPercentage: 80
  # targetMemoryUtilizationPercentage: 80

nodeSelector: {}

tolerations: []

affinity: {}

env: {}

initContainers: {}

volumeMounts: {}

volumes: {}

功能使用

Brokers

点击brokers可以查看三个节点的使用情况、存储情况、分区倾斜等信息。

同样也可以看到Kafka日志的存储目录显示。

可视化更改Kafka的配置。

Topics

在topics页面可以查看更详细的信息。

概述页面。进入某topics可以查看该topics更详细的信息,也可以直接创建消息。

消息页面。可以根据offset偏移量查询、也可以根据时间点查询、可以选择查询哪个分区里的消息、key/value的类型、以及按照最新、最老排序,也支持强大的实时模式,在DEBUG场景下特别好用。

点击下边某条消息则可以查看具体值。

在消费者页面可以查看都有哪些消费者订阅了这个topic,以及各消费者的消费情况。

设置页可以查看topics的配置信息。

统计数据页面可以查看当前topics的状态信息。

Consumers

在消费者页面,点击具体的消费者,可以查看该消费者在每个partition的消费情况。

小结

综上所述,通过不同Kafka图形化工具,可以在排查问题的时候使得服务与服务的中间消息链路不再黑盒,更快的定位问题,祝愿天下程序员们少加班,早回家,无BUG~

退出移动版