Kafka 1.1.0消费者分区总滞后计算脚本开发咨询
针对你的需求,我整理了两个实用的脚本方案,都是基于Kafka 1.1.0原生命令实现的,完全不需要依赖Kafka Manager或Burrow这类第三方工具,直接就能完成滞后量统计、日志记录和Elasticsearch推送。
先明确Kafka命令的输出格式
Kafka 1.1.0的bin/kafka-consumer-groups --describe输出大致如下(表头+数据行),我们需要解析其中的LAG和CONSUMER-ID字段来计算每个消费者的总滞后:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID my-group my-topic 0 100 150 50 consumer-1-xxx /192.168.1.100 consumer-1 my-group my-topic 1 200 220 20 consumer-1-xxx /192.168.1.100 consumer-1 my-group my-topic 2 300 310 10 consumer-2-xxx /192.168.1.101 consumer-2
方案1:轻量Bash脚本(适合快速部署)
这个脚本用原生命令组合实现,不需要额外依赖,适合在服务器上直接定时执行:
#!/bin/bash # 配置参数 BOOTSTRAP_SERVERS="kafka1:9092,kafka2:9092,kafka3:9092" GROUP_NAME="your-target-group" LOG_FILE="/var/log/kafka_consumer_lag.log" ES_ENDPOINT="http://es-node:9200/kafka-consumer-lag-$(date -u +%Y.%m.%d)/_doc" TIMESTAMP=$(date -u +"%Y-%m-%dT%H:%M:%SZ") # 执行命令并解析输出 bin/kafka-consumer-groups --bootstrap-server $BOOTSTRAP_SERVERS --describe --group $GROUP_NAME | awk 'NR>1 { # 提取核心字段 group=$1 lag=$6 consumer_id=$7 host=$8 client_id=$9 # 累加每个消费者的总滞后 total_lag[consumer_id] += lag # 保存消费者的元信息 meta[consumer_id] = "{\"host\":\"" substr(host,2) "\",\"client_id\":\"" client_id "\"}" } END { # 遍历结果,输出日志并推送ES for (cid in total_lag) { # 写入本地日志 log_line = sprintf("%s | group:%s | consumer_id:%s | host:%s | client_id:%s | total_lag:%d", "'$TIMESTAMP'", group, cid, substr(host,2), client_id, total_lag[cid]) print log_line >> "'$LOG_FILE'" # 生成JSON并推送Elasticsearch json = sprintf("{\"timestamp\":\"%s\",\"group\":\"%s\",\"consumer_id\":\"%s\",\"%s\",\"total_lag\":%d}", "'$TIMESTAMP'", group, cid, substr(meta[cid],2,length(meta[cid])-1), total_lag[cid]) system("curl -X POST '"$ES_ENDPOINT"' -H \"Content-Type: application/json\" -d '"json"'") } }'
脚本说明:
- 用
awk跳过表头,逐行解析字段 - 自动累加每个消费者ID的所有分区滞后量
- 生成UTC时间戳,保证时间一致性
- 同时写入本地日志和推送ES(ES索引按日期拆分,方便管理)
- 处理了
HOST字段开头的斜杠(比如/192.168.1.100转为192.168.1.100)
方案2:Python脚本(适合复杂场景)
如果需要更灵活的异常处理、细节记录或者后续扩展,Python脚本会更合适,还能直接用官方ES库推送数据:
import subprocess from datetime import datetime from elasticsearch import Elasticsearch # 配置参数 BOOTSTRAP_SERVERS = "kafka1:9092,kafka2:9092,kafka3:9092" GROUP_NAME = "your-target-group" ES_HOSTS = ["http://es-node1:9200", "http://es-node2:9200"] LOG_FILE = "/var/log/kafka_consumer_lag.log" def calculate_consumer_lag(): # 执行Kafka消费者组命令 cmd = [ "/path/to/kafka/bin/kafka-consumer-groups", "--bootstrap-server", BOOTSTRAP_SERVERS, "--describe", "--group", GROUP_NAME ] try: result = subprocess.run(cmd, capture_output=True, text=True, check=True) except subprocess.CalledProcessError as e: with open(LOG_FILE, "a") as f: f.write(f"{datetime.utcnow().isoformat()}Z | ERROR: {e.stderr}\n") return # 解析输出(跳过表头) lines = result.stdout.strip().split("\n")[1:] consumer_lag_map = {} timestamp = datetime.utcnow().isoformat() + "Z" for line in lines: parts = line.split() if len(parts) < 9: continue # 跳过格式异常的行 group, topic, partition, _, _, lag, consumer_id, host, client_id = parts[:9] lag = int(lag) # 初始化消费者记录 if consumer_id not in consumer_lag_map: consumer_lag_map[consumer_id] = { "timestamp": timestamp, "group": group, "consumer_id": consumer_id, "host": host.lstrip("/"), "client_id": client_id, "total_lag": 0, "partition_details": [] } # 累加总滞后,同时记录每个分区的细节 consumer_lag_map[consumer_id]["total_lag"] += lag consumer_lag_map[consumer_id]["partition_details"].append({ "topic": topic, "partition": int(partition), "lag": lag }) # 写入本地日志 with open(LOG_FILE, "a") as f: for data in consumer_lag_map.values(): log_line = (f"{data['timestamp']} | group:{data['group']} | consumer_id:{data['consumer_id']} " f"| host:{data['host']} | client_id:{data['client_id']} | total_lag:{data['total_lag']}\n") f.write(log_line) # 推送至Elasticsearch es = Elasticsearch(ES_HOSTS) index_name = f"kafka-consumer-lag-{datetime.utcnow().strftime('%Y.%m.%d')}" for data in consumer_lag_map.values(): es.index(index=index_name, document=data) if __name__ == "__main__": calculate_consumer_lag()
脚本说明:
- 包含异常捕获,命令执行失败时会记录错误日志
- 不仅统计总滞后,还保存每个分区的滞后细节,方便后续排查问题
- 支持ES集群多节点连接,推送更稳定
- 时间戳统一使用UTC格式,避免时区问题
额外注意事项
- 可以用
cron定时执行脚本(比如每分钟一次),持续收集滞后数据 - 确保执行脚本的用户有Kafka消费者组的描述权限,以及ES的写入权限
- 如果消费者组处于rebalance状态,
CONSUMER-ID会显示null,可以根据需求调整脚本,改为按GROUP+TOPIC统计总滞后 - 对于Bash脚本,建议把Kafka命令的绝对路径写全,避免PATH环境变量问题
内容的提问来源于stack exchange,提问作者oraclept
相关产品推荐
相关产品推荐

