作为 MySQL 开发者,你已经熟悉了关系型数据库的建表、写 SQL、用 ORM 框架操作数据的完整流程。但当数据量从百万级增长到亿级、十亿级时,MySQL 的单机架构会遭遇性能瓶颈——慢查询、锁表、CPU/IO 打满。这就是你需要学习 Spark 的根本原因。
MySQL 开发者的习惯:写 SQL → 发给数据库 → 数据库算好返回结果
Spark 开发者的思维:写 SQL → Spark 自己解析优化 → 拆分任务分发到集群多节点并行计算 → 汇总结果
差距不在"怎么写 SQL"(语法几乎一样),而在"谁来执行这条 SQL、在哪台机器上算、数据存在哪里"。
Spark 是通用分布式大数据计算引擎,支持离线批处理、实时流、SQL 分析、机器学习,替代老旧 MapReduce,内存计算速度提升 10~100 倍。它本身不持久化存储任何数据——数据存在 HDFS / 对象存储 / Hive / MySQL / HBase 等外部系统中,Spark 只负责把数据读出来、在集群内存中并行计算、再把结果写回去。
下面这张表是你最需要的——把已有的 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 集群模式 |
MySQL 开发者习惯用 mysql CLI 或 Navicat 连数据库。Spark 的环境稍有不同——你需要一个 Spark 运行环境来执行代码。以下是三种方案,按推荐程度排序。
最简单的方式,无需下载 Spark 安装包,一条命令搞定。
# 安装 PySpark
pip install pyspark==3.5.0
# 启动交互式终端(类似 mysql CLI)
pyspark
出现 Welcome to PySpark 3.5.0 即成功。启动后自动生成两个全局变量:
spark:SparkSession(2.0+ 统一入口,相当于你的数据库连接)sc:SparkContext(底层 RDD 上下文)# 下载解压
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
D:\sparkwinutils.exe,配置 HADOOP_HOMESPARK_HOME,Path 追加 %SPARK_HOME%\binpyspark# Jupyter Notebook 或 Google Colab 中直接执行
pip install pyspark
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("demo").getOrCreate()
下面用同一个需求——"统计用户行为日志中各页面的 UV"——对比 MySQL 写法和 Spark 写法,让你直观感受差异。
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 只接收结果。
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 优化器解析执行。
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
这是从 MySQL 转向 Spark 最核心的知识点——理解 SQL 在两种环境下"谁来执行"的根本差异。
MySQL 路径:全部计算压在数据库单机
Spark 路径:计算收拢到 Spark 集群,存储只提供原始数据
| 对比维度 | SQLAlchemy(普通数据库 SQL) | Spark SQL |
|---|---|---|
| 计算载体 | 数据库服务端(单节点) | Spark 集群多机器分布式并行 |
| SQL 作用域 | 仅能查询单一关系库 | 支持联邦查询:Hive + MySQL + ES + Kafka |
| 优化器 | 数据库自带优化器,受单机限制 | Catalyst 分布式优化器,谓词下推、分区裁剪 |
| 数据处理上限 | 千万级性能急剧衰减 | 十亿级常规处理 |
| 执行模式 | 单进程串行,无分布式分片 | 自动分片、多 Executor 并行 |
| 内存机制 | 数据库缓存,Python 客户端不参与 | 中间结果缓存到集群内存 cache() |
| 适用场景 | 业务系统增删改查、小维度查询 | 数仓 ETL、海量日志分析、跨源统计 |
| 事务支持 | 完整 ACID 事务 | 有限事务,面向批量 |
很多人误以为 spark.read.jdbc() 等价于 SQLAlchemy,实际完全不同:
简言之:Spark JDBC 是"把 MySQL 数据搬到 Spark 集群算",SQLAlchemy 是"让 MySQL 自己算"。
对 MySQL 开发者来说,DataFrame 就是一张分布式数据表——带字段名、数据类型,兼容标准 SQL,性能最优。这是生产环境的核心,必学。
带 Schema 的结构化表,支持 SQL 和 Catalyst 自动优化,性能远高于 RDD。相当于 MySQL 中的表。
无结构的原始分布式集合,底层 API。复杂底层开发才用,日常不需要。理解原理即可。
HDFS / 对象存储、Hive 元数据。相当于 MySQL 的磁盘存储,但 Spark 不拥有数据。
# 类似 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|
# +-----+----+----------+------+
# 类似 SQLAlchemy ORM 的链式查询
df.filter(df.city == "北京") \
.groupBy("page") \
.count() \
.show()
# +----+-----+
# |page|count|
# +----+-----+
# |home| 1|
# |cart| 1|
# +----+-----+
将 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()
# 读 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(已存在则报错) |
这是 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 中这不可能——你没法在一条 SQL 里 JOIN 另一个数据库系统的表。通常做法是分别查两张表,在 Python 本地内存 JOIN,数据量大直接 OOM。
# 开启 Hive 支持,自动读取 Hive Metastore 中的表
spark = SparkSession.builder \
.appName("app") \
.enableHiveSupport() \
.getOrCreate()
# 直接查询,无需额外配置
spark.sql("SELECT * FROM ods.user_log LIMIT 10").show()
# 方式 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'
)
""")
df = spark.read.csv("hdfs://xxx/data/log.csv", header=True)
df.createOrReplaceTempView("file_log")
df = spark.read \
.format("org.apache.phoenix.spark") \
.option("table", "user_profile") \
.option("zkUrl", "host:2181") \
.load()
df.createOrReplaceTempView("hbase_user")
MySQL 开发者看到 ods.user_log、mysql_dim.user 这样的表名,自然会问:这些前缀是 Spark 的标准约定吗?看到 ods. 就知道是 Hive?看到 mysql_dim. 就知道是 MySQL?
Apache Spark 源码、官方文档没有强制约定表名前缀与数据源的对应关系。Spark 完全不通过库名前缀判断数据源类型。
ods 是离线数仓分层标准中的贴源原始数据层(Operational Data Store),所有从 MySQL、日志同步过来的原始数据统一存在 Hive 的 ods 库下。这是行业通用的人工命名规范,不是 Spark 内核逻辑。
| 分层 | 全称 | 含义 |
|---|---|---|
ods | Operational Data Store | 原始贴源层 |
dwd | Data Warehouse Detail | 明细清洗层 |
dws | Data Warehouse Summary | 汇总宽表层 |
dim | Dimension | 公共维度表(Hive 内部维度) |
ads | Application Data Store | 应用指标层 |
企业为了一眼区分"Hive 维度"和"MySQL 业务维度"做的命名区分:
dim.xxx:Hive 内部存储的静态维度表(HDFS)mysql_dim.xxx:映射 MySQL 库的维度表hbase_xxx:映射 HBase 表file_temp:文件注册的临时视图误区:Spark 靠小数点前的前缀自动路由数据源
事实:你完全可以把 MySQL 视图注册成 ods.test,Spark 依旧读 MySQL,不会读 Hive。前缀只是给人看的,程序逻辑完全不依赖。
既然不靠表名前缀,那 Spark 怎么知道 ods.user_log 是 Hive 表、mysql_dim.user 是 MySQL 表?答案是:靠"这张表是在哪种元数据载体里注册的"。
前提:代码开启了 enableHiveSupport()。哪怕你把库改名成 aaa.bbb,只要元数据在 Hive Metastore,Spark 依旧走 HDFS 读取。
同样是临时视图模式,Spark 内存记录该视图绑定的是 HBase 连接器,查询时走 HBase ZK 读取 HFile 数据。
企业正式环境不靠视图和命名区分,而是配置多 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 与对应数据源,是程序层面的官方区分方案。
需求:读取用户行为日志(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()
如果你之前用 SQLAlchemy 做同样的分析:数据全量拉到 Python 内存 → 本地 pandas 聚合 → 千万行数据直接 OOM。Spark 的做法是:数据分散在 HDFS 多台服务器,Spark 拆分任务并行统计,几分钟出结果。
不是。Spark 是分布式计算框架,你用 Python/PySpark SQL/Scala 编写逻辑。类似 MySQL 不是语言,SQL 才是语言。
生产全部用 DataFrame。自带 Catalyst 优化器,速度比原生 RDD 快数倍。RDD 只在理解底层原理时接触。
本地测试保留 .master("local[*]");提交 YARN 集群运行时删除,改用 spark-submit --master yarn。
底层存储 HDFS / 对象存储,Spark 只负责计算,不持久化数据。这和 MySQL"数据在自己磁盘上"完全不同。
不能。ods. mysql_dim. 只是行业/企业命名规范,Spark 靠元数据(Metastore / Catalog / 临时视图绑定)识别数据源。
基本相同 标准 SELECT/JOIN/GROUP BY 语法一致。差异:Spark 不支持行级 UPDATE/DELETE(需 Hudi/Iceberg),写入是批量覆盖/追加模式。
| 组件 | 定位 | 适用场景 |
|---|---|---|
| Spark | 通用分布式计算引擎 | 离线 ETL 分层、实时流、复杂大表 JOIN、机器学习 |
| Trino/Presto | 交互式联邦 SQL 引擎 | 分析师 Ad-hoc 查询、跨 MySQL/Hive/ES 多源关联、BI 报表 |
| Hive | 数据仓库元数据管理 | 表结构、分区管理、HiveQL 兼容、元数据中心 |
| MySQL | 关系型数据库 | 业务系统增删改查、小维度查询、后台业务接口 |
分区机制、缓存 cache、广播变量、累加器
窗口函数、分区写入、Hive 分区表操作
Structured Streaming 消费 Kafka
YARN 提交任务、参数调优
Iceberg / Hudi 与 Spark 集成
Jupyter/Zeppelin 交互式分析工作流
spark.sql() 跑通分析任务,再深入 DataFrame DSL。