从pyodbc查询结果创建Python Spark DataFrame时遇类型错误
解决Spark无法从pyodbc.Row创建DataFrame的问题
你遇到的报错原因很明确:Spark的createDataFrame()方法没办法直接解析pyodbc.Row这种自定义对象类型,它只能识别Python的基本数据类型(比如元组、列表、字典,或者pandas DataFrame这类结构)。
解决方案:转换pyodbc.Row为可识别的类型
只需要把cursor.fetchall()返回的每一行pyodbc.Row转换成元组或者列表就行,Spark就能顺利推断出数据结构了。
修改后的完整代码如下:
import pyodbc import sys import csv connection = pyodbc.connect("DSN=MySQL") cursor = connection.cursor() cursor.execute("SELECT * FROM customers") # 将每个pyodbc.Row转换为元组 rows = [tuple(row) for row in cursor.fetchall()] column_names = [x[0] for x in cursor.description] print(rows) print(column_names) # 现在可以正常创建Spark DataFrame了 df = spark.createDataFrame(rows, column_names)
另一种可选方案:借助pandas中转
如果你的数据量不大,也可以先把数据转成pandas DataFrame,再转换成Spark DataFrame,这种方式有时候更直观:
import pandas as pd # 前面的连接和查询代码不变 cursor.execute("SELECT * FROM customers") df_pandas = pd.DataFrame.from_records(cursor.fetchall(), columns=[x[0] for x in cursor.description]) df_spark = spark.createDataFrame(df_pandas)
两种方法都能解决你遇到的类型推断问题,选哪种看你的具体场景需求~
内容的提问来源于stack exchange,提问作者NEO
相关产品推荐
相关产品推荐

