如何将数据从 Pandas 数据帧分块加载到 Spark 数据帧

Gau*_*ama 3 python pandas apache-spark pyspark

我已经使用以下内容通过 pyodbc 连接以块的形式读取数据:

import pandas as pd
import pyodbc
conn = pyodbc.connect("Some connection Details")
sql = "SELECT * from TABLES;"
df1 = pd.read_sql(sql,conn,chunksize=10)
Run Code Online (Sandbox Code Playgroud)

现在我想使用以下内容将所有这些块读入一个单一的火花数据帧:

i = 0
for chunk in df1:
    if i==0:
        df2 = sqlContext.createDataFrame(chunk)
    else:
        df2.unionAll(sqlContext.createDataFrame(chunk))
    i = i+1
Run Code Online (Sandbox Code Playgroud)

问题是当我做 a 时,df2.count()我得到的结果是 10,这意味着只有 i=0 的情况在工作。这是 unionAll 的错误吗?我在这里做错了吗??

ber*_*nie 5

文档.unionAll()说明它返回一个新的数据帧,因此您必须重新分配给数据帧df2:

i = 0
for chunk in df1:
    if i==0:
        df2 = sqlContext.createDataFrame(chunk)
    else:
        df2 = df2.unionAll(sqlContext.createDataFrame(chunk))
    i = i+1
Run Code Online (Sandbox Code Playgroud)

此外,您可以改为使用enumerate()以避免必须自己管理i变量:

for i,chunk in enumerate(df1):
    if i == 0:
        df2 = sqlContext.createDataFrame(chunk)
    else:
        df2 = df2.unionAll(sqlContext.createDataFrame(chunk))
Run Code Online (Sandbox Code Playgroud)

此外.unionAll(),.unionAll()已弃用状态的文档,现在您应该.union()在 SQL 中使用类似于 UNION ALL 的行为:

for i,chunk in enumerate(df1):
    if i == 0:
        df2 = sqlContext.createDataFrame(chunk)
    else:
        df2 = df2.union(sqlContext.createDataFrame(chunk))
Run Code Online (Sandbox Code Playgroud)

编辑:
此外,我将停止进一步说,但在我进一步说之前不会说:正如@zero323 所说,让我们不要.union()在循环中使用。让我们做一些类似的事情:

def unionAll(*dfs):
    ' by @zero323 from here: http://stackoverflow.com/a/33744540/42346 '
    first, *rest = dfs  # Python 3.x, for 2.x you'll have to unpack manually
    return first.sql_ctx.createDataFrame(
        first.sql_ctx._sc.union([df.rdd for df in dfs]),
        first.schema
    )

df_list = []
for chunk in df1:
    df_list.append(sqlContext.createDataFrame(chunk))

df_all = unionAll(df_list)
Run Code Online (Sandbox Code Playgroud)