Dask Bag存储Pandas DataFrame到磁盘及Featuretools并行特征工程问题
Hey there! Let's tackle your issues one by one, starting with the urgent CSV output problem since that's blocking you right now.
Urgent Issue: Saving Dask Bag of DataFrames to Partitioned CSVs
The error you're hitting with to_textfiles() makes total sense—this method is designed for Bags containing strings or bytes, not Pandas DataFrames. Dask doesn't know how to automatically convert a DataFrame into raw text for this method, hence the TypeError.
Here are two straightforward solutions to save your partitioned DataFrames with proper partition identifiers:
Solution 1: Convert the Bag to a Dask DataFrame (Simplest)
Since each element in your dfs Bag is a single DataFrame (after concatenating 1k features), you can convert it to a Dask DataFrame. Dask's DataFrame to_csv() method automatically handles partitioned output with numbered filenames, exactly what you need:
# Grab a sample DataFrame to define metadata (required for Dask DataFrame) sample_df = dfs.take(1)[0] # Convert Bag to Dask DataFrame ddf = dd.from_bag(dfs, meta=sample_df) # Save to CSV—Dask will generate files like 0000_test.csv, 0001_test.csv, etc. ddf.to_csv('feature_matrices/calculate_matrix/*_test.csv', index=False)
Solution 2: Add Partition Indices to the Bag and Map Save Operations
If you prefer to stick with the Bag interface, you can use enumerate() to attach partition numbers to each DataFrame, then map a save function:
def save_partition(partition_item): partition_num, df = partition_item # Save with partition number in filename df.to_csv(f'feature_matrices/calculate_matrix/{partition_num}_test.csv', index=False) return f"Successfully saved partition {partition_num}" # Add partition indices to each DataFrame in the Bag dfs_with_indices = dfs.enumerate() # Execute save operations (no need to compute to local unless you want the status messages) save_results = dfs_with_indices.map(save_partition).compute()
Background Issues & Optimizations
Now let's address the slower-than-expected execution and "Too many open files" errors you ran into earlier:
1. Fixing "Too many open files" Errors
This is likely caused by either:
- Unclosed file handles in your feature generation code (check if any hidden file operations aren't being cleaned up)
- Worker processes hitting system file descriptor limits
To mitigate this:
- Adjust worker limits: When setting up your LSFCluster, add worker arguments to increase file descriptor limits:
cluster = LSFCluster( # Your existing config worker_extra_args=["ulimit -n 65535"] # Raise open file limit ) - Batch feature saves: You're already grouping 1k features per partition, which reduces the number of files created—keep doing this instead of saving individual feature CSVs.
2. Speeding Up Feature Generation
Your original issue with slow single-feature generation points to overhead in serializing and distributing your EntitySet (es) to workers. Here's how to fix that:
- Avoid global EntitySets: If
esis a global variable in yourmake_featurefunction, Dask has to serialize and send it to workers every time the function runs. Instead, passesas an argument and pre-scatter it to workers:# Assume you have a list of EntitySets (one per data partition) called es_list # Pre-scatter EntitySets to workers to avoid repeated serialization scattered_es = client.scatter(es_list, broadcast=False) # Create a sequence of (feature, es) pairs feature_es_pairs = zip(feature_list, scattered_es) b = db.from_sequence(feature_es_pairs, partition_size=1000) def make_feature(item): feature, es = item feature_name = feature.generate_name() try: feature_matrix = ft.calculate_feature_matrix(feature, entityset=es, n_jobs=1, verbose=1) print(f"Finished generating feature {feature_name}") return feature_matrix except Exception as e: print(f"Could not make feature: {feature_name}") print(f"Error: {str(e)}") return None b = b.map(make_feature) - Use Featuretools' Dask Integration: If you haven't tried it already, Featuretools has built-in support for Dask that can handle parallel feature generation without manual partitioning of features or data. You can use a Dask DataFrame as your entity data source and let Featuretools handle the parallelism:
This eliminates the need for manual feature partitioning and reduces overhead from serializing EntitySets repeatedly.# Convert your dataset to a Dask DataFrame first dask_df = dd.from_pandas(your_full_df, npartitions=...) # Create EntitySet with Dask DataFrame es = ft.EntitySet(id='your_dataset') es = es.add_dataframe( dataframe_name='data', dataframe=dask_df, index='PAT_ENC_CSN_ID', # Your other entity config ) # Generate features and calculate matrix in parallel feature_matrix = ft.calculate_feature_matrix(features, entityset=es, dask_kwargs={'client': client})
3. Avoid Pulling All Data to Local
Your current dfs.compute() pulls all partitioned DataFrames to your local machine, which could use a lot of memory with 200k+ features. Instead, save the data directly from workers using the methods we covered earlier—this keeps data distributed and avoids local memory bottlenecks.
Hope these solutions get you unstuck! Let me know if you need clarification on any part.
内容的提问来源于stack exchange,提问作者An Hoang

