Python脚本向GCP Bucket传文件时Cloud Logging日志未显示求助
解决多进程环境下GCP Cloud Logging不显示日志的问题
看起来你的脚本已经能成功同步文件到GCP Bucket,但日志没出现在Log Viewer里——核心问题出在多进程(multiprocessing.Pool)环境下CloudLoggingHandler的继承失效:父进程初始化的日志客户端在子进程中无法正常发送日志,因为fork后的进程会继承父进程的资源,但GCP日志依赖的网络连接/状态在子进程中是无效的。
下面是具体的修复步骤和修改后的代码:
1. 核心问题分析
当你用Pool创建子进程时,子进程会复制父进程的内存空间,包括已经初始化的CloudLoggingHandler。但GCP的日志客户端基于的HTTP连接或gRPC通道在fork后无法复用,导致子进程的日志无法被发送到GCP日志服务。
2. 解决方案:在每个子进程中重新初始化日志
你需要让每个子进程单独初始化CloudLoggingHandler,而不是继承父进程的logger。这里提供两种可行的实现方式:
方式一:在子进程任务中重新初始化Logger
修改execute_jobs方法,在任务执行前重新创建logger(同时清空可能继承的旧handler,避免重复日志或失效连接):
from multiprocessing import Pool from subprocess import Popen, PIPE, TimeoutExpired, CalledProcessError import os import sys import logging as lg import google.cloud.logging as gcl from google.cloud.logging.handlers import CloudLoggingHandler os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "/home/sant/key.json" ftp_path1 = "/home/sant" GCS_DATA_INGEST_BUCKET_URL = "dev2-ingest-manual" PROJECT_ID = "your-gcp-project-id" # 替换为你的GCP项目ID class GcsMover: def __init__(self): self.folder_list = ["raw_amr", "osr_data"] # 父进程的logger仅用于父进程日志,子进程会重新初始化 self.logger = self.create_logger() @staticmethod def create_logger(log_name="Root_Logger", log_level=lg.INFO): try: log_format = lg.Formatter("%(levelname)s %(asctime)s - %(message)s") # 显式指定项目ID,避免环境变量歧义 client = gcl.Client(project=PROJECT_ID) log_handler = CloudLoggingHandler(client) log_handler.setFormatter(log_format) logger = lg.getLogger(log_name) # 清空子进程继承的旧handler,避免重复发送或失效连接 if logger.handlers: logger.handlers.clear() logger.setLevel(log_level) logger.addHandler(log_handler) # 禁用日志传播,避免重复输出到根logger logger.propagate = False return logger except Exception as e: sys.exit(f"WARNING - Invalid cloud logging: {str(e)}") def execute_jobs(self, cmd): # 子进程执行任务前,重新初始化logger self.logger = self.create_logger() try: gs_sp = Popen(cmd, stdin=PIPE, stdout=PIPE, stderr=PIPE, shell=True) print(f"starting process with Pid {str(gs_sp.pid)} for command {cmd}") self.logger.info(f"starting process with Pid {str(gs_sp.pid)} for command {cmd}") sp_out, sp_err = gs_sp.communicate(timeout=int(3600)) except OSError as e: self.logger.error(f"Processing aborted for Pid {str(gs_sp.pid)}: {str(e)}") except TimeoutExpired: gs_sp.kill() self.logger.error(f"Processing aborted for Pid {str(gs_sp.pid)} due to timeout") else: if gs_sp.returncode != 0: # 解码stderr为字符串,避免二进制日志内容 err_msg = sp_err.decode('utf-8').strip() self.logger.error(f"Failure for Pid {str(gs_sp.pid)} (cmd: {cmd}): {err_msg}") else: success_msg = f"Loading successful for Pid {str(gs_sp.pid)}" print(success_msg) self.logger.info(success_msg) finally: # 强制刷新日志,确保所有日志都被发送到GCP for handler in self.logger.handlers: handler.flush() def move_files(self): command_list = [] for folder in self.folder_list: gs_command = f"gsutil -m rsync -r {ftp_path1}/{folder} gs://{GCS_DATA_INGEST_BUCKET_URL}/{folder}" command_list.append(gs_command) pool = Pool(processes=2, maxtasksperchild=1) pool.map(self.execute_jobs, iterable=command_list) pool.close() pool.join() def main(): gsu = GcsMover() gsu.move_files() if __name__ == "__main__": main()
方式二:用Pool的initializer初始化子进程
另一种更高效的方式是利用Pool的initializer参数,在每个子进程启动时一次性初始化logger:
# 新增子进程初始化函数 def init_worker(): os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "/home/sant/key.json" # 为子进程初始化logger GcsMover.create_logger() # 修改move_files方法 def move_files(self): command_list = [] for folder in self.folder_list: gs_command = f"gsutil -m rsync -r {ftp_path1}/{folder} gs://{GCS_DATA_INGEST_BUCKET_URL}/{folder}" command_list.append(gs_command) # 传入initializer,每个子进程启动时执行init_worker pool = Pool(processes=2, maxtasksperchild=1, initializer=init_worker) pool.map(self.execute_jobs, iterable=command_list) pool.close() pool.join()
3. 额外检查项
除了代码修改,还要确认以下配置:
- 服务账号权限:确保
key.json对应的服务账号拥有Logs Writer(roles/logging.logWriter)角色,否则无法写入日志。 - Log Viewer过滤条件:在GCP Log Viewer中,选择正确的项目,然后用查询语句过滤:
logName="projects/[你的项目ID]/logs/Root_Logger",避免因过滤条件错误看不到日志。 - 显式指定项目ID:在创建
gcl.Client时显式传入project=PROJECT_ID,避免依赖环境变量导致的项目不匹配。
内容的提问来源于stack exchange,提问作者Santanu Ghosh
相关产品推荐
相关产品推荐

