You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Google Colab中用Python/PySpark高效筛选大数据特定记录

解决方案:Google Colab下高效筛选Azure时序数据集

方案一:纯Python(Pandas逐文件流式处理)

无需额外依赖,完全利用Colab默认环境,核心思路是逐文件读取-筛选-增量写入,全程不加载全量数据到内存,从根源避免OOM。

步骤说明

  1. 基础配置

    • 定义目标VMID列表
    • 设置临时文件存储目录(优先用/tmp,不占用主磁盘空间)
    • 指定最终输出格式(优先选Parquet,比CSV省空间且读写效率更高)
  2. 内存优化

    • 读取文件时仅加载需要的字段(usecols)
    • 为字段指定紧凑数据类型:
      • vmid设为category(若VMID是固定值,比字符串更省内存)
      • CPU字段用float32替代默认的float64
      • timestamp设为datetime64[ns]
  3. 逐文件处理

    • 遍历所有gzip文件,逐个读取
    • 筛选目标VMID的记录
    • 将筛选结果追加到对应VMID的临时CSV文件(避免一次性加载大量数据)
    • 处理完单个文件后立即释放内存,清理临时变量
  4. 最终存储与清理

    • 读取每个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,利用分布式计算能力处理大文件,缓存小数据集后效率极高,完全适配默认资源配置。

步骤说明

  1. 安装并初始化Spark

    • 安装PySpark与Java依赖
    • 配置SparkSession,设置内存参数适配Colab默认资源(分配4-6GB即可)
  2. 加载数据并指定Schema

    • 直接读取所有gzip文件,Spark自动并行处理
    • 手动指定Schema,避免自动推断占用额外内存
  3. 筛选与缓存

    • 筛选目标VMID的记录(数据量极小)
    • 缓存筛选后的数据集,避免重复计算
  4. 按VMID存储

    • 用partitionBy("vmid")按VMID分文件夹存储为Parquet
    • 或手动写出每个VMID的独立文件
  5. 清理资源

    • 取消缓存,停止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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 19:47:31