You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 11:37:05