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

