本节摘要:Spark 读 NoSQL 行存(HBase、Cassandra、MongoDB 等)时,连接器把表按分区键切成 scan 区间分派给任务。本节以 HBase 与 Cassandra 为样本,看切片从哪来、并行度怎么定、写回时的权衡,并归纳连接器的三种读模式。
上一节的数据躺在文件系统里,块边界就是天然的分区边界。NoSQL 没有文件概念,但有自己的切分逻辑——HBase 按 region、Cassandra 按 token 环。连接器的全部工作,就是把数据库的切分翻译成 Spark 的分区。
# 读:以 OpenTSDB 风格的时序表为例 scan = {"table": "tsdb", "catalog": """ {"table":{"namespace":"default","name":"tsdb"}, "rowkey":"key", "columns":{ "metric":{"cf":"t","col":"metric","type":"string"}, "value":{"cf":"t","col":"v","type":"double"} }} """} df = spark.read.format("org.apache.hadoop.hbase") \ .options(**scan).load() df.filter(df.metric == "cpu.load").count()
连接器先向 HBase 的 Master 查询表的 region 清单——每个 region 负责一段连续 rowkey 并落在某台 region 服务器上。这份清单直接成为 Spark 的分区计划:一个 region 一个任务,任务尽量派到 region 所在节点(数据本地性在行存世界以另一种形式复活)。巡检推论:region 数就是并行度上限,一张没预分裂的小表只有一个 region,Spark 再多核也只能一个任务慢慢扫。
df = spark.read.format("org.apache.spark.sql.cassandra") \ .options(table="events", keyspace="analytics").load() # 单列主键表按 token 区间切,并发度由两个参数控制 spark.conf.set("spark.cassandra.input.split.size", "16000") df.where(df.tenant_id == 7).select("event", "ts").collect()
Cassandra 把键空间哈希成一个 token 环,连接器按 token 区间切 scan,区间大小由输入切片参数控制——并行度是可调的,不像 HBase 那样被 region 数锁死。更要紧的是谓词下推:where 里的分区键条件会被改写成 Cassandra 的范围查询,select 的列裁剪也一并下推。看执行计划的 PushedFilters 行,能确认过滤真的发生在数据库侧而不是拉回来再滤。
| 模式 | 切片来源 | 下推能力 | 典型代表 |
|---|---|---|---|
| 元数据切分型 | 数据库分区结构 | 分区键与行键可下推 | HBase、Cassandra |
| 抽样切分型 | 驱动列的值域抽样 | 只能整段扫描 | JDBC 长表 |
| 单连接型 | 无切片 | 无 | 通用 REST 源 |
JDBC 的例子最能说明第三类问题:不设分区参数时,连接器开一条连接一把梭,全部数据流过一个任务——Spark 沦为昂贵的单线程客户端。解法是指定分区列、上下界与分片数,把长表切成多段并行拉取。
df = spark.read.format("jdbc") \ .option("url", "jdbc:mysql://db:3306/dw") \ .option("dbtable", "orders") \ .option("partitionColumn", "order_id") \ .option("lowerBound", "1").option("upperBound", "100000000") \ .option("numPartitions", "64").load()
df.write.format("org.apache.spark.sql.cassandra") \ .option("keyspace", "analytics").option("table", "events_summary") \ .mode("append").save()
写路径的关键词是批量:连接器把每个分区的行攒成批次再提交,逐行写会把 NoSQL 的协调节点打垮。两个巡检点:其一是并发控制——写任务数乘批次大小要与数据库的写入容量匹配,参数可调;其二是 Cassandra 一致性语义——写侧可声明一致性级别,写到一半任务重试时,不幂等的数据会重复,按业务选择去重或幂等键设计。
背景:分析师用 Spark 扫一张亿行 Cassandra 表做年报,未设切片参数,也未做列裁剪。操作:直接全表 load 后再过滤。结果:Cassandra 协调节点 CPU 打满、线上查询超时,Spark 侧只有 8 个任务在慢速爬行。解读:三个错叠加——无切片并行度低、无下推全字段拉取、拉回再滤把网络放大十几倍,数据库成了唯一的窄口。变式:select 两列加分区键过滤、切片调小后,同样数据 10 分钟跑完,数据库负载峰值降七成。结论:NoSQL 集成的第一守则是让数据库做它擅长的按键读取,Spark 只负责它做不了的聚合分析。
⚠️ 常见坑:把 Spark 当作 NoSQL 的在线查询代理。连接器面向批量扫描设计,毫秒级点查应该走数据库自己的客户端,走 Spark 要付任务调度与序列化的固定税。
💡 关键直觉:NoSQL 连接器的体检指标就两个——切片数够不够并行、谓词有没有下推。两者都在执行计划里一行可见。
下一节的数据不再安静地躺着,而是源源不断地流动——消息队列的集成要处理"读取位置"这个新变量。