前言
现所属公司使用了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~
