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

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 es is a global variable in your make_feature function, Dask has to serialize and send it to workers every time the function runs. Instead, pass es as 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:
    # 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})
    
    This eliminates the need for manual feature partitioning and reduces overhead from serializing EntitySets repeatedly.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 03:59:54