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

Kafka消费者向GCP Bucket存数据时的超时与重平衡问题排查

Kafka消费者数据拉取与GCP存储问题排查

我使用Kafka消费者从指定Topic拉取数据,按1分钟间隔将数据存储至GCP Bucket。以下是Kafka消费者的代码:

import os
import kafka
import json
from datetime import date
from io import BytesIO
from google.cloud import storage
import time

storage_client = storage.Client.from_service_account_json(
    os.path.dirname(os.path.abspath(__file__)) + "/XXXX.json"
)
bucket = storage_client.get_bucket("XXXX")


def get_filname(topic):
    return bucket.blob(topic + "_" + str(date.today()))


def retrieve_data_from_topic(topic, topic_name):
    data = []
    existing_data = []
    begin_time = time.time()
    blob = get_filname(topic_name)
    if not blob.exists():
        store_to_bucket(blob, [])
    for message in topic:
        blob = get_filname(topic_name)
        data.append(message.value)
        file_data = blob.download_as_text()
        existing_data = list(json.loads(file_data))
        new_data = []
        interval = 60  ## 1 min
        for fresh_d in data:
            if fresh_d not in existing_data:
                existing_data.append(fresh_d)
        current_time = time.time()
        period = current_time - begin_time
        if period >= interval:
            store_to_bucket(blob, existing_data)
            begin_time = current_time
            topic.commit()


def store_to_bucket(blob, msg):
    blob.upload_from_string(
        data=json.dumps(msg),
    )


def retrieve_topic(topic: str):
    consumer = kafka.KafkaConsumer(
        topic,
        bootstrap_servers=["kafka:9093", "kafka:9092"],
        value_deserializer=lambda m: json.loads(m.decode("utf-8")),
        group_id="metaowl-bucket" + topic,
        auto_offset_reset="earliest",
        max_poll_records=300,
        max_poll_interval_ms= 14400000,
        enable_auto_commit=False,
    )
    return consumer


def topic_list():
    consumer = kafka.KafkaConsumer(
        value_deserializer=lambda m: json.loads(m.decode("utf-8")),
        group_id="metaowl-bucket-list",
        auto_offset_reset="earliest",
        enable_auto_commit=True,
        bootstrap_servers=["kafka:9093", "kafka:9092"],
    )
    tokens = set()
    for topic in consumer.topics():
        tokens.add(topic.split("_")[0])
    return tokens

当前问题

消费者处理消息耗时过长,有时会超出设置的4小时max_poll_interval_ms限制,导致消费者组ID触发重平衡并失败(仅使用一个消费者组ID)。尝试调整过max_poll_records和max_poll_interval_ms参数,但由于消息本身并不复杂,调整参数并无意义,恳请排查问题原因。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 19:40:34