Kubeflow Pipeline(KFP)日志文件写入Google Cloud Storage的问题
Kubeflow Pipeline(KFP)日志文件写入Google Cloud Storage的问题
我明白你遇到的困扰了——在KFP组件里用logging打终端日志一切正常,但想把日志写到文件并存到GCS就卡壳了。其实核心原因很简单:KFP组件是在临时容器里运行的,你直接写本地文件的话,容器执行完就会被销毁,文件也跟着消失,而且没同步到GCS,所以根本看不到结果。下面给你两个实用的解决方法,直接改现有代码就能用:
方法一:本地写日志后手动上传到指定GCS路径
这个方法逻辑很清晰:先把日志写到容器的临时文件,任务完成后再把文件上传到你指定的GCS桶路径,完全可控:
- 先确认组件能访问GCS:如果你的base镜像没装
google-cloud-storage,就在组件装饰器里加上依赖(已经装了的话可以跳过); - 给logger添加文件处理器,把日志写到本地临时文件;
- 任务结束后调用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
相关产品推荐
相关产品推荐

