在Google Colab中用Python/PySpark高效筛选大数据特定记录
解决方案:Google Colab下高效筛选Azure时序数据集
方案一:纯Python(Pandas逐文件流式处理)
无需额外依赖,完全利用Colab默认环境,核心思路是逐文件读取-筛选-增量写入,全程不加载全量数据到内存,从根源避免OOM。
步骤说明
基础配置
- 定义目标VMID列表
- 设置临时文件存储目录(优先用
/tmp,不占用主磁盘空间) - 指定最终输出格式(优先选Parquet,比CSV省空间且读写效率更高)
内存优化
- 读取文件时仅加载需要的字段(
usecols) - 为字段指定紧凑数据类型:
vmid设为category(若VMID是固定值,比字符串更省内存)- CPU字段用
float32替代默认的float64 timestamp设为datetime64[ns]
- 读取文件时仅加载需要的字段(
逐文件处理
- 遍历所有gzip文件,逐个读取
- 筛选目标VMID的记录
- 将筛选结果追加到对应VMID的临时CSV文件(避免一次性加载大量数据)
- 处理完单个文件后立即释放内存,清理临时变量
最终存储与清理
- 读取每个VMID的临时CSV,转换为Parquet格式存储
- 删除临时CSV文件与目录,释放磁盘空间
代码示例
import pandas as pd import gc import os from glob import glob # 1. 基础配置 TARGET_VMIDS = ["vmid_1", "vmid_2", "vmid_3", "vmid_4"] # 替换为你的目标VMID TEMP_DIR = "/tmp/vm_temp_files" OUTPUT_DIR = "/content/vm_output" # 创建目录 os.makedirs(TEMP_DIR, exist_ok=True) os.makedirs(OUTPUT_DIR, exist_ok=True) # 初始化临时文件(写入表头) schema = pd.DataFrame(columns=["timestamp", "vmid", "mincpu", "maxcpu", "avgcpu"]) for vmid in TARGET_VMIDS: temp_path = os.path.join(TEMP_DIR, f"{vmid}.csv") schema.to_csv(temp_path, index=False, mode="w") # 2. 逐文件处理 gzip_files = glob("/path/to/your/gzip/files/*.gz") # 替换为你的文件路径 for file in gzip_files: # 读取文件,指定紧凑数据类型 df = pd.read_csv( file, compression="gzip", usecols=["timestamp", "vmid", "mincpu", "maxcpu", "avgcpu"], dtype={ "vmid": "category", "mincpu": "float32", "maxcpu": "float32", "avgcpu": "float32" }, parse_dates=["timestamp"] ) # 筛选目标VMID filtered = df[df["vmid"].isin(TARGET_VMIDS)] # 追加到对应临时文件 for vmid in TARGET_VMIDS: vmid_data = filtered[filtered["vmid"] == vmid] temp_path = os.path.join(TEMP_DIR, f"{vmid}.csv") vmid_data.to_csv(temp_path, index=False, mode="a", header=False) # 强制释放内存 del df, filtered, vmid_data gc.collect() # 3. 转换为Parquet并清理临时文件 for vmid in TARGET_VMIDS: temp_path = os.path.join(TEMP_DIR, f"{vmid}.csv") output_path = os.path.join(OUTPUT_DIR, f"{vmid}.parquet") final_df = pd.read_csv(temp_path, parse_dates=["timestamp"], dtype={"vmid": "category"}) final_df.to_parquet(output_path, index=False) # 删除临时文件 os.remove(temp_path) # 清理临时目录 os.rmdir(TEMP_DIR)
注意事项
- 用
!df -h查看Colab磁盘空间,不足时可删除无关文件释放空间 - 若单个gzip文件仍过大,可添加
chunksize=100000参数分块读取,逐块处理
方案二:PySpark(适合超大规模数据)
Colab可快速安装PySpark,利用分布式计算能力处理大文件,缓存小数据集后效率极高,完全适配默认资源配置。
步骤说明
安装并初始化Spark
- 安装PySpark与Java依赖
- 配置SparkSession,设置内存参数适配Colab默认资源(分配4-6GB即可)
加载数据并指定Schema
- 直接读取所有gzip文件,Spark自动并行处理
- 手动指定Schema,避免自动推断占用额外内存
筛选与缓存
- 筛选目标VMID的记录(数据量极小)
- 缓存筛选后的数据集,避免重复计算
按VMID存储
- 用
partitionBy("vmid")按VMID分文件夹存储为Parquet - 或手动写出每个VMID的独立文件
- 用
清理资源
- 取消缓存,停止SparkSession,释放内存
代码示例
# 1. 安装PySpark依赖 !pip install pyspark !apt-get install openjdk-8-jdk-headless -qq > /dev/null # 2. 初始化SparkSession import os os.environ["JAVA_HOME"] = "/usr/lib/jvm/java-8-openjdk-amd64" from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, TimestampType, FloatType spark = SparkSession.builder \ .appName("AzureVMFilter") \ .config("spark.driver.memory", "4g") \ .getOrCreate() # 3. 定义Schema schema = StructType([ StructField("timestamp", TimestampType(), True), StructField("vmid", StringType(), True), StructField("mincpu", FloatType(), True), StructField("maxcpu", FloatType(), True), StructField("avgcpu", FloatType(), True) ]) # 4. 加载所有gzip文件 df = spark.read.csv( "/path/to/your/gzip/files/*.gz", schema=schema, header=True, compression="gzip" ) # 5. 筛选目标VMID并缓存 TARGET_VMIDS = ["vmid_1", "vmid_2", "vmid_3", "vmid_4"] filtered_df = df.filter(df.vmid.isin(TARGET_VMIDS)) filtered_df.cache() # 缓存小数据集,提升后续操作速度 # 6. 按VMID存储为Parquet output_path = "/content/vm_spark_output" filtered_df.write \ .mode("overwrite") \ .partitionBy("vmid") \ .parquet(output_path) # 可选:手动写出每个VMID的独立文件 for vmid in TARGET_VMIDS: single_vmid_df = filtered_df.filter(filtered_df.vmid == vmid) single_vmid_df.write \ .mode("overwrite") \ .parquet(f"{output_path}/{vmid}.parquet") # 7. 清理资源 filtered_df.unpersist() spark.stop()
注意事项
partitionBy会自动为每个VMID创建独立文件夹,方便后续读取- 若内存仍不足,可调整
spark.driver.memory参数(Colab默认内存约12GB,不要超过8GB)
内容的提问来源于stack exchange,提问作者Mario
相关产品推荐
相关产品推荐

