如何基于OneDrive文件同步状态优化BigQuery定时更新脚本
解决方案:优化OneDrive同步检测+增量数据处理
针对你遇到的内存占用过高和OneDrive同步未完成就触发更新的问题,我给你分两个核心部分来解决:
一、可靠检测OneDrive文件是否完全同步
OneDrive同步时会先创建临时文件再替换目标文件,单纯看修改时间会提前触发。这里推荐两种靠谱的检测方式:
1. 文件大小+修改时间双重校验
同步过程中文件大小会持续变化,完成后大小和修改时间都会稳定。可以写一个函数,间隔几秒检查两次文件的大小和修改时间,如果两次结果一致,说明同步完成:
import os import time def is_file_synced(file_path, check_interval=2): # 第一次获取文件信息 first_stat = os.stat(file_path) first_size = first_stat.st_size first_mtime = first_stat.st_mtime time.sleep(check_interval) # 第二次获取文件信息 second_stat = os.stat(file_path) second_size = second_stat.st_size second_mtime = second_stat.st_mtime # 大小和修改时间都不变,视为同步完成 return first_size == second_size and first_mtime == second_mtime
2. 排除OneDrive临时文件
OneDrive同步时会生成带~$前缀的临时文件(比如~$sales.xlsx),或者.tmp后缀的文件,筛选文件时直接排除这些:
valid_files = [ file for file in os.listdir(input_directory) if file.count('-')<=1 and file.endswith('.xlsx') and not file.startswith('~$') and not file.endswith('.tmp') ]
二、增量处理避免全量加载内存
原来的方案全量读取所有文件到内存,导致内存紧张。我们可以通过记录已处理文件的元数据和基于业务主键的增量筛选来优化:
1. 记录已处理文件的状态
创建一个本地记录文件(比如processed_files.json),保存已经处理过的文件的文件名、最后修改时间、文件哈希值,每次运行只处理新增或修改过的文件:
import json import hashlib def get_file_hash(file_path): # 计算文件哈希,判断文件内容是否变化 hash_obj = hashlib.md5() with open(file_path, 'rb') as f: for chunk in iter(lambda: f.read(4096), b''): hash_obj.update(chunk) return hash_obj.hexdigest() # 加载已处理文件记录 try: with open('processed_files.json', 'r') as f: processed = json.load(f) except FileNotFoundError: processed = {} # 筛选需要处理的文件 to_process = [] for file in valid_files: file_path = os.path.join(input_directory, file) if not is_file_synced(file_path): continue # 跳过未同步完成的文件 file_stat = os.stat(file_path) current_mtime = file_stat.st_mtime current_hash = get_file_hash(file_path) # 文件未处理过,或者修改时间/哈希变化了 if file not in processed or processed[file]['mtime'] != current_mtime or processed[file]['hash'] != current_hash: to_process.append(file_path) # 更新记录 processed[file] = {'mtime': current_mtime, 'hash': current_hash} # 保存更新后的记录 with open('processed_files.json', 'w') as f: json.dump(processed, f)
2. 增量读取数据并去重
只读取需要处理的文件,并且如果销售数据有唯一标识(比如order_id)或者时间戳(比如sale_time),可以先从BigQuery获取已存在的最大ID/最新时间,然后只读取新文件中大于该值的数据,进一步减少内存占用:
import pandas as pd from google.cloud import bigquery # 初始化BigQuery客户端 client = bigquery.Client.from_service_account_json('credentials.json') def get_latest_sale_info(): # 查询BigQuery中已有的最新销售时间或最大订单ID query = f"SELECT MAX(sale_time) as latest_time, MAX(order_id) as max_id FROM `{project_id}.{table}`" result = client.query(query).to_dataframe() return result.iloc[0]['latest_time'], result.iloc[0]['max_id'] latest_time, max_id = get_latest_sale_info() incremental_data = [] for file_path in to_process: # 读取文件时只加载需要的列,并且筛选增量数据 df = pd.read_excel(file_path) # 假设用sale_time作为增量判断条件,根据你的实际业务调整 df_filtered = df[df['sale_time'] > latest_time] incremental_data.append(df_filtered) if incremental_data: all_incremental = pd.concat(incremental_data, ignore_index=True).drop_duplicates(subset=['order_id']) # 按唯一ID去重 # 增量上传到BigQuery,用append模式,避免全量替换 all_incremental.to_gbq( project_id=project_id, destination_table=table, credentials=service_account.Credentials.from_service_account_file('credentials.json'), progress_bar=True, if_exists='append' )
额外优化建议
- 如果你的销售数据没有唯一主键,也可以用
hashlib计算每一行的哈希值,把已上传的哈希值存在本地或BigQuery的辅助表中,筛选未上传的行。 - 可以把脚本改成定时任务(比如用Windows任务计划或Linux cron),每1-2小时运行一次,避免持续占用内存。
内容的提问来源于stack exchange,提问作者Hamza
相关产品推荐
相关产品推荐

