安装pymysql和MySQLdb的版本差异
Python 3下安装pymysql很简单,一把梭就能搞定:
pip3 install pymysql
如果服务器上有多个Python版本,用python3 -m pip install pymysql确保装到当前解释器对应的包里,装完后验证一下:
python3 -c "import pymysql; print(pymysql.__version__)"
顺利输出版本号就说明装好了,如果报错提示找不到pip,先装pip或改用系统包管理器安装,比如CentOS下yum install python3-pip,Ubuntu下apt install python3-pip。
Python 2环境则有些区别,Python 2时代的MySQL驱动叫MySQLdb,模块名是_mysql,安装包名是MySQL-python,但MySQLdb在Python 3下会编译失败,所以如果你用的Python 3,直接用pymysql就能覆盖所有核心操作,相当于pymysql是MySQLdb的纯Python替代品,少了一堆编译依赖,安装更稳。
| 环境 | 推荐驱动 | 安装命令 | 导入名 |
|---|---|---|---|
| Python 3 | pymysql | pip3 install pymysql | import pymysql |
| Python 2 | MySQL-python | pip install MySQL-python | import MySQLdb |
pymysql连接MySQL的完整实操流程
先说连接参数,pymysql的connect()函数需要主机地址、端口、用户名、密码、数据库名这五件套,端口默认3306,如果MySQL跑在Docker容器里,注意端口映射是否正确,这也是服务器上容易出问题的环节。
核心连接配置如下:
import pymysql
conn = pymysql.connect(
host='127.0.0.1',
port=3306,
user='root',
password='你的密码',
charset='utf8mb4',
database='spark_result'
)
cursor = conn.cursor()

接着执行查询并读取结果:
cursor.execute("SELECT FROM user_behavior LIMIT 10")
rows = cursor.fetchall()
for row in rows:
print(row)
写完insert或update语句后必须conn.commit(),不提交的话数据不会真正落库,这个坑很多人踩过一次才记住,批量插入时用executemany比逐条execute更快,代码层面性能差距就体现在这里:
data = [(1, 'a'), (2, 'b'), (3, 'c')]
cursor.executemany("INSERT INTO test_table (id, name) VALUES (%s, %s)", data)
conn.commit()
用完记得关闭游标和连接,防止连接数泄漏,连接数满了之后,其他脚本再连就会报Too many connections,直接把MySQL拒绝服务。
连接池解决Spark结果回读时的并发瓶颈
Spark调用Python脚本回读MySQL里的结果表,如果一次性启动多个并发任务,每个任务都新建连接,数据库很容易被压垮,常见做法是用DBUtils.PooledDB建一个连接池,把连接变成可复用的资源:
from dbutils.pooled_db import PooledDB
pool = PooledDB(
creator=pymysql,
maxconnections=10,
mincached=2,
host='127.0.0.1',
user='root',
password='你的密码',
database='spark_result',
charset='utf8mb4'
)
def query(sql):
conn = pool.connection()
cursor = conn.cursor()
cursor.execute(sql)
result = cursor.fetchall()
cursor.close()
conn.close()
return result

连接池的思想是当某个连接被close()之后,连接不真正断开,而是回到池里等待下一次复用,这样Spark的多个executor节点可以共享同一组数据库连接,而不需要为每个task新建一次TCP握手。
运行脚本过程中的典型报错与排查方向
Spark写MySQL和pymysql连接MySQL,虽然属于两条技术路线,但报错有相当一部分集中在连接层和驱动层。
驱动程序没加载
通常包含No suitable driver found,解决办法就是在Spark提交命令里加上--jars /path/to/mysql-connector-java.jar,或者把jar包丢进$SPARK_HOME/jars/,还有一个隐蔽的问题是Spark集群模式下的Driver和Executor不在同一台机器,jar包每台节点都要有。
连接超时
MySQL所在服务器防火墙没放行3306端口,或者MySQL绑定地址是0.0.1而不是0.0.0,远端Spark节点压根连不上,排查命令直接用telnet 服务器IP 3306验证端口通不通,MySQL的my.cnf里bind-address需要改成0.0.0,并确认用户权限允许远程登录:
GRANT ALL PRIVILEGES ON . TO 'root'@'%' IDENTIFIED BY '密码'; FLUSH PRIVILEGES;
中文乱码
连接URL里加上characterEncoding=utf8,MySQL建表时指定DEFAULT CHARSET=utf8mb4,pymysql连接参数里charset='utf8mb4',三处对齐基本不会乱码。
全局表名大小写问题
Linux MySQL默认区分表名大小写,Spark写入

UserBehavior表时,如果MySQL里建的表是user_behavior,会提示表不存在,最稳妥的方案是lower_case_table_names=1,或让Spark脚本和建表语句统一表名。
常见问题解答
Spark作业结果写MySQL时机选在行动算子位置
需要在截图或后续查询结果之前先落库,注意df.write本身是行动操作,触发后会立即执行,不要在foreach循环里再次调用写库,否则会重复写入,造成大量脏数据。
pymysql脚本执行SQL脚本文件
不想在命令行用source跑SQL文件,也可以在python脚本里读取文件内容再一次性执行:
with open('/opt/sql/init.sql', 'r', encoding='utf-8') as f:
sql = f.read()
cursor.execute(sql)
conn.commit()
多个语句混在一起时,部分pymysql版本会报语法错,此时需要按分号拆分逐条执行,或者改用pymysql.cursors的execute多次调用。
pymysql和SQLAlchemy怎么选
当数据量小、查询逻辑简单、只需要快速读写MySQL时,pymysql直接写SQL即可,当SQL逻辑复杂,或需要ORM映射时,SQLAlchemy是更好的选择,它底层默认不依赖pymysql,需要配合mysql+pymysql://这个连接串来用,Spark的DataFrame写入MySQL时,同样推荐直接用JDBC而非pymysql,因为JDBC是Spark原生支持的路径,性能和稳定性都有保证。
综合来看,Spark作业结果存储MySQL的链路并不复杂,关键的几块在于JDBC驱动配置、MySQL服务器脚本运行方式、以及Python侧pymysql模块的安装和连接池管理,把这几个环节逐一打通,Spark到MySQL的落库与回读流程就能顺畅跑起来。