Python结合Google Cloud提取远程文件新增数据行方法
增量提取新增数据实现方案
这个数据源的CSV文件只会追加新日期的记录,不会修改历史数据,所以不需要换JSON格式,直接用日期比对就能精准拿到新增行,逻辑简单可靠性高。
核心流程:
- 每次函数触发时,先读取GCS上已存储的旧文件,提取其中已经入库的最大日期值
- 下载最新的全量CSV到临时目录
- 从新文件中筛选出日期大于已存最大日期的行,即为纯新增数据
- 将新增行追加写入GCS的原有文件即可,无需全量覆盖
注:该CSV的
date列是标准YYYY-MM-DD格式的字符串,直接做字符串大小比较就能判断日期先后,不需要额外转日期类型,出错概率极低。
核心代码实现
给你两个版本,选一个用就行,逻辑完全一致。
版本1:pandas实现(代码最简洁,推荐)
Cloud Function的Python 3.9+运行时已经预装pandas,不需要额外打包依赖,适合编程基础薄弱的情况:
import io import pandas as pd from google.cloud import storage # 以下是你原有逻辑里已经初始化好的变量,直接复用即可 # client = storage.Client() # bucket = client.get_bucket(bucket_name) # cf_path = '/tmp/{}'.format(file_name) # 1. 读取GCS上已有的旧文件,获取已存储的最新日期 old_file_blob = bucket.get_blob(file_name) last_stored_date = "1900-01-01" # 初始默认值,对应第一次运行无历史文件的场景 if old_file_blob: old_file_content = old_file_blob.download_as_bytes() old_df = pd.read_csv(io.BytesIO(old_file_content)) last_stored_date = old_df["date"].max() # 2. 读取刚下载到临时目录的最新全量CSV new_full_df = pd.read_csv(cf_path) # 3. 筛选出所有新增行 new_rows_df = new_full_df[new_full_df["date"] > last_stored_date] # -------------------------- # 到这一步就拿到了所有需要追加的新增数据,你可以自己实现后续写入GCS的逻辑 # 转可直接写入的CSV格式参考:new_rows_df.to_csv(index=False, header=not old_file_blob) # header参数说明:第一次运行无旧文件时写表头,后续追加时不写表头,避免重复 # --------------------------
版本2:Python标准库实现(无第三方依赖)
如果不想用pandas,直接用Python内置的csv模块也能实现,不需要装任何额外包:
import io import csv from google.cloud import storage # 复用你原有逻辑里的初始化变量 # client = storage.Client() # bucket = client.get_bucket(bucket_name) # cf_path = '/tmp/{}'.format(file_name) old_file_blob = bucket.get_blob(file_name) last_stored_date = "1900-01-01" csv_headers = None # 1. 读取旧文件拿到最新已存日期 if old_file_blob: old_file_content = old_file_blob.download_as_text() old_reader = csv.DictReader(io.StringIO(old_file_content)) csv_headers = old_reader.fieldnames for row in old_reader: if row["date"] > last_stored_date: last_stored_date = row["date"] # 2. 读取新下载的全量文件,筛选新增行 new_rows = [] with open(cf_path, "r", encoding="utf-8") as f: new_reader = csv.DictReader(f) # 第一次运行没有拿到旧表头的话,用新文件的表头 if not csv_headers: csv_headers = new_reader.fieldnames for row in new_reader: if row["date"] > last_stored_date: new_rows.append(row) # -------------------------- # 这里new_rows就是所有新增的行字典列表,csv_headers是CSV表头 # 你可以自行实现追加写入GCS的逻辑 # --------------------------
注意事项
- 不需要使用JSON格式数据源:JSON文件体积比CSV大3倍以上,解析速度更慢,处理逻辑反而更复杂,用CSV是最优选择
- 追加写入时记得判断是否是第一次运行:第一次运行要写入表头,后续追加只写数据行,避免CSV里出现多个表头导致后续读取出错
- 该数据源不会修改历史数据,只会按天追加前一天的全球记录,用日期比对的方式不会漏数据也不会产生重复
- 如果担心极端情况同一天数据多次修正,可以把比对条件改成
>=,写入前对同日期数据做去重即可,日常定时每天跑一次的场景不需要额外处理。
内容的提问来源于stack exchange,提问作者reign
相关产品推荐
相关产品推荐

