7.2 与NoSQL数据库集成:把扫描并行化


7.2 与NoSQL数据库集成:把扫描并行化

本节摘要:Spark 读 NoSQL 行存(HBase、Cassandra、MongoDB 等)时,连接器把表按分区键切成 scan 区间分派给任务。本节以 HBase 与 Cassandra 为样本,看切片从哪来、并行度怎么定、写回时的权衡,并归纳连接器的三种读模式。

上一节的数据躺在文件系统里,块边界就是天然的分区边界。NoSQL 没有文件概念,但有自己的切分逻辑——HBase 按 region、Cassandra 按 token 环。连接器的全部工作,就是把数据库的切分翻译成 Spark 的分区。

HBase:region 就是分区清单

# 读:以 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 再多核也只能一个任务慢慢扫。

Cassandra:token 环切成区间

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 一致性语义——写侧可声明一致性级别,写到一半任务重试时,不幂等的数据会重复,按业务选择去重或幂等键设计。

巡检案例:一场把 Cassandra 打挂的分析

背景:分析师用 Spark 扫一张亿行 Cassandra 表做年报,未设切片参数,也未做列裁剪。操作:直接全表 load 后再过滤。结果:Cassandra 协调节点 CPU 打满、线上查询超时,Spark 侧只有 8 个任务在慢速爬行。解读:三个错叠加——无切片并行度低、无下推全字段拉取、拉回再滤把网络放大十几倍,数据库成了唯一的窄口。变式:select 两列加分区键过滤、切片调小后,同样数据 10 分钟跑完,数据库负载峰值降七成。结论:NoSQL 集成的第一守则是让数据库做它擅长的按键读取,Spark 只负责它做不了的聚合分析。

⚠️ 常见坑:把 Spark 当作 NoSQL 的在线查询代理。连接器面向批量扫描设计,毫秒级点查应该走数据库自己的客户端,走 Spark 要付任务调度与序列化的固定税。

💡 关键直觉:NoSQL 连接器的体检指标就两个——切片数够不够并行、谓词有没有下推。两者都在执行计划里一行可见。

本节要点回顾

  • 数据库分区即 Spark 分区:HBase 被 region 数锁死,Cassandra 可调区间
  • 下推看 PushedFilters:分区键过滤与列裁剪必须发生在数据库侧
  • JDBC 要显式切:分区列加上下界,否则单连接一把梭
  • 写回靠批量:批次与并发要匹配数据库容量,重试需幂等设计
  • 别拿 Spark 做点查:固定调度税决定了它只属于批量分析

下一节的数据不再安静地躺着,而是源源不断地流动——消息队列的集成要处理"读取位置"这个新变量。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U