拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

EFAK Kafka可视化管理:从命令行运维到实时数据流洞察

EFAK Kafka可视化管理:从命令行运维到实时数据流洞察

1. 为什么今天还要认真对待 Kafka 可视化管理——从“查不到消息”到“一眼看清数据流”的真实转变

EFAK(Elastic Kafka Manager),也就是大家更熟悉的 Kafka Eagle,不是又一个花哨的监控面板。它是我过去三年在金融、电商和物联网三条业务线里,反复验证过唯一能真正解决“Kafka 运维盲区”的开源工具。很多人第一次接触它,是因为在 CDH 集群上跑 KSQL 任务时突然卡住,日志里只有一行Failed to commit offset,却根本不知道是哪个 consumer group 滞后了 200 万条,也不知道 lag 是均匀分布还是集中在某个 partition。这时候打开 EFAK 的 Consumer Lag 页面,3 秒内就能定位到order-processor-v3这个 group 在topic_order_events-7上 lag 达到 198 万——而其他 15 个 partition 都是 0。这种“所见即所得”的能力,不是靠猜,也不是靠写一堆 shell 脚本轮询kafka-consumer-groups.sh,而是 EFAK 把 Kafka 原生协议层的元数据、JMX 指标、ZooKeeper 状态、以及消费位点的真实计算逻辑,全部做了结构化聚合与可视化映射。

它和 Apache APISIX、Apache JMeter 这类工具的本质区别在于:APISIX 是网关,JMeter 是压测,它们不碰 Kafka 的内部状态;而 EFAK 是 Kafka 的“X 光机”——它不转发流量,不生成负载,只做一件事:把 Kafka 集群里那些藏在命令行、埋在 JMX、散落在 ZooKeeper 节点里的碎片信息,拼成一张可交互、可下钻、可告警的实时拓扑图。比如你看到某个 topic 的 ISR 数量从 3 掉到 2,EFAK 不会只显示一个红色感叹号,它会直接关联到 broker 列表页,告诉你 broker.id=5 的磁盘使用率已达 94%,且该 broker 正在执行 unclean leader election;再点开这个 broker 的 JVM 监控,发现 GC 时间已连续 5 分钟超过 2s。这才是真正的“问题链路穿透”,而不是孤立指标报警。

所以这不是一个“装了就能用”的玩具型 UI。它的价值恰恰体现在你已经熟悉 Kafka 原生命令、能手写 KSQL、甚至自己搭过 Prometheus+Grafana 监控体系之后——当你发现 Grafana 里看到的kafka_server_BrokerTopicMetrics_OneMinuteRate曲线异常,却无法快速判断是 producer 发送失败、还是 consumer 拉取超时、抑或是 broker 内部线程池阻塞时,EFAK 就成了那个能帮你“切片归因”的手术刀。它不替代你的运维知识,而是把你已有的 Kafka 认知,转化成可操作的界面动作:点击 lag 值跳转到对应 partition 的 message preview,拖动时间轴查看历史 offset 变化,右键 topic 触发自动 rebalance,甚至一键导出当前所有 consumer group 的完整 offset map 用于灾备比对。这些功能背后,是它对 Kafka 0.10.x 至 3.6.x 全版本协议兼容的扎实实现,也是它在 CDH 6.3.2、Cloudera Runtime 7.2.10 等企业级发行版中完成深度适配的结果。如果你还在用kafka-topics.sh --describe查 partition 分布,用kafka-consumer-groups.sh --group xxx --describe查 offset,那你不是在管理 Kafka,你是在翻译 Kafka 的二进制语言。EFAK 的存在,就是让 Kafka 管理回归“人话”。

2. EFAK 的底层架构不是黑盒——它如何绕过 Kafka AdminClient 的限制获取真实状态

很多团队在评估 EFAK 时第一个疑问是:“它是不是只是把 Kafka 命令行包装了一层 Web?”答案是否定的。EFAK 的核心能力来源于它对 Kafka 协议栈的三重穿透:元数据层、运行时层、存储层。这决定了它为什么能在不依赖 Kafka Broker 开放额外端口、不修改任何 Kafka 配置的前提下,获取到比 AdminClient 更全、更准、更及时的状态信息。

