Azure Databricks PySpark:获取外部MySQL表列表并批量创建DataFrame
批量迁移MySQL数据库到Azure Databricks的解决方案
1. 获取MySQL ABC库的表名称列表
你的代码执行show tables后仅返回执行结果数量,是因为没有提取游标中的查询结果。修改如下:
首先在Databricks Notebook中安装依赖(推荐用%pip魔法命令):
%pip install pymysql
提取表名列表的代码:
import pymysql # 建立MySQL连接 conn = pymysql.connect( host='XXXXX', password='XXXXX', port=3306, user='my_username', charset='utf8', # 注意:MySQL无utf6编码,改为utf8/utf8mb4 database='ABC' ) with conn.cursor() as cursor: cursor.execute('SHOW TABLES') # 提取所有表名,fetchall()返回元组列表,取每个元组的第一个元素 table_list = [tbl[0] for tbl in cursor.fetchall()] conn.close() # 验证结果 print(f"共获取到 {len(table_list)} 张表:") print(table_list)
2. 批量创建PySpark DataFrame并写入Databricks开发库
在Databricks中,使用PySpark的JDBC连接MySQL是更高效的批量迁移方式,无需依赖pymysql。具体步骤:
配置JDBC连接参数
# MySQL JDBC连接地址 jdbc_url = "jdbc:mysql://XXXXX:3306/ABC" # 连接属性 jdbc_props = { "user": "my_username", "password": "XXXXX", "driver": "com.mysql.cj.jdbc.Driver" } # 目标Databricks数据库名称 target_database = "your_dev_database"
批量迁移表
for table in table_list: # 从MySQL读取表数据到DataFrame df = spark.read.jdbc(url=jdbc_url, table=table, properties=jdbc_props) # 写入Databricks数据库(默认存储为Delta格式,支持overwrite/append/ignore模式) df.write.mode("overwrite").saveAsTable(f"{target_database}.{table}") print(f"表 {table} 已成功写入Databricks库 {target_database}")
关键注意事项
- JDBC驱动:确保Databricks集群已安装MySQL JDBC驱动,可在集群的"库"页面添加
mysql:mysql-connector-java:8.0.33(或对应版本)。 - 写入模式:
mode("overwrite")会覆盖目标库中同名表,如需保留原有数据可改为mode("append"),或用mode("ignore")跳过已存在的表。 - 字符集修正:原代码中的
charset='utf6'是错误配置,MySQL正确的UTF-8编码为utf8或utf8mb4。
内容的提问来源于stack exchange,提问作者Emma
相关产品推荐
相关产品推荐

