GCP云函数中Pandas导出Parquet至GCS时遇文件不存在错误
GCP云函数导出Parquet至GCS时FileNotFoundError问题排查与解决
我在GCP云函数中尝试将MySQL数据导出为Parquet文件并存储到GCS云存储桶,执行chunk.to_parquet(parquet_file_path, engine='fastparquet', compression='snappy')时触发如下错误:
FileNotFoundError: [Errno 2] No such file or directory: 'new_folder_20230206_065500/table1-20230206_065638.parquet'
桶内的文件夹已成功创建,但Parquet文件无法生成,相关代码如下:
import mysql.connector import pandas as pd from google.cloud import storage from datetime import datetime, timedelta import os def extract_data_to_gcs(request): connection = mysql.connector.connect( host=os.getenv('..'), user=os.getenv('...'), password=os.getenv('...'), database='....' ) cursor = connection.cursor(buffered=True) tables = ["table1", "table2", "table3"] client = storage.Client() bucket = client.bucket('data-lake-archive') # Create a timestamp-based folder name now = datetime.now() folder_name = now.strftime("new_folder_%Y%m%d_%H%M%S") folder_path = f"{folder_name}/" # Create the folder in the GCS bucket blob = bucket.blob(folder_path) blob.upload_from_string("", content_type="application/octet-stream") for table in tables: cursor.execute("SELECT * FROM {}".format(table)) chunks = pd.read_sql_query("SELECT * FROM {}".format(table), connection, chunksize=5000000) for i, chunk in enumerate(chunks): chunk.columns = [str(col) for col in chunk.columns] ingestion_timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") parquet_file_path = folder_path + f"{table}-{i}.parquet" timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") # parquet_file_path = folder_path + f'abc.parquet' print(f'folder path is {folder_path}') print(f'parquet file path is {parquet_file_path}') chunk.to_parquet(parquet_file_path, engine='fastparquet', compression='snappy') # blob = bucket.blob(folder_path + f'{table}-{i}.parquet') # blob.upload_from_filename(folder_path + f'{table}-{i}.parquet') cursor.execute("SELECT table_name, column_name FROM information_schema.key_column_usage WHERE referenced_table_name = '{}'".format(table)) referenced_tables = cursor.fetchall() for referenced_table in referenced_tables: chunks = pd.read_sql_query("SELECT * FROM {}".format(referenced_table[0]), connection, chunksize=5000000) for i, chunk in enumerate(chunks): chunk.columns = [str(col) for col in chunk.columns] ingestion_timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") chunk.to_parquet(f"{folder_path}{referenced_table[0]}-{ingestion_timestamp}-{i}.parquet", engine='fastparquet', compression='snappy') blob = bucket.blob(folder_path + f'{referenced_table[0]}-{ingestion_timestamp}-{i}.parquet') blob.upload_from_filename(folder_path + f'{referenced_table[0]}-{ingestion_timestamp}-{i}.parquet') return 'Data extracted and uploaded to GCS'
问题原因
GCP云函数的运行环境是无服务器临时文件系统,你代码中直接使用folder_path(如new_folder_20230206_065500/)作为本地路径,但这个路径并未在云函数的本地文件系统中创建——你之前创建的只是GCS桶里的"虚拟文件夹"(GCS本质是对象存储,没有真正的文件夹,仅通过对象前缀模拟),本地系统不存在该目录,所以to_parquet无法写入文件。
修复方案
方案1:使用临时目录中转后上传
云函数提供/tmp目录作为可写临时存储空间,步骤如下:
- 将Parquet文件写入
/tmp下的临时路径 - 把临时文件上传到GCS目标路径
- 上传完成后删除临时文件释放空间
修改后的核心代码片段:
import tempfile # 替换原有文件写入及上传逻辑 for table in tables: chunks = pd.read_sql_query("SELECT * FROM {}".format(table), connection, chunksize=5000000) for i, chunk in enumerate(chunks): chunk.columns = [str(col) for col in chunk.columns] # 生成临时文件 temp_file = tempfile.NamedTemporaryFile(suffix='.parquet', delete=False) temp_file_path = temp_file.name temp_file.close() # 写入临时文件 chunk.to_parquet(temp_file_path, engine='fastparquet', compression='snappy') # 上传至GCS gcs_blob_path = f"{folder_path}{table}-{i}.parquet" blob = bucket.blob(gcs_blob_path) blob.upload_from_filename(temp_file_path) # 删除临时文件 os.remove(temp_file_path)
方案2:直接写入GCS(无需本地文件)
利用fastparquet支持写入文件对象的特性,直接将Parquet数据写入GCS的Blob对象,省去本地文件步骤:
from io import BytesIO # 替换原有文件写入及上传逻辑 for table in tables: chunks = pd.read_sql_query("SELECT * FROM {}".format(table), connection, chunksize=5000000) for i, chunk in enumerate(chunks): chunk.columns = [str(col) for col in chunk.columns] # 创建内存缓冲区 parquet_buffer = BytesIO() # 将数据写入缓冲区 chunk.to_parquet(parquet_buffer, engine='fastparquet', compression='snappy') parquet_buffer.seek(0) # 重置缓冲区指针到开头 # 上传至GCS gcs_blob_path = f"{folder_path}{table}-{i}.parquet" blob = bucket.blob(gcs_blob_path) blob.upload_from_file(parquet_buffer, content_type='application/octet-stream')
额外注意事项
- 云函数
/tmp目录默认有512MB存储空间限制,若单块数据超过该限制,建议使用方案2或调小chunksize - 确保云函数服务账号拥有GCS存储桶的
storage.objects.create权限 - 代码中重复执行了
SELECT * FROM {table}(cursor.execute和pd.read_sql_query各一次),可删除冗余的cursor.execute调用优化性能
内容的提问来源于stack exchange,提问作者Tushaar
相关产品推荐
相关产品推荐