首先看元数据层。Kafka AdminClient 默认只能通过describeTopics()、listConsumerGroups()等 API 获取快照式数据,且受max.in.flight.requests.per.connection=1等客户端参数影响,高并发调用时容易超时或返回不一致结果。EFAK 则采用双通道策略:一方面复用 AdminClient 获取 topic schema、partition count、replication factor 等静态元数据;另一方面,它会主动连接 ZooKeeper(或 KRaft 模式下的 metadata log),读取/brokers/ids、/consumers/<group>/offsets等原始节点。注意,这里不是简单地get一个 znode,而是监听WATCHER事件——当某个 broker 下线时,EFAK 能在 ZooKeeper 的NodeDeleted事件触发后 200ms 内更新 UI 状态,远快于 AdminClient 的默认 30s 心跳检测周期。我在某次生产环境演练中实测过:手动 kill broker.id=3 后,EFAK 的 broker 列表页在 1.2 秒内变灰并显示 “Disconnected”,而kafka-broker-api-versions.sh命令仍需等待 27 秒才报错超时。

其次是运行时层。AdminClient 对 consumer lag 的计算是近似值:它调用listConsumerGroupOffsets()获取每个 partition 的 committed offset,再用endOffsets()获取 log end offset,两者相减得出 lag。但endOffsets()本身有缓存机制,且在高吞吐场景下可能返回 stale 数据。EFAK 则绕过了这个瓶颈——它直接解析 broker 的__consumer_offsetstopic。这个 topic 存储了所有 consumer group 的 offset 提交记录,EFAK 启动一个专用的 internal consumer,以earliest策略订阅__consumer_offsets,持续拉取新提交的 offset record,并结合本地维护的 partition 分配映射表,实时计算每个 group 的精确 lag。这意味着即使某个 group 已停止消费,只要它最近一次提交过 offset,EFAK 就能算出其 lag;而 AdminClient 在 group 无活跃成员时,listConsumerGroupOffsets()会直接抛出UNKNOWN_MEMBER_ID异常,导致 lag 显示为 0(错误)。

最后是存储层。这是 EFAK 最被低估的能力。它内置了一个轻量级嵌入式数据库(默认 H2,可切换 MySQL/PostgreSQL),专门用于持久化三类关键数据:一是 topic 的 schema evolution 历史(每次kafka-topics.sh --alter --add-config操作都会被捕获并存档);二是 consumer group 的 offset 变化轨迹(每 5 分钟采样一次,形成时间序列);三是 alert rule 的触发日志(比如 “lag > 100000 持续 3 分钟” 的完整上下文)。这些数据不是为了炫技,而是支撑了 EFAK 的核心功能:比如 “Compare Offset History” 功能,你可以选择两个时间点(如故障前 1 小时 vs 故障后 5 分钟),对比同一个 group 在所有 partition 上的 offset 差值,从而精准定位是哪个 partition 出现了消费停滞;再比如 “Replay Message” 功能,它不是简单地kafka-console-consumer.sh --from-beginning,而是先从 H2 中查出该 message 的物理位置(segment file + offset),再调用 broker 的FetchRequest协议精准拉取,避免全量扫描。

提示:EFAK 的 ZooKeeper 依赖仅限于元数据同步,它不依赖 ZooKeeper 执行任何写操作。因此在 Kafka 3.3+ 的 KRaft 模式集群中,只需关闭efak.zk.enable=false并配置efak.metadata.bootstrap.servers=broker1:9092,broker2:9092,即可完全脱离 ZooKeeper 运行。这一点常被误读为“EFAK 不支持 KRaft”,实际恰恰相反——它的架构设计天然适配 Kafka 的演进方向。

3. 保姆级安装不是“解压启动”——CDH 环境下的 7 个必须校准环节

在 CDH 环境部署 EFAK,最大的陷阱不是“装不上”,而是“装上了但看不到数据”。我见过太多团队在 Cloudera Manager 里启用了 Kafka 服务、配置了 JMX、开放了 9092 端口,却在 EFAK UI 上看到一片空白的 topic 列表。问题往往不出在 EFAK 本身,而在于 CDH 对 Kafka 组件的封装方式与 EFAK 的探针逻辑之间存在 7 处隐性断点。下面按执行顺序,逐个拆解这些必须手动校准的环节:

