from pyspark.sql import SparkSession # 初始化spark會話 spark = SparkSession \ .builder \ .getOrCreate() spark_df = spark.createDataFrame(pandas_df)
import pandas as pd pandas_df = spark_df.toPandas()
因爲pandas
的方式是單機版的,即toPandas()
的方式是單機版的,因此參考breeze_lsw改爲分佈式版本:sql
import pandas as pd def _map_to_pandas(rdds): return [pd.DataFrame(list(rdds))] def topas(df, n_partitions=None): if n_partitions is not None: df = df.repartition(n_partitions) df_pand = df.rdd.mapPartitions(_map_to_pandas).collect() df_pand = pd.concat(df_pand) df_pand.columns = df.columns return df_pand pandas_df = topas(spark_df)