You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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"'")
    }
}'

脚本说明:

  1. 用awk跳过表头,逐行解析字段
  2. 自动累加每个消费者ID的所有分区滞后量
  3. 生成UTC时间戳,保证时间一致性
  4. 同时写入本地日志和推送ES(ES索引按日期拆分,方便管理)
  5. 处理了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()

脚本说明:

  1. 包含异常捕获,命令执行失败时会记录错误日志
  2. 不仅统计总滞后,还保存每个分区的滞后细节,方便后续排查问题
  3. 支持ES集群多节点连接,推送更稳定
  4. 时间戳统一使用UTC格式,避免时区问题

额外注意事项
  • 可以用cron定时执行脚本(比如每分钟一次),持续收集滞后数据
  • 确保执行脚本的用户有Kafka消费者组的描述权限,以及ES的写入权限
  • 如果消费者组处于rebalance状态,CONSUMER-ID会显示null,可以根据需求调整脚本,改为按GROUP+TOPIC统计总滞后
  • 对于Bash脚本,建议把Kafka命令的绝对路径写全,避免PATH环境变量问题

内容的提问来源于stack exchange,提问作者oraclept

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 10:32:46