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
相关产品推荐
相关产品推荐

