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

Google Cloud Storage Python上传脚本遇500错误,如何避免卡顿?

解决GCS上传500错误的重试与无卡顿处理方案

嘿,我看到你遇到GCS上传时的500后端错误导致脚本卡住的问题了——这类临时错误其实很常见,通过合理的重试机制和错误处理就能解决。我帮你优化了代码,加上了可靠的重试策略,同时修复了原代码里的一些小问题,让脚本能稳定无卡顿地运行。

核心改进点

  • 正确捕获GCS的InternalServerError及相关临时异常
  • 添加指数退避重试策略(比固定间隔重试更友好,避免给服务器额外施压)
  • 复用Storage客户端(避免每次上传都重新初始化连接)
  • 调整文件清理逻辑(只有上传成功才删除源文件和压缩文件,防止数据丢失)
  • 修复文件路径拼接、格式解析的潜在问题

修改后的完整代码

#!/usr/bin/env python
import threading
import time
import signal
import argparse
import os
import gzip
import shutil
from datetime import timedelta
from google.cloud import storage
from google.api_core.exceptions import InternalServerError, ServiceUnavailable
from google.api_core.retry import Retry

# 配置重试策略:针对500/503临时错误,最多重试5次,间隔指数增长
DEFAULT_RETRY = Retry(
    predicate=lambda exc: isinstance(exc, (InternalServerError, ServiceUnavailable)),
    initial=1.0,  # 第一次重试等待1秒
    multiplier=2,  # 每次重试间隔翻倍
    maximum=10.0,  # 最大等待10秒
    deadline=60.0  # 总重试超时60秒
)

WAIT_TIME_SECONDS = 120

class ProgramKilled(Exception):
    pass

# 全局复用Storage客户端,避免重复初始化
storage_client = storage.Client()

def upload_blob(bucket_name, data_path):
    print(f"This program was started at {time.ctime()} ")
    print(f"Target bucket: {bucket_name}")
    
    # 确保路径末尾有斜杠,避免拼接错误
    if not data_path.endswith('/'):
        data_path += '/'
    
    files = []
    # 遍历目录下所有.log文件,保存完整路径
    for r, d, f in os.walk(data_path):
        for file in f:
            if file.endswith(".log"):
                files.append(os.path.join(r, file))
    
    for log_file_path in files:
        file_name = os.path.basename(log_file_path)
        gz_file_path = log_file_path + ".gz"
        
        # 步骤1:压缩日志文件
        try:
            with open(log_file_path, 'rb') as f_in:
                with gzip.open(gz_file_path, 'wb') as f_out:
                    shutil.copyfileobj(f_in, f_out)
            print(f"Compressed {log_file_path} to {gz_file_path}")
        except Exception as e:
            print(f"Failed to compress {log_file_path}: {str(e)}")
            continue  # 压缩失败,跳过该文件继续处理下一个
        
        # 步骤2:解析文件格式,生成GCS存储路径
        blob_path = file_name + ".gz"
        try:
            d, fn = file_name.split("__")
            s = d.replace("-", "/")
            blob_path = f"{s}/{blob_path}"
        except ValueError:
            print(f"File {file_name} doesn't match '__' format, using default path")
        
        # 步骤3:带重试的GCS上传
        try:
            bucket = storage_client.get_bucket(bucket_name, retry=DEFAULT_RETRY)
            blob = bucket.blob(blob_path)
            # 上传操作也应用重试策略
            blob.upload_from_filename(gz_file_path, retry=DEFAULT_RETRY)
            print(f"Successfully uploaded {file_name} to gs://{bucket_name}/{blob_path}")
            
            # 只有上传成功才清理文件
            cleanup(log_file_path, gz_file_path)
        except (InternalServerError, ServiceUnavailable) as e:
            print(f"Final upload failed for {file_name} after retries: {str(e)}")
            # 保留文件,下次任务继续尝试
        except Exception as e:
            print(f"Unexpected error uploading {file_name}: {str(e)}")
            # 其他异常也保留文件,避免数据丢失

def cleanup(source_file_name, gz_file_name):
    try:
        os.remove(source_file_name)
        os.remove(gz_file_name)
        print(f"Removed files: {source_file_name} and {gz_file_name}")
    except Exception as e:
        print(f"Failed to clean up files: {str(e)}")

def signal_handler(signum, frame):
    raise ProgramKilled

class Job(threading.Thread):
    def __init__(self, interval, execute, args):
        threading.Thread.__init__(self)
        self.daemon = False
        self.stopped = threading.Event()
        self.interval = interval
        self.execute = execute
        self.args = args

    def stop(self):
        self.stopped.set()
        self.join()

    def run(self):
        # 启动后立即执行一次,再按间隔循环
        self.execute(*self.args)
        while not self.stopped.wait(self.interval.total_seconds()):
            self.execute(*self.args)

if __name__ == "__main__":
    signal.signal(signal.SIGTERM, signal_handler)
    signal.signal(signal.SIGINT, signal_handler)
    
    parser = argparse.ArgumentParser(
        description="Upload log files to Google Cloud Storage with automatic compression and retry",
        formatter_class=argparse.RawDescriptionHelpFormatter)
    parser.add_argument('bucket_name', help='Your cloud storage bucket name')
    subparsers = parser.add_subparsers(dest='command', required=True)
    upload_parser = subparsers.add_parser('upload', help='Upload log files from specified path')
    upload_parser.add_argument('data_path', help='Path to directory containing log files')
    
    args = parser.parse_args()
    
    job = Job(
        interval=timedelta(seconds=WAIT_TIME_SECONDS),
        execute=upload_blob,
        args=(args.bucket_name, args.data_path)
    )
    job.start()
    
    while True:
        try:
            time.sleep(1)
        except ProgramKilled:
            print("Program killed: stopping scheduled job")
            job.stop()
            break

关键改进说明

  1. 指数退避重试:使用Google官方的Retry类,专门针对临时的500/503错误进行重试,间隔从1秒开始翻倍,最多等待10秒,总超时60秒——这种策略既保证了重试成功率,又不会持续给GCS后端施压。
  2. 正确的异常捕获:原代码写错了异常类路径,现在改为从google.api_core.exceptions导入正确的异常类,确保能精准捕获目标错误。
  3. 客户端复用:全局只创建一次storage_client,避免每次上传都重新建立连接,提升运行效率。
  4. 安全的文件管理:只有当上传完全成功后才删除源文件和压缩文件,避免上传失败导致数据丢失;如果压缩或上传出错,文件会保留,下次任务会自动重试。
  5. 路径处理优化:用os.path.join处理文件路径,避免手动拼接的错误;自动补全路径末尾的斜杠,确保路径拼接正确。
  6. 即时首次执行:原代码会先等待120秒才执行第一次上传,修改后启动后立即执行一次,更符合实际使用需求。

这样修改后,脚本遇到临时的500错误会自动重试,不会卡住,同时能保证数据安全,运行更稳定。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:46:19