从 MySQL 开发到 Spark 实战

思维转换与入门到实践完整教程
PySpark 3.5.x JDK 11 Python 3.8+ 面向 MySQL 开发者

目录

  1. 为什么要从 MySQL 转向 Spark —— 能力边界与思维转换
  2. 概念映射表 —— 你已有的 MySQL 知识如何迁移到 Spark
  3. 环境搭建 —— 从 MySQL 客户端到 Spark 运行环境
  4. 第一个 Spark 程序 —— MySQL 开发者的习惯写法 vs Spark 写法
  5. SQL 执行引擎的本质区别 —— SQLAlchemy vs SparkSQL 深度对比
  6. DataFrame 与 SparkSQL —— 生产核心
  7. 多数据源联邦查询 —— Spark 的跨库能力
  8. 表名前缀的真相 —— ods. / mysql_dim. 是什么
  9. Spark 如何识别数据源 —— 元数据驱动而非命名约定
  10. 完整实战 —— 离线日志统计(端到端)
  11. 常见误区与 FAQ
  12. 进阶学习路线

01为什么要从 MySQL 转向 Spark

作为 MySQL 开发者,你已经熟悉了关系型数据库的建表、写 SQL、用 ORM 框架操作数据的完整流程。但当数据量从百万级增长到亿级、十亿级时,MySQL 的单机架构会遭遇性能瓶颈——慢查询、锁表、CPU/IO 打满。这就是你需要学习 Spark 的根本原因。

MySQL 的世界

  • 计算和存储耦合在同一台服务器
  • SQL 语句发给数据库,数据库自己执行
  • Python/Java 只是"传话筒",拼接 SQL、接收结果
  • 百万级数据表现优秀,千万级开始吃力
  • 完整 ACID 事务,支持行级增删改
  • 单库查询,跨库 JOIN 需要应用层处理

Spark 的世界

  • 计算与存储分离,Spark 只负责计算
  • SQL 语句交给 Spark Catalyst 优化器解析执行
  • Spark 集群多台机器并行计算,数据拉到 Executor 内存
  • 十亿级离线数据常规处理,靠横向扩容提升性能
  • 面向批量,不适合高频在线事务
  • 原生支持跨 Hive / MySQL / HBase / 文件 联邦查询
核心认知转换

MySQL 开发者的习惯:写 SQL → 发给数据库 → 数据库算好返回结果

Spark 开发者的思维:写 SQL → Spark 自己解析优化 → 拆分任务分发到集群多节点并行计算 → 汇总结果

差距不在"怎么写 SQL"(语法几乎一样),而在"谁来执行这条 SQL、在哪台机器上算、数据存在哪里"。

Spark 是什么?

Spark 是通用分布式大数据计算引擎,支持离线批处理、实时流、SQL 分析、机器学习,替代老旧 MapReduce,内存计算速度提升 10~100 倍。它本身不持久化存储任何数据——数据存在 HDFS / 对象存储 / Hive / MySQL / HBase 等外部系统中,Spark 只负责把数据读出来、在集群内存中并行计算、再把结果写回去。

02概念映射表

下面这张表是你最需要的——把已有的 MySQL 知识直接映射到 Spark 概念上,快速建立认知框架。

