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

Kubeflow Pipeline(KFP)日志文件写入Google Cloud Storage的问题

Kubeflow Pipeline(KFP)日志文件写入Google Cloud Storage的问题

我明白你遇到的困扰了——在KFP组件里用logging打终端日志一切正常,但想把日志写到文件并存到GCS就卡壳了。其实核心原因很简单:KFP组件是在临时容器里运行的,你直接写本地文件的话,容器执行完就会被销毁,文件也跟着消失,而且没同步到GCS,所以根本看不到结果。下面给你两个实用的解决方法,直接改现有代码就能用:

方法一:本地写日志后手动上传到指定GCS路径

这个方法逻辑很清晰:先把日志写到容器的临时文件,任务完成后再把文件上传到你指定的GCS桶路径,完全可控:

  1. 先确认组件能访问GCS:如果你的base镜像没装google-cloud-storage,就在组件装饰器里加上依赖(已经装了的话可以跳过);
  2. 给logger添加文件处理器,把日志写到本地临时文件;
  3. 任务结束后调用GCS客户端把日志文件上传到目标路径。

修改后的代码示例:

from kfp.dsl import Artifact, Dataset, Input, Metrics, Model, Output, component, Markdown
from google.cloud import storage
import os
import logging

@component(
    base_image="europe-west2-docker.pkg.dev/ml-repo/bigquery",
    # 如果base镜像没有google-cloud-storage,就加上这行
    # packages_to_install=["google-cloud-storage"]
)
def read_data_from_big_query_component(
    raw_dataset: Output[Dataset],
    gcs_bucket_name: str = "your-gcs-bucket-name",
    gcs_log_path: str = "pipeline-logs/bigquery_read_logs.log"
):
    logger = logging.getLogger(__name__)
    logger.setLevel(logging.INFO)

    # 配置本地日志文件处理器
    local_log_path = "/tmp/bigquery_read_logs.log"
    file_handler = logging.FileHandler(local_log_path)
    file_handler.setLevel(logging.INFO)
    # 给日志加时间戳,方便后续排查问题
    formatter = logging.Formatter("%(asctime)s - %(levelname)s - %(message)s")
    file_handler.setFormatter(formatter)
    logger.addHandler(file_handler)

    # --- 你原来的业务代码完全保留 ---
    logger.info("LOG - INFO 1: Starting the BigQuery data reading process.")

    # Connect to BigQuery
    client = bigquery.Client(project="my project ID goes here")
    
    logger.info("LOG - INFO 2: Connected to the BigQuery client.")

    query = f"""      
      SELECT * FROM MY_TABLE
      WHERE customer_name = "MY CUSTOMER NAME" and DATETIME >= "DATE"
      ORDER BY DATETIME;
      """
    
    logger.info("LOG - INFO 3: Query prepared.")

    job_config = bigquery.QueryJobConfig()
    query_job = client.query(query=query, job_config=job_config)
    df = query_job.result().to_dataframe()

    logger.info("LOG - INFO 4: Fetched data successfully.")
    logger.info(f"LOG - INFO 5: Number of rows fetched: {len(df)}")
    logger.info(f"LOG - INFO 6: Data preview:\n{df.head()}")
    logger.info("LOG - INFO 7: Data converted to dataframe.")
    
    if df.empty:
        logger.warning("LOG - WARNING 1: The fetched data from BigQuery is empty.")

    df.to_csv(raw_dataset.path, index=False)
    
    logger.info(f'LOG - INFO 8: Data saved to a CSV file as {raw_dataset}.')
    logger.info("LOG - INFO 9: BigQuery data reading process completed.")
    # --- 业务代码结束 ---

    # 上传本地日志到GCS
    def upload_logs():
        storage_client = storage.Client()
        bucket = storage_client.bucket(gcs_bucket_name)
        blob = bucket.blob(gcs_log_path)
        blob.upload_from_filename(local_log_path)
        logger.info(f"✅ Log file uploaded to GCS: gs://{gcs_bucket_name}/{gcs_log_path}")
    
    try:
        upload_logs()
    except Exception as e:
        logger.error(f"❌ Failed to upload logs to GCS: {str(e)}")
        raise

方法二:利用KFP Artifact自动托管日志文件

如果不想手动处理GCS上传,可以用KFP的Output[Artifact]类型,KFP会自动把你写入这个artifact路径的文件同步到GCS的pipeline artifact存储中,还能在KFP UI里直接查看日志文件,非常省心:

修改后的代码示例:

from kfp.dsl import Artifact, Dataset, Input, Metrics, Model, Output, component, Markdown
import logging

@component(
    base_image="europe-west2-docker.pkg.dev/ml-repo/bigquery",
)
def read_data_from_big_query_component(
    raw_dataset: Output[Dataset],
    process_log: Output[Artifact],  # 新增日志artifact参数
):
    logger = logging.getLogger(__name__)
    logger.setLevel(logging.INFO)

    # 配置日志处理器到KFP Artifact的路径
    file_handler = logging.FileHandler(process_log.path)
    file_handler.setLevel(logging.INFO)
    formatter = logging.Formatter("%(asctime)s - %(levelname)s - %(message)s")
    file_handler.setFormatter(formatter)
    logger.addHandler(file_handler)

    # --- 你原来的业务代码完全不用改 ---
    logger.info("LOG - INFO 1: Starting the BigQuery data reading process.")

    # Connect to BigQuery
    client = bigquery.Client(project="my project ID goes here")
    
    logger.info("LOG - INFO 2: Connected to the BigQuery client.")

    query = f"""      
      SELECT * FROM MY_TABLE
      WHERE customer_name = "MY CUSTOMER NAME" and DATETIME >= "DATE"
      ORDER BY DATETIME;
      """
    
    logger.info("LOG - INFO 3: Query prepared.")

    job_config = bigquery.QueryJobConfig()
    query_job = client.query(query=query, job_config=job_config)
    df = query_job.result().to_dataframe()

    logger.info("LOG - INFO 4: Fetched data successfully.")
    logger.info(f"LOG - INFO 5: Number of rows fetched: {len(df)}")
    logger.info(f"LOG - INFO 6: Data preview:\n{df.head()}")
    logger.info("LOG - INFO 7: Data converted to dataframe.")
    
    if df.empty:
        logger.warning("LOG - WARNING 1: The fetched data from BigQuery is empty.")

    df.to_csv(raw_dataset.path, index=False)
    
    logger.info(f'LOG - INFO 8: Data saved to a CSV file as {raw_dataset}.')
    logger.info("LOG - INFO 9: BigQuery data reading process completed.")

用这个方法的话,你在运行pipeline后,进入KFP UI的对应run详情页,找到process_log这个artifact,就能直接下载或查看日志文件了——KFP会自动把它存在GCS的pipeline默认artifact桶里,不用你手动管理路径。

最后提两个注意点:

  • 确保你的KFP组件使用的服务账号有GCS的写入权限(自动托管的GKE集群默认服务账号通常有这个权限,自定义账号的话需要添加roles/storage.objectCreator权限);
  • 日志级别要对应,比如你用logger.info(),就要确保handler的级别是logging.INFO或更低,不然日志会被过滤掉。

备注:内容来源于stack exchange,提问作者Alireza Bolhari

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 17:48:09