3.1 确认 Kafka Broker 的 JMX 端口暴露策略

CDH 默认将 Kafka 的 JMX 服务绑定在127.0.0.1:9999,这是一个典型的“本地回环绑定”。EFAK 的监控模块需要远程连接此端口获取kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec等指标,但127.0.0.1对外部不可达。解决方案不是简单改 bind address,而是利用 CDH 的安全机制:进入 Cloudera Manager → Kafka Service → Configuration → Filter “jmx” → 找到KAFKA_JMX_OPTS参数,在其值末尾追加-Dcom.sun.management.jmxremote.host=<broker-hostname>和-Dcom.sun.management.jmxremote.port=9999。注意<broker-hostname>必须是集群内其他节点能 DNS 解析的 FQDN(如kafka-broker01.prod.example.com),不能填 IP 或localhost。实测发现,若此处填0.0.0.0,CDH 的 Java 安全策略会拦截 JMX RMI 连接,导致 EFAK 报错java.rmi.ConnectException: Connection refused to host: 0.0.0.0。

3.2 修正 Kafka 的 advertised.listeners 配置

这是最隐蔽也最致命的一环。CDH 的 Kafka 配置中,advertised.listeners默认值常为PLAINTEXT://$HOSTNAME:9092。EFAK 的 internal consumer 需要根据此配置构建 bootstrap servers 列表来拉取消息,但$HOSTNAME是 CDH agent 的主机名,不一定能被 EFAK 所在服务器解析。例如,CDH broker 主机名为cdh-kafka-01.internal,而 EFAK 部署在monitoring-prod-03服务器上,后者 DNS 中并无cdh-kafka-01.internal记录。此时 EFAK 会尝试连接cdh-kafka-01.internal:9092并超时。正确做法是:在 Cloudera Manager → Kafka Service → Configuration → Filter “advertised” → 修改advertised.listeners为PLAINTEXT://cdh-kafka-01.example.com:9092(使用全局可解析的域名),并确保该域名在 EFAK 服务器/etc/hosts中有静态映射。我们曾因此问题排查了 17 小时,最终发现是 DNS 服务器未同步这条 A 记录。

3.3 配置 EFAK 的 ZooKeeper 连接字符串(CDH 6.3+ 特别注意)

CDH 6.3 及以后版本默认启用 Kerberos 认证的 ZooKeeper。EFAK 的efak.zk.connect参数若只填zookeeper01:2181,zookeeper02:2181,会因认证失败而无法读取/brokers/ids。必须启用 SASL 认证:在efak.properties中添加:

efak.zk.connect=zookeeper01:2181,zookeeper02:2181,zookeeper03:2181 efak.zk.sasl.enable=true efak.zk.sasl.kerberos.principal=zkclient@EXAMPLE.COM efak.zk.sasl.kerberos.keytab=/etc/security/keytabs/zkclient.keytab

其中zkclient@EXAMPLE.COM是 CDH 自动创建的 ZooKeeper 客户端 principal,keytab 文件路径需从 Cloudera Manager 的 ZooKeeper 服务配置页中复制。漏掉这一项,EFAK 将无法发现任何 broker,topic 列表永远为空。

3.4 调整 EFAK 的 JVM 内存参数以匹配 CDH 集群规模

EFAK 默认启动脚本bin/startup.sh设置-Xms512m -Xmx1g,这在单 broker 测试环境可行,但在 50+ broker 的 CDH 生产集群中会频繁 OOM。原因在于 EFAK 的 metadata cache 会为每个 topic 的每个 partition 创建独立对象,一个含 200 个 partition 的 topic 就占用约 12MB 堆内存。我们线上集群有 387 个 topic,平均 partition 数 42,粗略估算 metadata 对象需 200MB+。建议按公式调整:-Xms = (broker_count * 10 + topic_count * 5) MB,-Xmx = Xms * 1.5。对于 60 broker + 400 topic 的集群,应设为-Xms3.2g -Xmx4.8g。同时添加-XX:+UseG1GC -XX:MaxGCPauseMillis=200,避免 CMS GC 导致 UI 响应延迟。

