BigQuery拆分导出CSV跨文件时间序列分组应用自定义函数的实现方案
解决方案
方案1:仅用基础pandas实现(无需修改导出规则)
基于你给出的约束条件:数据全局按TimeseriesID、TimeID排序,且单个时间序列最多跨2个文件,我们可以通过仅缓存上一个文件末尾的跨文件时间序列数据实现低内存处理,全程不会加载全量数据到内存:
实现逻辑
- 仅保留上一个文件最后一个
TimeseriesID的所有数据作为缓存,其余数据处理完成后直接释放内存 - 逐个读取文件时先拼接缓存的跨文件数据,再做分组计算
- 当前文件除最后一个
TimeseriesID外的所有分组直接调用自定义函数输出结果,最后一个ID的所有数据存入缓存待下一个文件读取后合并处理 - 最后一个文件的所有分组直接全部处理
示例代码
import pandas as pd import numpy as np # 你的自定义处理函数 def custom_func(x): return np.mean(x) # 按导出顺序填写所有CSV文件路径 file_list = ["gs://your_bucket/part-00000.csv", "gs://your_bucket/part-00001.csv", "gs://your_bucket/part-00002.csv"] # 存储上一个文件末尾跨文件的时间序列数据 cross_file_cache = pd.DataFrame() # 存储计算结果 result = [] for idx, file_path in enumerate(file_list): # 读取当前文件 curr_df = pd.read_csv(file_path) # 拼接上一轮缓存的跨文件数据 if not cross_file_cache.empty: curr_df = pd.concat([cross_file_cache, curr_df], axis=0, ignore_index=True) # 按TimeseriesID分组(注意sort=False保留原有排序,避免额外性能开销) groups = curr_df.groupby("TimeseriesID", sort=False) all_group_ids = list(groups.groups.keys()) is_last_file = (idx == len(file_list) - 1) for group_id in all_group_ids: group_data = groups.get_group(group_id) # 非最后一个文件的最后一个分组,存入缓存待下一轮处理 if not is_last_file and group_id == all_group_ids[-1]: cross_file_cache = group_data.copy() else: # 应用自定义函数,保存结果 calc_res = custom_func(group_data["value"]) result.append({"TimeseriesID": group_id, "calc_result": calc_res}) # 转换为最终结果DataFrame final_result = pd.DataFrame(result)
方案2:调整BigQuery导出规则(从根源避免跨文件问题)
直接在BigQuery导出时按TimeseriesID分区导出,保证同一个时间序列的所有数据都落在同一个文件中,后续处理无需考虑跨文件逻辑,实现更简单:
导出SQL示例
EXPORT DATA OPTIONS ( uri = 'gs://你的存储桶路径/export_prefix_*.csv', format = 'CSV', overwrite = TRUE, header = TRUE, field_delimiter = ',' ) AS ( SELECT * FROM `你的项目.你的数据集.你的表名` ORDER BY TimeseriesID, TimeID ) -- 按TimeseriesID分区,保证同个ID的所有数据在同一个文件 PARTITION BY TimeseriesID;
导出完成后每个CSV文件仅包含同一个TimeseriesID的全部数据,逐个读取文件直接调用custom_func处理即可。
内容的提问来源于stack exchange,提问作者Vinson Ciawandy
相关产品推荐
相关产品推荐

