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

8GB Twitter压缩JSON数据处理:内存优化与高效存储方案咨询

内存高效的Twitter数据处理优化方案

你的核心问题是将所有筛选后的DataFrame暂存于内存,最终合并时引发内存溢出。以下是兼顾内存与时间效率的针对性优化方案:

方案1:逐文件写入Feather并追加(无需全量内存)

Feather是轻量级列式存储格式,读写速度快,且pandas 1.4.0及以上版本支持追加写入。处理每个压缩文件后直接将筛选结果追加到最终文件,避免全量数据占用内存。

修改后的代码:

import pandas as pd
import pathlib
import timeit
import zipfile

def main():
    output_path = "./tweet_data.feather"
    first_write = not pathlib.Path(output_path).exists()
    
    twitter_files = list(pathlib.Path("TwitterData").iterdir())
    
    for tfilename in twitter_files:
        with zipfile.ZipFile(tfilename, "r") as zfile:
            for filename in zfile.namelist():
                print(f"Processing {filename}")
                with zfile.open(filename) as jfile:
                    # 指定数据类型减少内存占用
                    df = pd.read_json(jfile, lines=True, dtype={
                        "id_str": str,
                        "in_reply_to_status_id_str": str,
                        "quoted_status_id_str": str
                    })
                    # 筛选字段并清理数据
                    sub_df = df[[
                        "id_str", "created_at", "in_reply_to_status_id_str",
                        "text", "entities", "extended_entities",
                        "quoted_status_id_str", "retweeted", "truncated"
                    ]].dropna(subset=["id_str", "created_at", "text"], how="any")\
                      .drop_duplicates(subset="text")
                    
                    # 写入文件:首次覆盖,后续追加
                    sub_df.to_feather(
                        output_path,
                        mode="w" if first_write else "a",
                        compression="zstd"
                    )
                    first_write = False
                    
                    # 及时释放内存
                    del df, sub_df

if __name__ == "__main__":
    duration = timeit.timeit(main, number=1)
    print(f"Total time: {duration:.2f} seconds")

优势:彻底避免全量数据驻留内存,Feather读写速度快,压缩后文件体积更小。
注意:确保所有子DataFrame的列名、数据类型完全一致,否则追加会报错。

方案2:用Dask处理超内存数据集

Dask是专门针对大数据的并行计算框架,自动将数据分块处理,无需一次性加载全量数据到内存。API与pandas高度兼容,适合后续需进行复杂分析的场景。

先安装依赖:pip install dask[complete]

修改后的代码:

import dask.dataframe as dd
import pathlib
import timeit
import zipfile
from dask.diagnostics import ProgressBar

def main():
    twitter_files = list(pathlib.Path("TwitterData").iterdir())
    dfs = []
    
    for tfilename in twitter_files:
        with zipfile.ZipFile(tfilename, "r") as zfile:
            for filename in zfile.namelist():
                print(f"Processing {filename}")
                # Dask读取压缩包内的JSON文件
                df = dd.read_json(
                    f"zip://{filename}::{tfilename}",
                    lines=True,
                    dtype={
                        "id_str": str,
                        "in_reply_to_status_id_str": str,
                        "quoted_status_id_str": str
                    }
                )
                # 筛选清理数据,API与pandas一致
                sub_df = df[[
                    "id_str", "created_at", "in_reply_to_status_id_str",
                    "text", "entities", "extended_entities",
                    "quoted_status_id_str", "retweeted", "truncated"
                ]].dropna(subset=["id_str", "created_at", "text"], how="any")\
                  .drop_duplicates(subset="text")
                dfs.append(sub_df)
    
    # 合并所有数据块
    tweets_dd = dd.concat(dfs)
    # 并行写入Feather文件
    with ProgressBar():
        tweets_dd.to_feather("./tweet_data_dask.feather", compression="zstd")

if __name__ == "__main__":
    duration = timeit.timeit(main, number=1)
    print(f"Total time: {duration:.2f} seconds")

优势:自动分块控制内存占用,后续分析可直接用Dask完成(如统计转发量、全局去重),无需重新加载全量数据。
注意:部分复杂操作的API与pandas略有差异,需参考Dask文档。

方案3:逐行处理+批量写入Parquet

若需极致控制内存,可逐行读取JSON但批量写入Parquet(列存格式,读写效率远高于JSON)。批量写入能减少I/O次数,避免瓶颈。

先安装依赖:pip install pyarrow

代码示例:

import json
import pandas as pd
import pathlib
import timeit
import zipfile
from pyarrow import parquet as pq
from pyarrow import Table

def main():
    output_path = "./tweet_data.parquet"
    batch_size = 10000  # 可根据内存调整批量大小
    batch_data = []
    
    twitter_files = list(pathlib.Path("TwitterData").iterdir())
    
    for tfilename in twitter_files:
        with zipfile.ZipFile(tfilename, "r") as zfile:
            for filename in zfile.namelist():
                print(f"Processing {filename}")
                with zfile.open(filename) as jfile:
                    for line in jfile:
                        try:
                            tweet = json.loads(line)
                            # 跳过关键字段缺失的Tweet
                            if not all(key in tweet for key in ["id_str", "created_at", "text"]):
                                continue
                            # 提取所需字段
                            filtered_tweet = {
                                "id_str": tweet["id_str"],
                                "created_at": tweet["created_at"],
                                "in_reply_to_status_id_str": tweet.get("in_reply_to_status_id_str"),
                                "text": tweet["text"],
                                "entities": tweet.get("entities"),
                                "extended_entities": tweet.get("extended_entities"),
                                "quoted_status_id_str": tweet.get("quoted_status_id_str"),
                                "retweeted": tweet.get("retweeted", False),
                                "truncated": tweet.get("truncated", False)
                            }
                            batch_data.append(filtered_tweet)
                            
                            # 达到批量大小则写入
                            if len(batch_data) >= batch_size:
                                df = pd.DataFrame(batch_data).drop_duplicates(subset="text")
                                table = Table.from_pandas(df)
                                # 首次写入覆盖,后续追加
                                if not pathlib.Path(output_path).exists():
                                    pq.write_table(table, output_path, compression="snappy")
                                else:
                                    pq.write_to_dataset(table, root_path=output_path, compression="snappy")
                                batch_data = []
                                del df, table
                        except json.JSONDecodeError:
                            # 跳过无效JSON行
                            continue
                    # 处理剩余数据
                    if batch_data:
                        df = pd.DataFrame(batch_data).drop_duplicates(subset="text")
                        table = Table.from_pandas(df)
                        if not pathlib.Path(output_path).exists():
                            pq.write_table(table, output_path, compression="snappy")
                        else:
                            pq.write_to_dataset(table, root_path=output_path, compression="snappy")
                        batch_data = []

if __name__ == "__main__":
    duration = timeit.timeit(main, number=1)
    print(f"Total time: {duration:.2f} seconds")

优势:内存占用极低(仅保留一个批次的数据),Parquet压缩率高,适合超大规模数据处理。
注意:批量大小需根据自身内存情况调整,过小会增加I/O次数,过大则提升内存压力。

额外优化建议

  • 指定数据类型:读取时为字符串类字段(如id_str)指定dtype=str,避免pandas自动推断为object类型,减少内存消耗。
  • 启用压缩:Feather和Parquet均支持zstd、snappy等压缩算法,既能减小文件体积,也能提升读写速度。
  • 分阶段去重:先按单个文件去重,后续再通过Dask或分块处理完成全局去重,避免重复数据占用内存。

内容的提问来源于stack exchange,提问作者cicciodevoto

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 19:05:55