如何通过Apache Beam读取Google Drive文件并写入GCS存储桶?
实现Apache Beam读取Google Drive文件并写入GCS
1. 准备依赖与认证
- 安装必要依赖:
pip install apache-beam[gcp] google-api-python-client google-auth-httplib2 google-auth-oauthlib - 配置认证:
- 创建Google Cloud服务账号,授予Google Drive文件读取权限和GCS存储桶写入权限
- 下载服务账号密钥JSON文件,设置环境变量:
export GOOGLE_APPLICATION_CREDENTIALS="/path/to/your/service-account-key.json"
2. 自定义Drive文件读取DoFn
Beam没有内置的Google Drive读取组件,需用Drive API自定义DoFn实现文件下载:
import apache_beam as beam from googleapiclient.discovery import build from googleapiclient.http import MediaIoBaseDownload import io class ReadDriveFile(beam.DoFn): def setup(self): # 初始化Drive API客户端 self.drive_service = build('drive', 'v3', cache_discovery=False) def process(self, file_id): # 根据file_id下载文件内容 request = self.drive_service.files().get_media(fileId=file_id) file_handle = io.BytesIO() downloader = MediaIoBaseDownload(file_handle, request) done = False while done is False: status, done = downloader.next_chunk() # 重置指针并返回内容(文本文件解码,二进制文件直接返回bytes) file_handle.seek(0) yield file_handle.read().decode('utf-8')
3. 构建完整Pipeline
整合读取、写入逻辑,替换为你的实际参数:
def run_pipeline(): # 替换为你的Drive文件ID和GCS输出路径 DRIVE_FILE_ID = "<your-file-id>" GCS_OUTPUT_PATH = "gs://your-bucket-name/output-file.txt" with beam.Pipeline() as pipeline: ( pipeline | "传入Drive文件ID" >> beam.Create([DRIVE_FILE_ID]) | "读取Drive文件" >> beam.ParDo(ReadDriveFile()) | "写入GCS" >> beam.io.WriteToText(GCS_OUTPUT_PATH) ) if __name__ == "__main__": run_pipeline()
关键说明
- Drive文件ID提取:从共享链接中获取,比如
https://drive.google.com/file/d/123abc/view里的123abc就是file_id - 大文件处理:若处理大文件,建议分块读取或用临时文件中转,避免内存溢出
- 权限校验:确保服务账号已被添加为Drive文件协作者(拥有读取权限),同时具备GCS存储桶的写入权限
- 二进制文件适配:处理图片、视频等二进制文件时,去掉
decode('utf-8'),改用beam.io.WriteToFiles并指定二进制模式
内容的提问来源于stack exchange,提问作者codebot
相关产品推荐
相关产品推荐

