如何使用Python获取GCS存储桶中最新存入的文件用于DAG开发
Python获取GCS指定路径下最新存入文件的实现方案
前置依赖
首先需要安装GCP官方的Storage SDK:
pip install google-cloud-storage
同时需要提前配置好GCP身份认证,本地调试可以通过gcloud auth application-default login生成密钥,部署到Airflow环境的话可以通过服务账号密钥文件或者Workload Identity完成认证。
核心实现代码
from google.cloud import storage def get_latest_gcs_file(bucket_name: str = "name", prefix: str = "neth/") -> str: storage_client = storage.Client() blobs = storage_client.list_blobs(bucket_name, prefix=prefix) # 过滤符合file_yymmdd_命名格式的文件,排除路径文件夹 valid_blobs = [] for blob in blobs: file_name = blob.name.split("/")[-1] if not blob.name.endswith("/") and file_name.startswith("file_"): valid_blobs.append(blob) if not valid_blobs: raise FileNotFoundError(f"路径gs://{bucket_name}/{prefix}下未找到符合命名规则的文件") # 按文件上传时间倒序排序,取第一个即为最新存入的文件 valid_blobs.sort(key=lambda x: x.updated, reverse=True) latest_blob = valid_blobs[0] # 返回完整GCS路径 return f"gs://{bucket_name}/{latest_blob.name}" # 本地调用示例 if __name__ == "__main__": latest_file = get_latest_gcs_file() print(f"最新存入的文件为:{latest_file}")
Airflow DAG中使用说明
如果要在DAG中调用该逻辑,可以直接把上述函数封装为PythonOperator的执行逻辑,也可以使用Airflow自带的GCSHook简化认证流程,示例如下:
from airflow.providers.google.cloud.hooks.gcs import GCSHook def get_latest_file(**context): gcs_hook = GCSHook(gcp_conn_id="你配置的GCP连接ID") blobs = gcs_hook.list(bucket_name="name", prefix="neth/") # 后续过滤、排序逻辑与上述实现一致,拿到最新文件路径后可通过XCom传递给下游任务
注意事项
- 代码默认按文件的实际上传时间排序取最新文件,如果你需要优先按文件名中的日期后缀匹配,可以修改排序逻辑,提取文件名中的
yymmdd部分转成datetime对象后排序即可 - 批量文件同时上传的场景下,建议给排序逻辑加二级规则,避免时间相同时取到不符合预期的文件
内容的提问来源于stack exchange,提问作者Marthe Sheila Gnitedemg
相关产品推荐
相关产品推荐