3.5 启用 EFAK 的 KSQL 集成模块(非默认开启)

EFAK 的 KSQL 支持不是开箱即用的。CDH 的 KSQL Server 默认监听http://ksql-server-01:8088,但 EFAK 需要额外配置才能连接。在efak.properties中必须显式声明:

efak.ksql.enable=true efak.ksql.servers=http://ksql-server-01:8088,http://ksql-server-02:8088 efak.ksql.timeout=30000

且 KSQL Server 的listeners配置必须包含http://0.0.0.0:8088(而非http://127.0.0.1:8088),否则 EFAK 无法访问。此外,CDH 的 KSQL Server 默认关闭了 REST API 的 CORS 支持,需在 KSQL Server 的ksql-server.properties中添加ksql.rest.api.cors.origins=*(生产环境建议限定为 EFAK 的域名)。

3.6 配置 EFAK 的 LDAP/AD 认证对接 CDH 的 Sentry 权限体系

CDH 集群通常使用 Sentry 或 Ranger 做 Kafka ACL 管理。EFAK 本身不接管权限,但可通过 LDAP 同步用户组,再映射到 Sentry 的 role。在efak.properties中:

efak.security.auth.type=ldap efak.ldap.url=ldaps://ad.example.com:636 efak.ldap.base.dn=OU=KafkaUsers,DC=example,DC=com efak.ldap.user.dn.pattern=uid={0},OU=KafkaUsers,DC=example,DC=com efak.ldap.group.search.base=OU=KafkaGroups,DC=example,DC=com efak.ldap.group.search.filter=(member={0})

然后在 EFAK 的 Web UI → Settings → Role Mapping 中,将 AD 组kafka-admins映射到 EFAK 的ADMIN角色,kafka-developers映射到USER角色。这样用户登录后,EFAK 会自动调用 Sentry 的listPermissionsAPI,过滤其可见的 topic 列表,实现权限继承。

3.7 验证 EFAK 的 Metrics Collector 是否与 Cloudera Manager 的 Agent 共存

CDH 的 Metrics Collector(由 Cloudera Management Service 提供)默认监听7184端口。EFAK 的efak.metrics.collector.enable=true时,也会尝试启动内置 collector,若端口冲突会导致 EFAK 启动失败。解决方案是:在efak.properties中显式禁用内置 collector,改用 CDH 的标准指标源:

efak.metrics.collector.enable=false efak.metrics.source.type=cloudera efak.cloudera.cm.host=cm-server.example.com efak.cloudera.cm.port=7183 efak.cloudera.cm.username=admin efak.cloudera.cm.password=your_password efak.cloudera.cm.cluster.name=ProductionCluster

这样 EFAK 就能直接读取 Cloudera Manager 的 Kafka 指标(如kafka_broker_messages_in_total_rate),无需重复采集。

4. 从零开始的实操安装流程——基于 CDH 6.3.2 的完整命令链

现在我们把上述 7 个校准点,转化为一份可直接执行、带解释说明的安装清单。以下所有命令均在 EFAK 部署服务器(假设为efak-prod-01.example.com)上执行,操作系统为 CentOS 7.9,CDH 版本为 6.3.2,Kafka 版本为 2.3.0。

4.1 下载与解压 EFAK 发行包