概念维度 MySQL(你熟悉的) Spark(需要建立的)
计算载体 数据库服务端单机执行 Spark 集群多机器分布式并行计算
数据存储 数据库本地磁盘 / SSD,行式存储 HDFS / 对象存储,列式 Parquet/ORC
SQL 执行者 MySQL 引擎自己解析执行 Spark Catalyst 优化器解析,生成分布式执行计划
SQL 语法 标准 SQL 标准 SQL(几乎相同
客户端连接 JDBC/ORM 连接单库 Connector 统一对接多种数据源,支持联邦查询
数据规模上限 千万级性能急剧衰减 十亿级常规处理,横向扩容
表概念 数据库中的物理表 DataFrame(逻辑表)/ Hive 元数据表 / 临时视图
Schema CREATE TABLE 定义列 DataFrame 自带 Schema,或从 Hive Metastore 读取
查询入口 mysql CLI / Navicat / SQLAlchemy SparkSession(2.0+ 统一入口)
事务支持 完整 ACID 事务 有限事务,面向批量,不支持高频单行更新
数据写入 INSERT/UPDATE/DELETE 行级操作 批量写入:overwrite 覆盖 / append 追加
缓存机制 数据库自身 Buffer Pool cache() 将中间结果缓存到集群内存
优化器 MySQL 优化器(单机,受内存 CPU 限制) Catalyst 分布式优化器(谓词下推、分区裁剪、Join 优化)
运行模式 单进程串行 local / Standalone / YARN / K8s 集群模式

03环境搭建

MySQL 开发者习惯用 mysql CLI 或 Navicat 连数据库。Spark 的环境稍有不同——你需要一个 Spark 运行环境来执行代码。以下是三种方案,按推荐程度排序。

方案 1:pip 安装 PySpark(推荐新手)

最简单的方式,无需下载 Spark 安装包,一条命令搞定。

前置依赖

# 安装 PySpark
pip install pyspark==3.5.0

# 启动交互式终端(类似 mysql CLI)
pyspark

出现 Welcome to PySpark 3.5.0 即成功。启动后自动生成两个全局变量:

方案 2:完整 Spark 安装包(Linux/Mac)

# 下载解压
wget https://archive.apache.org/dist/spark/spark-3.5.0/spark-3.5.0-bin-hadoop3.tgz
tar -xzf spark-3.5.0-bin-hadoop3.tgz -C /opt/
ln -s /opt/spark-3.5.0-bin-hadoop3 /opt/spark

# 环境变量
export SPARK_HOME=/opt/spark
export PATH=$SPARK_HOME/bin:$PATH

# 验证
pyspark

方案 3:Windows(避坑要点)

Windows 避坑
  1. 解压 Spark 到无中文、无空格路径,如 D:\spark
  2. 下载对应 Hadoop 版本的 winutils.exe,配置 HADOOP_HOME
  3. 环境变量添加 SPARK_HOME,Path 追加 %SPARK_HOME%\bin
  4. cmd 执行 pyspark

方案 4:在线免安装(Jupyter/Colab)

# Jupyter Notebook 或 Google Colab 中直接执行
pip install pyspark

from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("demo").getOrCreate()

04第一个 Spark 程序

下面用同一个需求——"统计用户行为日志中各页面的 UV"——对比 MySQL 写法和 Spark 写法,让你直观感受差异。

MySQL + SQLAlchemy 写法(你熟悉的)

from sqlalchemy import create_engine

engine = create_engine(
    "mysql+pymysql://user:pass@host/db"
)

# SQL 交给 MySQL 执行
sql = """
SELECT page, COUNT(DISTINCT uid) AS uv
FROM user_log
GROUP BY page
"""

result = engine.execute(sql)
for row in result:
    print(row)

MySQL 服务端执行全部计算,Python 只接收结果。

Spark SQL 写法(需要掌握的)

from pyspark.sql import SparkSession

# 1. 创建 SparkSession(固定模板)
spark = SparkSession.builder \
    .appName("UV统计") \
    .master("local[*]") \
    .enableHiveSupport() \
    .getOrCreate()

# 2. SQL 交给 Spark 集群执行
spark.sql("""
SELECT page, COUNT(DISTINCT uid) AS uv
FROM ods.user_log
GROUP BY page
""").show()

# 3. 关闭
spark.stop()

Spark 自己解析 SQL,分发到集群并行计算。

关键发现

SQL 语句本身几乎一模一样!区别在于:engine.execute(sql) 把 SQL 发给 MySQL 执行;spark.sql(sql) 把 SQL 交给 Spark 自己的 Catalyst 优化器解析执行。

Spark 程序固定模板(可复制到 .py 文件运行)

from pyspark.sql import SparkSession

# ===== 固定模板:所有 PySpark 代码的开头 =====
spark = SparkSession.builder \
    .appName("MyApp") \
    .master("local[*]") \
    .getOrCreate()

# ===== 你的业务逻辑 =====
df = spark.read.csv("data.csv", header=True, inferSchema=True)
df.createOrReplaceTempView("my_table")
spark.sql("SELECT * FROM my_table LIMIT 10").show()

# ===== 固定结尾 =====
spark.stop()

运行命令:

spark-submit --master local[*] my_app.py

05SQL 执行引擎的本质区别

这是从 MySQL 转向 Spark 最核心的知识点——理解 SQL 在两种环境下"谁来执行"的根本差异。

执行流程对比

1
Python 代码
拼接 SQL 字符串
2
SQLAlchemy
转发 SQL
3
MySQL 服务端
解析 + 执行 + 返回

MySQL 路径:全部计算压在数据库单机

1
Python 代码
拼接 SQL 字符串
2
Spark Catalyst
解析优化生成计划
3
Spark 集群
多 Executor 并行计算
4
外部存储
HDFS/MySQL 只提供数据

Spark 路径:计算收拢到 Spark 集群,存储只提供原始数据

分维度详细对比

对比维度 SQLAlchemy(普通数据库 SQL) Spark SQL
计算载体 数据库服务端(单节点) Spark 集群多机器分布式并行
SQL 作用域 仅能查询单一关系库 支持联邦查询:Hive + MySQL + ES + Kafka
优化器 数据库自带优化器,受单机限制 Catalyst 分布式优化器,谓词下推、分区裁剪
数据处理上限 千万级性能急剧衰减 十亿级常规处理
执行模式 单进程串行,无分布式分片 自动分片、多 Executor 并行
内存机制 数据库缓存,Python 客户端不参与 中间结果缓存到集群内存 cache()
适用场景 业务系统增删改查、小维度查询 数仓 ETL、海量日志分析、跨源统计
事务支持 完整 ACID 事务 有限事务,面向批量
最容易混淆的点:Spark 通过 JDBC 读取 MySQL 时

很多人误以为 spark.read.jdbc() 等价于 SQLAlchemy,实际完全不同:

  • SQLAlchemy:一次性拉取全量数据到 Python 内存,单机处理
  • Spark JDBC:支持分区下推,多并发分片从 MySQL 拉取 → 数据落到 Spark 集群分布式内存 → 过滤、聚合、Join 全部在集群并行执行 → 仅把过滤条件下推给 MySQL 减少传输

简言之:Spark JDBC 是"把 MySQL 数据搬到 Spark 集群算",SQLAlchemy 是"让 MySQL 自己算"。

06DataFrame 与 SparkSQL

对 MySQL 开发者来说,DataFrame 就是一张分布式数据表——带字段名、数据类型,兼容标准 SQL,性能最优。这是生产环境的核心,必学。

Spark 三层数据抽象

DataFrame / Dataset 生产首选

带 Schema 的结构化表,支持 SQL 和 Catalyst 自动优化,性能远高于 RDD。相当于 MySQL 中的表。

RDD 底层原理

无结构的原始分布式集合,底层 API。复杂底层开发才用,日常不需要。理解原理即可。

共享存储

HDFS / 对象存储、Hive 元数据。相当于 MySQL 的磁盘存储,但 Spark 不拥有数据。

6.1 手动创建 DataFrame

# 类似 MySQL 的 CREATE TABLE ... VALUES
data = [
    ("Alice", "北京", "2024-01-01", "home"),
    ("Bob", "上海", "2024-01-01", "search"),
    ("Alice", "北京", "2024-01-02", "cart"),
]

df = spark.createDataFrame(data, ["name", "city", "date", "page"])
df.show()
# +-----+----+----------+------+
# | name|city|      date|  page|
# +-----+----+----------+------+
# |Alice|北京|2024-01-01|  home|
# |  Bob|上海|2024-01-01|search|
# |Alice|北京|2024-01-02|  cart|
# +-----+----+----------+------+

6.2 DSL 链式操作(类 pandas / ORM 写法)

# 类似 SQLAlchemy ORM 的链式查询
df.filter(df.city == "北京") \
  .groupBy("page") \
  .count() \
  .show()

# +----+-----+
# |page|count|
# +----+-----+
# |home|    1|
# |cart|    1|
# +----+-----+

6.3 SparkSQL 标准 SQL 查询(最常用)

核心模式

将 DataFrame 注册为临时视图 → 直接写 SQL 查询。这和 MySQL 中写 SQL 的体验几乎一致。

# 注册临时视图(类似给表起个别名)
df.createOrReplaceTempView("user_log")

# 直接写标准 SQL,语法和 MySQL 基本一致
result = spark.sql("""
SELECT page, COUNT(DISTINCT name) AS uv
FROM user_log
GROUP BY page
ORDER BY uv DESC
""")
result.show()

6.4 读写外部文件

# 读 CSV(类似 MySQL 的 LOAD DATA INFILE)
df = spark.read.csv("hdfs://xxx/data/user_log.csv",
                     header=True, inferSchema=True)

# 读 Parquet(列式存储,Spark 性能最优格式)
df = spark.read.parquet("hdfs://xxx/data/user_log.parquet")

# 写入结果(类似 INSERT INTO ... SELECT)
result.write.mode("overwrite") \
       .parquet("hdfs://xxx/output/uv_result")
写入模式 mode 行为 MySQL 类比
overwrite覆盖已有数据TRUNCATE + INSERT
append追加到已有数据INSERT INTO
ignore存在则不写入INSERT IGNORE
errorifexists存在直接报错(默认)CREATE TABLE(已存在则报错)

07多数据源联邦查询

这是 MySQL 开发者最需要理解的能力差异——MySQL 只能查自己库里的表,而 Spark 可以在一条 SQL 中同时关联 Hive 表、MySQL 表、HBase 表、文件,所有计算在 Spark 集群内完成。

场景:三表关联查询

-- Hive 离线日志 + MySQL 业务维度 + HBase 用户画像 三表关联
SELECT
    a.page,
    b.user_name,
    c.user_tag,
    COUNT(DISTINCT a.uid) AS uv
FROM ods.user_log a
JOIN mysql_dim.user b ON a.uid = b.uid
JOIN hbase_user c ON a.uid = c.uid
GROUP BY a.page, b.user_name, c.user_tag
MySQL 开发者的第一反应

在 MySQL 中这不可能——你没法在一条 SQL 里 JOIN 另一个数据库系统的表。通常做法是分别查两张表,在 Python 本地内存 JOIN,数据量大直接 OOM。

Spark 的做法
  1. 分别通过不同 Connector 拉取三类数据
  2. 数据全部加载到 Spark 集群内存 / 磁盘
  3. 在 Spark 集群内部完成分布式 JOIN、去重、聚合
  4. MySQL / HBase / Hive 只负责提供原始数据,不参与复杂计算

四种数据源的接入方式

① Hive 表(最常用)

# 开启 Hive 支持,自动读取 Hive Metastore 中的表
spark = SparkSession.builder \
    .appName("app") \
    .enableHiveSupport() \
    .getOrCreate()

# 直接查询,无需额外配置
spark.sql("SELECT * FROM ods.user_log LIMIT 10").show()

② MySQL 表(JDBC 外部数据源)

# 方式 1:临时读取后注册视图
mysql_df = spark.read \
    .format("jdbc") \
    .option("url", "jdbc:mysql://host:3306/db") \
    .option("dbtable", "user") \
    .option("user", "root") \
    .option("password", "xxx") \
    .load()

mysql_df.createOrReplaceTempView("mysql_dim.user")

# 方式 2:创建永久外部表(映射关系存入 Hive Metastore)
spark.sql("""
CREATE TABLE IF NOT EXISTS mysql_dim.user
USING jdbc
OPTIONS (
    url 'jdbc:mysql://host:3306/db',
    dbtable 'user',
    user 'root',
    password 'xxx'
)
""")

③ 原始文件(CSV/Parquet/JSON)

df = spark.read.csv("hdfs://xxx/data/log.csv", header=True)
df.createOrReplaceTempView("file_log")

④ HBase 表

df = spark.read \
    .format("org.apache.phoenix.spark") \
    .option("table", "user_profile") \
    .option("zkUrl", "host:2181") \
    .load()
df.createOrReplaceTempView("hbase_user")

08表名前缀的真相

MySQL 开发者看到 ods.user_logmysql_dim.user 这样的表名,自然会问:这些前缀是 Spark 的标准约定吗?看到 ods. 就知道是 Hive?看到 mysql_dim. 就知道是 MySQL?

答案:不是!

Apache Spark 源码、官方文档没有强制约定表名前缀与数据源的对应关系。Spark 完全不通过库名前缀判断数据源类型。

ods. 的来源:数仓分层标准(Hive 专属)

ods 是离线数仓分层标准中的贴源原始数据层(Operational Data Store),所有从 MySQL、日志同步过来的原始数据统一存在 Hive 的 ods 库下。这是行业通用的人工命名规范,不是 Spark 内核逻辑。

分层 全称 含义
odsOperational Data Store原始贴源层
dwdData Warehouse Detail明细清洗层
dwsData Warehouse Summary汇总宽表层
dimDimension公共维度表(Hive 内部维度)
adsApplication Data Store应用指标层

mysql_dim. 是企业自定义规则

企业为了一眼区分"Hive 维度"和"MySQL 业务维度"做的命名区分:

关键误区

误区:Spark 靠小数点前的前缀自动路由数据源

事实:你完全可以把 MySQL 视图注册成 ods.test,Spark 依旧读 MySQL,不会读 Hive。前缀只是给人看的,程序逻辑完全不依赖。

09Spark 如何识别数据源

既然不靠表名前缀,那 Spark 怎么知道 ods.user_log 是 Hive 表、mysql_dim.user 是 MySQL 表?答案是:靠"这张表是在哪种元数据载体里注册的"

先分清两种表对象

永久表(托管 / 外部)

  • 元数据持久存在 Hive Metastore 或 Spark Catalog
  • 重启 Spark 后还能直接查询
  • 相当于 MySQL 中的物理表

临时视图(TempView)

  • 只存在当前 Spark 会话内存
  • 重启失效,必须代码先读数据源再注册
  • 相当于 SQL 中的 WITH 临时表

逐表识别逻辑

ods.user_log → Hive 数据源

1
SQL 解析
识别 ods.user_log
2
查 Metastore
返回存储路径和 Schema
3
Hive Connector
读取 HDFS 上的文件

前提:代码开启了 enableHiveSupport()。哪怕你把库改名成 aaa.bbb,只要元数据在 Hive Metastore,Spark 依旧走 HDFS 读取。

mysql_dim.user → MySQL 数据源

1
SQL 解析
识别临时视图
2
查会话内存
发现绑定了 JDBC
3
JDBC Connector
从 MySQL 拉取数据

hbase_user → HBase 数据源

同样是临时视图模式,Spark 内存记录该视图绑定的是 HBase 连接器,查询时走 HBase ZK 读取 HFile 数据。

统一执行流程
  1. Spark SQL 解析器拆分三张表
  2. 分别查表元数据来源:Hive Metastore / 临时视图 / 外部表定义
  3. Catalyst 优化器生成分布式计划,分别下发任务拉取数据(支持谓词下推)
  4. 所有数据拉到 Spark 集群 Executor 内存
  5. 集群内部完成 JOIN、DISTINCT、GROUP BY
  6. 汇总结果输出

Spark 3.x 正规方案:Multi Catalog

企业正式环境不靠视图和命名区分,而是配置多 Catalog,语法天然隔离:

# spark-defaults.conf 配置
spark.sql.catalog.hive_catalog  org.apache.spark.sql.hive.HiveCatalog
spark.sql.catalog.mysql_catalog com.mysql.jdbc.MySQLCatalog
spark.sql.catalog.mysql_catalog.url jdbc:mysql://host:3306

# SQL 中通过 Catalog 名天然隔离数据源
SELECT *
FROM hive_catalog.ods.user_log a
JOIN mysql_catalog.dim.user b ON a.uid = b.uid

此时 . 前面第一层是 Catalog 名称,Spark 配置文件绑定了 Catalog 与对应数据源,是程序层面的官方区分方案。

10完整实战

需求:读取用户行为日志(CSV),统计各页面访问 UV,写入结果文件。这是从 MySQL 开发者视角的端到端 Spark 实战。

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

# ===== 1. 创建 SparkSession =====
spark = SparkSession.builder \
    .appName("PageUV") \
    .master("local[*]") \
    .getOrCreate()

# ===== 2. 读取数据(类似 LOAD DATA INFILE)=====
df = spark.read.csv(
    "hdfs://xxx/data/user_log.csv",
    header=True,
    inferSchema=True
)

# ===== 3. 注册临时视图 =====
df.createOrReplaceTempView("user_log")

# ===== 4. 执行 SQL(和 MySQL 写法一致)=====
result = spark.sql("""
SELECT
    page,
    COUNT(DISTINCT uid) AS uv,
    COUNT(*) AS pv
FROM user_log
WHERE date >= '2024-01-01'
GROUP BY page
ORDER BY uv DESC
""")

# ===== 5. 查看结果 =====
result.show()

# +------+---+---+
# |  page| uv| pv|
# +------+---+---+
# |  home|150|300|
# |search|120|250|
# |  cart| 80|100|
# +------+---+---+

# ===== 6. 写入结果文件(类似 INSERT INTO ... SELECT)=====
result.write.mode("overwrite") \
       .parquet("hdfs://xxx/output/uv_result")

# ===== 7. 关闭 =====
spark.stop()
MySQL 开发者注意

如果你之前用 SQLAlchemy 做同样的分析:数据全量拉到 Python 内存 → 本地 pandas 聚合 → 千万行数据直接 OOM。Spark 的做法是:数据分散在 HDFS 多台服务器,Spark 拆分任务并行统计,几分钟出结果。

11常见误区与 FAQ

Spark 是一种语言吗?

不是。Spark 是分布式计算框架,你用 Python/PySpark SQL/Scala 编写逻辑。类似 MySQL 不是语言,SQL 才是语言。

RDD 和 DataFrame 选哪个?

生产全部用 DataFrame。自带 Catalyst 优化器,速度比原生 RDD 快数倍。RDD 只在理解底层原理时接触。

local[*] 什么时候删掉?

本地测试保留 .master("local[*]");提交 YARN 集群运行时删除,改用 spark-submit --master yarn

数据存在哪里?

底层存储 HDFS / 对象存储,Spark 只负责计算,不持久化数据。这和 MySQL"数据在自己磁盘上"完全不同。

表名前缀能判断数据源吗?

不能ods. mysql_dim. 只是行业/企业命名规范,Spark 靠元数据(Metastore / Catalog / 临时视图绑定)识别数据源。

SparkSQL 和 MySQL SQL 语法一样吗?

基本相同 标准 SELECT/JOIN/GROUP BY 语法一致。差异:Spark 不支持行级 UPDATE/DELETE(需 Hudi/Iceberg),写入是批量覆盖/追加模式。

Spark 与生态组件分工总结

组件 定位 适用场景
Spark 通用分布式计算引擎 离线 ETL 分层、实时流、复杂大表 JOIN、机器学习
Trino/Presto 交互式联邦 SQL 引擎 分析师 Ad-hoc 查询、跨 MySQL/Hive/ES 多源关联、BI 报表
Hive 数据仓库元数据管理 表结构、分区管理、HiveQL 兼容、元数据中心
MySQL 关系型数据库 业务系统增删改查、小维度查询、后台业务接口

12进阶学习路线

1. Spark 核心

分区机制、缓存 cache、广播变量、累加器

2. Spark SQL 进阶

窗口函数、分区写入、Hive 分区表操作

3. 实时计算

Structured Streaming 消费 Kafka

4. 集群部署

YARN 提交任务、参数调优

5. 数据湖

Iceberg / Hudi 与 Spark 集成

6. Spark + Notebook

Jupyter/Zeppelin 交互式分析工作流

给 MySQL 开发者的建议
  1. 不要试图用 Spark 替代 MySQL——它们是互补关系,不是替代关系。业务系统用 MySQL,海量离线分析用 Spark。
  2. 从 SparkSQL 开始——你的 SQL 技能直接复用,先学会 spark.sql() 跑通分析任务,再深入 DataFrame DSL。
  3. 重点理解执行模型差异——不是语法难学,而是"谁来算"这个思维转变最重要。
  4. 配套分工:业务后台用 SQLAlchemy + MySQL;海量离线加工用 Spark SQL;分析师自助查询用 Trino/Presto;元数据管理用 Hive。