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

Confluent Kafka消费者组当前偏移量返回-1001的解决问询

问题描述

我尝试用Python代码计算Confluent Kafka上消费者组的Lag,但运行后所有分区的current-offset始终返回-1001,而end-offset能正常获取。但用Shell命令kafka-consumer-groups却能正确拿到当前偏移量和末尾偏移量。请问Python代码需要做哪些调整才能得到准确的current-offset?

我的Python代码:

from confluent_kafka.admin import AdminClient, NewTopic
from confluent_kafka import KafkaException, KafkaError, Consumer
from confluent_kafka import TopicPartition

import json

# Set up the configuration for the Confluent Cluster
conf = {'bootstrap.servers': 'pkc-43332.us-west1.gcp.confluent.cloud:9092',
        'security.protocol': 'SASL_SSL',
        'sasl.mechanism': 'PLAIN',
        'sasl.username': '<user-name>',
        'sasl.password': '<pswd>'}

# Create the AdminClient using the configuration
admin_client = AdminClient(conf)

# Get the consumer group description
group_metadata = admin_client.list_groups()

# Check if the consumer group is active
group_name = 'connect-consumer-group'
if group_name not in [group.id for group in group_metadata]:
    print(f"No consumer group with name {group_name} found.")
    exit()

# Get the consumer group details
group_info = admin_client.describe_consumer_groups([group_name])
group_info = group_info[group_name].result()

# Get the topic partitions for the consumer group
topic_partitions = {}
for member in group_info.members:
    for tp in member.assignment.topic_partitions:
        topic_partitions[tp.topic] = topic_partitions.get(tp.topic, []) + [tp.partition]

# Create a Consumer object
consumer_conf = {'bootstrap.servers': 'pkc-43332.us-west1.gcp.confluent.cloud:9092',
                 'security.protocol': 'SASL_SSL',
                 'sasl.mechanism': 'PLAIN',
                 'sasl.username': '<user-name>',
                 'sasl.password': '<pswd>',
                 'group.id': group_name,
                 'auto.offset.reset': 'earliest'}
consumer = Consumer(consumer_conf)

# Calculate lag for each topic partition
for topic, partitions in topic_partitions.items():
    for partition in partitions:
        tp = TopicPartition(topic, partition)
        current_offset = consumer.position([tp])[0].offset
        end_offset = consumer.get_watermark_offsets(tp)[1]

        # Calculate lag
        lag = end_offset - current_offset
        print(f"Lag for {topic}-partition-{partition}: {lag}, end offset is {end_offset}, current offset is {current_offset}")

能正常返回结果的Shell命令:

kafka-consumer-groups 
--bootstrap-server pkc-43332.us-west1.gcp.confluent.cloud:9092 
--command-config /home/dbuser/client-config.properties 
--describe --group connect-consumer-group 
--timeout 10000
解决方案

返回-1001(对应KafkaError.OFFSET_INVALID)的核心原因是:你创建的Consumer实例未完成消费者组加入与分区偏移量同步流程,直接调用position()会返回无效偏移量。

以下是两种可行的调整方案:

方案一:通过AdminClient直接获取已提交偏移量(推荐)

这种方式和kafka-consumer-groups --describe的逻辑完全一致,无需创建Consumer实例,更轻量且不会干扰原有消费者组状态:

from confluent_kafka.admin import AdminClient
from confluent_kafka import KafkaException

# 集群配置
conf = {
    'bootstrap.servers': 'pkc-43332.us-west1.gcp.confluent.cloud:9092',
    'security.protocol': 'SASL_SSL',
    'sasl.mechanism': 'PLAIN',
    'sasl.username': '<user-name>',
    'sasl.password': '<pswd>'
}

admin_client = AdminClient(conf)
group_name = 'connect-consumer-group'

# 检查消费者组是否存在
try:
    group_metadata = admin_client.list_groups().result()
    if group_name not in [g.id for g in group_metadata]:
        print(f"未找到消费者组 {group_name}")
        exit(1)
except KafkaException as e:
    print(f"查询消费者组失败: {e}")
    exit(1)

# 获取消费者组详情(包含已提交偏移量)
group_desc = admin_client.describe_consumer_groups([group_name]).result()[group_name]

# 获取分区末尾偏移量
def get_end_offset(admin, topic, partition):
    return admin.list_topics(topic).result().topics[topic].partitions[partition].high_watermark

# 计算Lag
for member in group_desc.members:
    for tp in member.assignment.topic_partitions:
        # 匹配当前分区的已提交偏移量
        committed_offset = next((o.offset for o in group_desc.partitions_assigned if o.topic == tp.topic and o.partition == tp.partition), None)
        if committed_offset is None:
            print(f"{tp.topic}-partition-{tp.partition}: 无已提交偏移量")
            continue
        # 获取末尾偏移量并计算Lag
        end_offset = get_end_offset(admin_client, tp.topic, tp.partition)
        lag = end_offset - committed_offset
        print(f"Lag for {tp.topic}-partition-{tp.partition}: {lag}, end offset: {end_offset}, current offset: {committed_offset}")

方案二:修复Consumer的使用流程

如果必须使用Consumer,需先触发组协调与偏移量同步:

from confluent_kafka import Consumer, TopicPartition
from confluent_kafka.admin import AdminClient

# 集群配置
conf = {
    'bootstrap.servers': 'pkc-43332.us-west1.gcp.confluent.cloud:9092',
    'security.protocol': 'SASL_SSL',
    'sasl.mechanism': 'PLAIN',
    'sasl.username': '<user-name>',
    'sasl.password': '<pswd>'
}

admin_client = AdminClient(conf)
group_name = 'connect-consumer-group'

# 收集消费者组的分区信息
group_info = admin_client.describe_consumer_groups([group_name]).result()[group_name]
topic_partitions = []
for member in group_info.members:
    for tp in member.assignment.topic_partitions:
        topic_partitions.append(TopicPartition(tp.topic, tp.partition))

# 配置Consumer
consumer_conf = conf.copy()
consumer_conf.update({
    'group.id': group_name,
    'auto.offset.reset': 'earliest',
    'enable.auto.commit': False  # 禁止自动提交,避免干扰原有组偏移量
})
consumer = Consumer(consumer_conf)

# 分配分区并触发组协调流程
consumer.assign(topic_partitions)
consumer.poll(timeout=5.0)  # 必须调用poll完成偏移量同步

# 计算Lag
for tp in topic_partitions:
    committed_offset = consumer.committed([tp])[0].offset
    end_offset = consumer.get_watermark_offsets(tp)[1]
    lag = end_offset - committed_offset
    print(f"Lag for {tp.topic}-partition-{tp.partition}: {lag}, end offset: {end_offset}, current offset: {committed_offset}")

consumer.close()

关键说明

  • 方案一直接读取消费者组的已提交偏移量,是最贴合Shell命令的实现方式,无需额外创建消费者。
  • 方案二中的consumer.poll()是核心:Confluent Kafka的Consumer必须通过poll操作完成组加入、分区分配和偏移量同步,否则无法获取有效偏移量。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 02:38:14