EFAK 官方 GitHub Release 页面(https://github.com/kefeng-wang/EFAK/releases)提供预编译包。截至 2024 年,推荐使用efak-3.0.1-bin.tar.gz(兼容 Kafka 2.0+,修复了 CDH 6.3 的 Kerberos 兼容问题):

# 创建部署目录 sudo mkdir -p /opt/efak sudo chown kafka:kafka /opt/efak # 下载(使用国内镜像加速) curl -L https://ghproxy.com/https://github.com/kefeng-wang/EFAK/releases/download/v3.0.1/efak-3.0.1-bin.tar.gz \ -o /tmp/efak-3.0.1-bin.tar.gz # 校验 SHA256(官方发布页提供) echo "a1b2c3d4e5f6... /tmp/efak-3.0.1-bin.tar.gz" | sha256sum -c - # 解压并设置权限 tar -xzf /tmp/efak-3.0.1-bin.tar.gz -C /opt/efak sudo chown -R kafka:kafka /opt/efak

注意:不要使用unzip解压,EFAK 的 tar.gz 包内部结构依赖tar的符号链接处理。曾有团队用unzip导致bin/startup.sh中的../conf路径解析失败。

4.2 编辑核心配置文件 efak.properties

进入/opt/efak/conf/efak.properties,按 CDH 环境定制以下关键参数(其余保持默认):

# 【必填】EFAK 服务监听地址(CDH 环境建议绑定内网IP) efak.webui.host=10.20.30.40 efak.webui.port=8042 # 【必填】Kafka 集群配置(使用 CDH 提供的 FQDN) efak.kafka.cluster.alias=CDH-PROD efak.kafka.cluster.bootstrap.servers=cdh-kafka-01.example.com:9092,cdh-kafka-02.example.com:9092,cdh-kafka-03.example.com:9092 # 【必填】ZooKeeper 配置(CDH Kerberos 环境) efak.zk.connect=zookeeper01.example.com:2181,zookeeper02.example.com:2181,zookeeper03.example.com:2181 efak.zk.sasl.enable=true efak.zk.sasl.kerberos.principal=zkclient@EXAMPLE.COM efak.zk.sasl.kerberos.keytab=/etc/security/keytabs/zkclient.keytab # 【必填】JMX 配置(指向 CDH Kafka 的 JMX 端口) efak.kafka.jmx.enable=true efak.kafka.jmx.servers=cdh-kafka-01.example.com:9999,cdh-kafka-02.example.com:9999,cdh-kafka-03.example.com:9999 # 【选填】KSQL 集成(如果启用 KSQL Server) efak.ksql.enable=true efak.ksql.servers=http://cdh-ksql-01.example.com:8088,http://cdh-ksql-02.example.com:8088 # 【选填】LDAP 认证(对接 CDH 的 AD) efak.security.auth.type=ldap efak.ldap.url=ldaps://ad.example.com:636 efak.ldap.base.dn=OU=KafkaUsers,DC=example,DC=com efak.ldap.user.dn.pattern=uid={0},OU=KafkaUsers,DC=example,DC=com # 【性能调优】JVM 参数(写入 startup.sh,非 properties 文件) # 此处仅配置 EFAK 自身参数,JVM 参数在 startup.sh 中修改

4.3 修改启动脚本以适配 CDH 环境

编辑/opt/efak/bin/startup.sh,找到JAVA_OPTS行,替换为:

# 原始行(注释掉) # JAVA_OPTS="-Xms512m -Xmx1g -server" # 替换为(根据集群规模调整,此处为 60 broker 示例) JAVA_OPTS="-Xms3200m -Xmx4800m -server -XX:+UseG1GC -XX:MaxGCPauseMillis=200 \ -Djava.security.auth.login.config=/opt/efak/conf/jaas.conf \ -Dsun.net.inetaddr.ttl=30"

其中jaas.conf是 Kerberos 认证必需的 JAAS 配置文件,需在/opt/efak/conf/下创建:

cat > /opt/efak/conf/jaas.conf << 'EOF' Client { com.sun.security.auth.module.Krb5LoginModule required useKeyTab=true storeKey=true keyTab="/etc/security/keytabs/kafka-client.keytab" principal="kafka-client@EXAMPLE.COM"; }; EOF sudo chown kafka:kafka /opt/efak/conf/jaas.conf sudo chmod 600 /opt/efak/conf/jaas.conf

4.4 创建 systemd 服务单元文件

为实现开机自启和日志管理,创建/etc/systemd/system/efak.service:

[Unit] Description=EFAK Kafka Manager After=network.target [Service] Type=simple User=kafka Group=kafka WorkingDirectory=/opt/efak ExecStart=/opt/efak/bin/startup.sh Restart=on-failure RestartSec=10 StandardOutput=journal StandardError=journal SyslogIdentifier=efak # 防止 OOM killer 杀死进程 OOMScoreAdjust=-500 [Install] WantedBy=multi-user.target

然后启用服务:

sudo systemctl daemon-reload sudo systemctl enable efak sudo systemctl start efak # 查看启动日志 sudo journalctl -u efak -f

启动成功标志:日志中出现INFO [main] o.a.e.EFAKApplication - Started EFAKApplication in X.XXX seconds,且无ZooKeeper connection failed或JMX connection timeout类错误。

4.5 首次访问与基础验证

在浏览器中访问http://10.20.30.40:8042(即efak.webui.host:efak.webui.port)。首次访问会跳转到登录页。若配置了 LDAP,则输入 AD 账号密码;若未配置,使用默认账号admin/admin登录。

登录后立即验证三项核心功能:

  1. Broker 列表:顶部导航栏 →Brokers→ 应显示所有 CDH Kafka broker(如cdh-kafka-01.example.com:9092),状态为Online,且右侧显示 CPU、Memory、Disk Usage 实时曲线。
  2. Topic 列表:导航栏 →Topics→ 应列出所有 CDH Kafka 中创建的 topic(如topic_user_clicks),点击任一 topic,右侧应显示 Partition 分布、ISR 状态、Message Rate 等。
  3. Consumer Lag:导航栏 →Consumers→ 选择任意 group(如flink-processor),应显示每个 partition 的Current Offset、Log End Offset、Lag值,且Lag列有颜色编码(绿色 < 1000,黄色 1000-10000,红色 > 10000)。

若以上任一环节失败,请立即检查对应环节的配置(如 Broker 列表为空 → 检查 3.1 和 3.3;Topic 列表为空 → 检查 3.2 和 3.4)。

5. 高可用部署与生产级调优——让 EFAK 成为 Kafka 运维的“心脏监护仪”

EFAK 在生产环境绝不能单点部署。我们线上集群采用“双活+自动故障转移”架构,其设计逻辑不是简单地多起几个实例,而是让每个 EFAK 实例承担明确角色,并通过外部组件协调状态。这套方案已在 3 个千万级 TPS 的 Kafka 集群中稳定运行 18 个月,年可用率达 99.997%。

5.1 双实例热备架构:Active-Standby 模式

我们不采用传统 Nginx 负载均衡,因为 EFAK 的 UI 状态(如用户 session、alert rule 编辑草稿、message preview 的 offset 位置)是强状态化的,简单轮询会导致用户操作丢失。取而代之的是基于 Consul 的服务注册与健康检查:

  • 部署两台 EFAK 服务器:efak-prod-01(主)、efak-prod-02(备)
  • 每台服务器运行 Consul Agent,注册自身为efak-web服务,健康检查脚本为:
    # /usr/local/bin/check-efak.sh # 检查 EFAK 进程存活且端口可连 if pgrep -f "EFAKApplication" > /dev/null && nc -z 127.0.0.1 8042; then exit 0 else exit 1 fi
  • 在 Consul 上配置service "efak-web"的check,并设置passing状态阈值为2(连续 2 次检查通过才标记为 healthy)
  • 外部 DNS(如 CoreDNS)配置efak.prod.example.com的 SRV 记录,只返回passing状态的实例 IP

这样,当efak-prod-01因硬件故障宕机时,Consul 在 30 秒内将其标记为critical,DNS 查询efak.prod.example.com将自动返回efak-prod-02的 IP,用户浏览器刷新即可无缝切换,session cookie 仍有效(因为两台实例共享同一 Redis session store)。

5.2 外部 Session Store:Redis 集群托管用户状态

EFAK 默认使用内存存储 session,这在双实例下必然导致 session 不一致。我们将其迁移到 Redis Cluster:

# 在 efak.properties 中添加 efak.session.store.type=redis efak.redis.host=redis-cluster.prod.example.com efak.redis.port=6379 efak.redis.password=your_redis_password efak.redis.database=1 efak.redis.timeout=2000

Redis Cluster 的 key 命名空间为efak:session:*,TTL 设为 30 分钟。实测表明,即使 Redis Cluster 某个 shard 故障,EFAK 仍能降级为本地 session(通过efak.session.fallback=true配置),保证基本功能可用。

5.3 告警引擎的分级推送策略

EFAK 内置的告警(Alert)模块不是简单的邮件轰炸器。我们配置了三级响应机制:

  • Level 1(P0):Lag > 1000000 AND Duration > 60s→ 触发电话告警(通过 PagerDuty webhook),同时自动执行kafka-consumer-groups.sh --bootstrap-server ... --group xxx --reset-offsets --to-earliest --execute重置 offset(需提前授权)。
  • Level 2(P1):UnderReplicatedPartitions > 0 AND Duration > 300s→ 发送企业微信消息到 Kafka 运维群,附带自动诊断链接(如http://efak.prod.example.com/brokers?filter=under_replicated)。
  • Level 3(P2):MessageRate < 1000 AND Duration > 300s(针对低频 topic)→ 生成工单(Jira webhook),分配给对应业务线负责人。

告警规则在efak.properties中定义为:

efak.alert.rule.1.name=HighLag efak.alert.rule.1.expression=lag > 1000000 && duration > 60 efak.alert.rule.1.action=phone,pagerduty efak.alert.rule.1.action.phone.number=+8613800138000 efak.alert.rule.2.name=UnderReplicated efak.alert.rule.2.expression=under_replicated_partitions > 0 && duration > 300 efak.alert.rule.2.action=wechat,webhook efak.alert.rule.2.action.wechat.group=Kafka-Ops

5.4 数据持久化迁移:从 H2 到 PostgreSQL

H2 数据库在单机环境下足够,但双实例需共享 schema history 和 offset 轨迹。我们迁移到 PostgreSQL 12:

-- 创建数据库 CREATE DATABASE efak_prod OWNER kafka; -- 创建表(EFAK 3.0.1 的 DDL) CREATE TABLE efak_topic_history ( id BIGSERIAL PRIMARY KEY, topic_name VARCHAR(255) NOT NULL, config_key VARCHAR(255), config_value TEXT, created_time TIMESTAMP WITH TIME ZONE DEFAULT NOW() ); CREATE TABLE efak_offset_history ( id BIGSERIAL PRIMARY KEY, group_id VARCHAR(255) NOT NULL, topic_name VARCHAR(255) NOT NULL, partition_id INT NOT NULL, offset_value BIGINT NOT NULL, timestamp TIMESTAMP WITH TIME ZONE DEFAULT NOW() );

然后在efak.properties中切换:

efak.db.type=postgresql efak.db.url=jdbc:postgresql://pg-prod-01.example.com:5432/efak_prod efak.db.username=kafka efak.db.password=your_pg_password efak.db.driver=org.postgresql.Driver

迁移后,Compare Offset History功能的查询响应时间从 8.2s 降至 0.3s(索引优化后),且支持跨月数据对比。

5.5 性能压测与容量规划

我们对 EFAK 进行了真实流量压测:模拟 200 个并发用户,每人每秒执行 1 次 topic describe、1 次 consumer lag refresh、1 次 message preview。结果如下:

配置平均响应时间95% 延迟CPU 使用率内存占用
4c8g + H21240ms2100ms78%3.2GB
8c16g + PostgreSQL320ms580ms42%5.1GB
16c32g + PostgreSQL + Redis180ms310ms28%8.7GB

结论:对于 100+ topic、50+ broker 的 CDH 集群,推荐 EFAK 服务器配置为 8c16g,数据库单独部署 4c8g PostgreSQL,Redis Cluster 至少 3 shard × 2 replica。低于此规格,UI 会出现明显卡顿,特别是Message Preview的分页加载。

注意:EFAK 的Message Preview功能默认只拉取 100 条消息,但若用户手动修改为 10000 条,会触发全量 scan,导致 broker 线程阻塞。我们在efak.properties中强制限制:efak.message.preview.max.count=1000,并在 UI 上隐藏“自定义数量”输入框,避免误操作。

6. 常见问题排查手册——从“页面空白”到“数据错乱”的 12 个真实案例

EFAK 的安装不是一劳永逸,日常运维中会遇到各种“看似奇怪、实则有因”的问题。以下是我在生产环境中记录的 12 个高频问题及其根因分析,每个都附带可验证的诊断命令和修复步骤。

6.1 问题:Topic 列表为空,但kafka-topics.sh --list能正常返回

现象:EFAK UI 的 Topics 页面显示 “No topics found”,而终端执行kafka-topics.sh --bootstrap-server cdh-kafka-01:9092 --list返回 237 个 topic。

根因分析:EFAK 的

返回列表