如何在Spark 2.0(Python2.7.13)中读取Python3.X生成的.pkl字典文件
Got it, let's tackle this problem step by step. The main hurdles here are the Python 2 vs Python 3 pickle compatibility and how to load a single pickle file in a distributed PySpark environment. Here are two feasible solutions depending on your file size and use case:
Solution 1: Directly Load in PySpark (Small to Medium Files)
Since pickle files are single-node by nature, we can distribute the file to all Spark workers using SparkContext.addFile, then load it and convert the dictionary into an RDD/DataFrame. The key fix here is handling the Python 3-to-Python 2 pickle encoding mismatch.
Code Example:
from pyspark import SparkContext import pickle # Initialize SparkContext (adjust config as needed for your cluster) sc = SparkContext.getOrCreate() # Step 1: Distribute the pickle file to all worker nodes # Replace with your actual file path (local or HDFS) sc.addFile("/path/to/your/data.pkl") # Step 2: Define a function to load the pickle file on each worker def load_pkl(_): from pyspark import SparkFiles # Use encoding='latin1' to resolve Python 3 -> Python 2 pickle compatibility with open(SparkFiles.get("data.pkl"), 'rb') as fp: data_dict = pickle.load(fp, encoding='latin1') # Convert dictionary to key-value pairs for RDD compatibility return data_dict.items() # Step 3: Load the data into an RDD # We use a single-element RDD to trigger the load function once per worker (but since it's a single file, it'll load once) key_value_rdd = sc.parallelize([0]).flatMap(load_pkl) # Optional: Convert RDD to DataFrame for SQL operations from pyspark.sql import SparkSession spark = SparkSession(sc) df = spark.createDataFrame(key_value_rdd, schema=["key", "value"]) # Verify the data df.show(5)
Why encoding='latin1'?
Python 3's pickle uses bytes for string serialization, which Python 2's pickle doesn't handle natively. Using latin1 encoding bridges this gap by converting bytes to strings in a way that preserves data integrity for most common data types.
Solution 2: Pre-Convert to Python 2-Compatible Pickle (Large Files or Better Compatibility)
If your pickle file is large, or you want to avoid encoding workarounds, first re-save the file in Python 3 using a pickle protocol that Python 2 supports (Protocol 2 is the highest protocol compatible with Python 2.7).
Step 1: Re-save in Python 3
import pickle # Load the original Python 3 pickle with open("data.pkl", 'rb') as fp: data_dict = pickle.load(fp) # Save with Protocol 2 (Python 2-compatible) with open("data_py2_compatible.pkl", 'wb') as fp: pickle.dump(data_dict, fp, protocol=2)
Step 2: Load in PySpark
Now you can load the converted file without encoding parameters:
sc.addFile("/path/to/data_py2_compatible.pkl") def load_compatible_pkl(_): from pyspark import SparkFiles with open(SparkFiles.get("data_py2_compatible.pkl"), 'rb') as fp: data_dict = pickle.load(fp) return data_dict.items() key_value_rdd = sc.parallelize([0]).flatMap(load_compatible_pkl)
Important Notes
- If your dictionary contains complex objects (like custom classes), you'll need to ensure those classes are available on all Spark workers (e.g., via adding the module to
sc.addPyFile). - For very large dictionaries, consider converting the data to a distributed format like Parquet or CSV in Python 3 first—this will be more efficient for Spark's distributed processing model than using pickle.
内容的提问来源于stack exchange,提问作者Emily Johnson

