Spark 读写外部数据源
概述
Spark 通过 DataSource API 统一读写外部系统。数据源选型与参数配置直接决定 ETL 性能:JDBC 要控分区避免压垮数据库,Kafka 要调吞吐,HBase/Redis 要防热点。本文逐个讲解六类常用数据源的读写方式与调优要点。
一、数据源抽象
1.1 统一接口
scala
spark.read.format("jdbc")
.option("url", ...).option("dbtable", "orders").load()
df.write.format("parquet").save("hdfs:///out")| 方法 | 说明 |
|---|---|
format | 数据源类型 |
option | 数据源参数 |
load/save | 读写入口 |
jdbc/parquet/... | 便捷方法 |
1.2 分区与并行度
| 概念 | 说明 |
|---|---|
| 分区数 | 读数据源的 Task 数 |
| 分区策略 | 按列范围 / 按分区键 |
| 数据本地性 | 尽量就近读 |
二、JDBC(MySQL/PostgreSQL/Oracle)
2.1 读
scala
spark.read.jdbc(
url = "jdbc:mysql://host:3306/db",
table = "orders",
column = "id", // 分区列
lowerBound = 1, upperBound = 10000000,
numPartitions = 10,
connectionProperties = props
)2.2 分区调优
| 参数 | 说明 |
|---|---|
partitionColumn | 分区列(需有序,如主键) |
lowerBound/upperBound | 分区范围 |
numPartitions | 分区数(并行度) |
fetchsize | 每次拉取行数 |
queryTimeout | 查询超时 |
| 风险 | 处理 |
|---|---|
| 并发太高压垮数据库 | numPartitions 控制并发度(≤10-20) |
| 分区列分布不均 | 选均匀列或自定义查询 |
| 大表全量拉取 | 按时间增量同步 |
2.3 写
scala
df.write.mode("append").option("batchsize", "1000").jdbc(url, "orders", props)| 注意 | 说明 |
|---|---|
| batchsize | 批量提交(默认 1000) |
| 主键冲突 | 用 upsert 或先删后插 |
| 写压力 | 控制写入并行度 |
三、Hive
3.1 读写
scala
spark.read.table("dw.ods_orders") // 读 Hive 表
df.write.saveAsTable("dw.dws_report") // 写 Hive 表| 说明 | 值 |
|---|---|
| 分区表 | 自动识别分区列 |
| 动态分区 | spark.sql.sources.partitionOverwriteMode 控制覆盖 |
| 格式 | 默认 Parquet |
3.2 调优
| 参数 | 作用 |
|---|---|
spark.sql.hive.convertMetastoreParquet | 用原生 Parquet 读 |
| 向量化读取 | 大表读加速 |
| 分区裁剪 | 只读必要分区 |
四、HBase
4.1 读
scala
spark.read.format("org.apache.hadoop.hbase.spark")
.option("hbase.table", "user")
.option("hbase.columns.mapping", "id f:name, f:age")
.load()4.2 要点
| 要点 | 说明 |
|---|---|
| RowKey 扫描 | 读时用 RowKey 范围(startRow/stopRow) |
| 列族映射 | 指定 family:qualifier |
| 缓存 | HBase 端 blockcache 调优 |
| 避免全表扫描 | 大表全扫极慢,用过滤器 |
4.3 调优
| 手段 | 说明 |
|---|---|
| 增大 scan 缓存 | setCaching 减少 RPC |
| 并行 Region 扫描 | 分区数与 Region 数匹配 |
| 批量写 | 用 HFile 批量生成或 Put 批量 |
五、Kafka
5.1 读(批量)
scala
spark.read.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "topic")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
.selectExpr("CAST(value AS STRING)")5.2 吞吐调优
| 参数 | 说明 |
|---|---|
minPartitions | 最小分区数(≥ Topic 分区数) |
maxOffsetsPerTrigger | 每次读取上限(批量场景) |
fetchOffset.numRetries | 拉 offset 重试 |
kafka.batch.size | 批量拉取大小 |
| 场景 | 配置方向 |
|---|---|
| 一次性全量读 | minPartitions 拉满,maxOffsetsPerTrigger 置大 |
| 流式增量 | Structured Streaming 用触发器控制 |
六、Elasticsearch
6.1 读写
scala
spark.read.format("org.elasticsearch.spark.sql")
.option("es.nodes", "es:9200")
.option("es.resource", "index/doc")
.option("es.query", """{"match_all":{}}""")
.load()
df.write.format("org.elasticsearch.spark.sql")
.option("es.resource", "index")
.option("es.mapping.id", "id")
.mode("append").save()6.2 调优
| 参数 | 作用 |
|---|---|
es.batch.size.entries | 批量写条目数 |
es.batch.size.bytes | 批量写字节数 |
es.index.read.missing.as.empty | 缺索引容错 |
es.scroll.size | 滚动读大小 |
es.nodes.wan.only | 跨网络访问 |
| 注意 | 说明 |
|---|---|
| 大批量写压垮 ES | 控制分区与批量大小 |
| 映射冲突 | 写前保证 mapping 兼容 |
| 深分页 | 用 scroll 而非深分页查询 |
七、Redis
7.1 读
scala
spark.read.format("org.apache.spark.sql.redis")
.option("host", "redis-host")
.option("port", "6379")
.option("table", "user_cache")
.load()7.2 要点
| 要点 | 说明 |
|---|---|
| key 设计 | 按前缀/模式读取 |
| 热点 key | 读压集中导致 Redis 热点 |
| pipeline | 批量读用 pipeline |
| 写入模式 | 批量 set/批量写 hash |
| 风险 | 处理 |
|---|---|
| 大 key | 拆分或换存储 |
| 热点 | 加缓存分层、随机后缀 |
| 连接数爆炸 | 控制并行度,用连接池 |
八、数据源选型对照
| 数据源 | 典型用途 | 读优化 | 写优化 |
|---|---|---|---|
| JDBC | 业务库同步 | 分区列+增量 | 批量提交 |
| Hive | 数仓主存储 | 分区裁剪+向量化 | saveAsTable |
| HBase | 实时 KV | RowKey 扫描 | 批量 Put |
| Kafka | 消息源 | 分区对齐 | 批量+ack |
| ES | 检索 | scroll | 批量+id |
| Redis | 缓存/维表 | pipeline | 批量 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| JDBC 拉取压垮库 | 减少分区与 fetchsize,走增量 |
| Kafka 读太快积压 | 用 maxOffsetsPerTrigger 限速 |
| ES 写超时 | 降批量与并行,调 es 集群 |
| HBase 全表扫慢 | 加 RowKey 范围过滤 |
| Redis 连接过多 | 控制并行度、复用连接 |