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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 11:57:29