2018年基于Dask的超内存数据集特征工程工作流优化方案问询
Great question—dealing with out-of-core datasets is way more common now than it was in 2018, so there are definitely better, more streamlined tools and workflows available today. Let me walk you through a practical, scalable pipeline that checks all your boxes:
First, forget the clunky pandas chunking from 2018. We’ll use Dask or Vaex—both are built specifically for out-of-core data processing, handle lazy execution, and play nicely with disk-backed datasets.
The first win is specifying tight data types upfront to cut down on memory bloat. Here’s how to do it with both tools:
Dask
import dask.dataframe as dd # Define optimized dtypes to minimize memory footprint dtype_spec = { "user_id": "int32", "transaction_amount": "float32", "category": "category", "timestamp": "datetime64[ns]" } # Load CSV as a lazy Dask DataFrame (no data loaded into memory yet) df = dd.read_csv( "your_large_dataset.csv", dtype=dtype_spec, usecols=["user_id", "transaction_amount", "category", "timestamp"] # Only load needed columns )
Vaex
Vaex converts CSV to a disk-backed format (like HDF5 or Arrow) on first load, making subsequent operations way faster:
import vaex dtype_spec = { "user_id": "int32", "transaction_amount": "float32" } # Load and convert to a memory-mapped DataFrame df = vaex.from_csv( "your_large_dataset.csv", convert=True, # Saves to a temp HDF5/Arrow file automatically dtype=dtype_spec, usecols=["user_id", "transaction_amount", "category", "timestamp"] )
Both tools handle multi-column groupbys without loading the full dataset into memory:
Dask (Lazy Execution)
# Calculate aggregated features grouped by user_id and category grouped_features = df.groupby(["user_id", "category"]).agg({ "transaction_amount": ["mean", "sum", "count"] }).reset_index() # Rename columns for clarity grouped_features.columns = [ "user_id", "category", "avg_transaction_per_cat", "total_spent_per_cat", "transaction_count_per_cat" ]
Dask won’t compute anything yet—this just defines the operation plan.
Vaex (Automatic Broadcast)
Vaex simplifies this even more: it automatically broadcasts groupby results back to the original dataset, no manual merge needed:
# Add groupby features directly to the original DataFrame df["avg_transaction_per_cat"] = df.groupby(["user_id", "category"])["transaction_amount"].mean() df["total_spent_per_cat"] = df.groupby(["user_id", "category"])["transaction_amount"].sum() df["transaction_count_per_cat"] = df.groupby(["user_id", "category"])["transaction_amount"].count()
Dask
Since both df and grouped_features are lazy Dask DataFrames, we can merge them directly and save to disk without loading everything into memory:
# Merge grouped features back to original data enhanced_df = df.merge( grouped_features, on=["user_id", "category"], how="left" ) # Save to disk (Parquet is way better than CSV for future use) enhanced_df.to_parquet( "enhanced_dataset.parquet", write_index=False, compression="snappy" # Balances speed and compression ) # Or split into multiple CSV files if you need that format # enhanced_df.to_csv("enhanced_dataset_part_*.csv", index=False)
Vaex
Since we already added features directly to the DataFrame, just export to a disk-backed format:
# Export to HDF5 (memory-mapped, great for future analysis) df.export_hdf5("enhanced_dataset.hdf5") # Or Parquet for cross-tool compatibility df.export_parquet("enhanced_dataset.parquet")
No need to run full pipelines every time you tweak a feature. Use subsets to test logic quickly:
Dask
# Grab a 1% random sample (computes only this subset into memory) sample_df = df.sample(frac=0.01, random_state=42).compute() # Or take the first 1000 rows sample_df = df.head(1000, compute=True) # Test your groupby logic on the sample first sample_grouped = sample_df.groupby(["user_id", "category"]).agg(...)
Vaex
# Create a lazy sample (no memory loaded until you convert to pandas) sample_df = df.sample(n=1000, random_state=42) # Convert to pandas for quick debugging sample_pd = sample_df.to_pandas_df()
- Ditch CSV for Parquet/Arrow: After your first processing pass, always save to Parquet. It’s compressed, preserves data types, and loads 10x faster than CSV.
- Adjust Block Sizes: In Dask, use
blocksize="64MB"(or higher/lower based on your RAM) to ensure each chunk fits in memory without wasting resources. - Avoid Global Operations: If possible, use rolling window features (supported by both Dask and Vaex) instead of full-dataset stats—they’re more memory-efficient.
- Use Category Types: For string columns with repeated values (like
category), converting tocategorycuts memory usage by 80-90%.
This pipeline is fully out-of-core—no part of the full dataset ever hits your RAM. It’s way more maintainable and efficient than the manual chunking approaches from 2018.
内容的提问来源于stack exchange,提问作者Florian Mutel

