☰
Cassandra与Spark集成实战:构建大数据ETL流水线全解析
2026/10/3 14:28:11 网站建设 项目流程

Cassandra 和 Spark 这对组合,在我这几年的数据工程实战里,几乎是绕不开的搭档。一个负责海量数据的高可用写入与存储,一个负责把沉睡在宽表里的数据拉出来做复杂计算和清洗。很多刚接触大数据处理流水线的朋友会问:Spark 不是有自己的数据源嘛,HDFS、Hive 都挺熟,为什么偏偏要接 Cassandra?

这个问题的答案其实很务实:Cassandra 天然适合写多读少的场景,但 CQL 对复杂聚合分析并不友好,Spark 恰好擅长分布式计算和批量 ETL。把两者集成起来,等于给数据仓库加了一个高性能前端计算引擎,既不用把数据搬来搬去,又能在原库上直接做清洗、聚合、建模。这篇文章就从一个真实项目的视角,带你把 Cassandra + Spark 集成流水线从零搭起来,讲清楚架构设计、连接器配置、读写代码怎么写、参数怎么调,以及那些官方文档里不会明说的坑。

1. 集成场景与整体设计思路

1.1 为什么要用 Spark 读 Cassandra

这里先说一个我反复跟团队强调的观点:不要为了集成而集成。如果你只是偶尔查几条记录,Cassandra 自己的 CQL 已经足够快,直接写单行查询就好,没必要上 Spark。真正需要 Spark 介入的,是下面这几类场景。

第一类是批量 ETL 与数据治理。比如日志系统往 Cassandra 里灌了上百张表、每天几十亿条数据,你要做字段裁剪、格式统一、敏感信息脱敏,还要把结果导给下游数据仓库。这种全表或大范围扫描操作,用 CQL 一条条查会慢到无法忍受,但用 Spark 的分布式扫描能力,可以把整个集群的 CPU 全部利用起来,一把梭完成全量计算。

第二类是复杂分析计算。Cassandra 的索引机制是为点查设计的,二级索引在跨分区聚合场景下几乎废掉,更别说 GROUP BY、JOIN、窗口函数。Spark 天然支持这些算子,通过连接器把数据拉进 DataFrame / RDD,就等于把 Cassandra 当成了一个分布式数据源,算完之后再写回或者落盘。

第三类是实时数仓的贴源层加工。当前很流行的 Lambda 架构中,Cassandra 往往承担“speed layer”和“serving layer”的角色,Spark Streaming 从 Kafka 拿到实时数据,经过结构化处理之后落到 Cassandra,后续的批处理再用 Spark 周期性读取 Cassandra 做离线重算。这种“读原库—加工—回写”的闭环,是 Cassandra + Spark 集成最典型的使用姿势。

1.2 流水线的两种常见架构选型

集成架构上,我见过的主流方案有两条路。

方案一:直连型。客户端直接让 Spark 通过连接器读写 Cassandra。数据不需要中转,计算节点并行扫表。优点是架构简单、运维成本低,缺点是会把 Cassandra 的 CPU 和 IO 压力拉起来,高峰期可能互相抢资源。

方案二:数据导出型。先用 Spark 把 Cassandra 数据批量导出成 Parquet / ORC 文件放到 HDFS 或 S3,再基于文件构建数仓。优点是计算与存储解耦,Cassandra 的压力可控,缺点是多了一份数据冗余,实时性变差。

从实践上看,如果 Cassandra 集群本身是独立资源池、CPU 不紧张,我倾向方案一直连;如果业务量非常大、集群不能承受全表扫描,就做方案二。你可以在同一条流水线里按表维度混用——大表导出小表直连,效果很不错。

2. 环境准备与连接器核心配置

2.1 连接器版本与依赖引入

Cassandra 和 Spark 之间最常用的桥是 DataStax 开源的spark-cassandra-connector。这东西的版本兼容性需要特别注意,很多初学者的集成失败都是栽在版本号上。

我的建议是:先确认 Spark 版本,再去连接器 GitHub 的 Release 页面查对应版本。以我常用的组合为例:

Spark 版本Scala 版本推荐连接器版本
3.0.x / 3.1.x2.123.0.1
3.2.x2.123.2.0
3.3.x2.123.2.0 以上
2.4.x2.112.4.3

注意连接器的 groupId 是com.datastax.spark,artifactId 是spark-cassandra-connector-assembly,也就是带所有依赖的打包版本。如果你用 Maven,可以这样引入:

<dependency> <groupId>com.datastax.spark</groupId> <artifactId>spark-cassandra-connector-assembly_2.12</artifactId> <version>3.2.0</version> </dependency>

这里有个容易踩的坑:连接器内部带的guava、netty等依赖经常和 Spark 自带的版本冲突。表现为启动时一堆 NoSuchMethodError 或 ClassNotFound,非常迷惑。强烈建议在提交任务时使用--conf spark.driver.userClassPathFirst=true和--conf spark.executor.userClassPathFirst=true,让用户 jar 优先加载,能解决大半冲突。碰到 guava 冲突时,我的土办法是在提交脚本里把spark.jars的路径排到前面,同时剔除冲突版本。

2.2 连接参数配置详解

连接器核心配置项不多,但每个都值得深究。我最常用的一套配置如下:

spark.cassandra.connection.host 172.16.1.10,172.16.1.11 # 至少要写两个种子节点,别只写一个 spark.cassandra.connection.port 9042 spark.cassandra.auth.username cassandra_user spark.cassandra.auth.password your_password spark.cassandra.connection.localDC DC1 # 同数据中心亲和,重要 spark.cassandra.connection.timeoutMS 10000

这里我要重点强调localDC。如果你的 Cassandra 是多数据中心部署,不设置这个参数,连接器会去随机挑一个 DC 连接,跨机房延迟直接拉垮整个计算任务。我亲眼见过一个线上任务慢 6 倍,最后发现就是没指定 localDC,连接器大部分请求打到了 200 公里外的节点上。

另外,spark.cassandra.connection.host建议写跟你 Spark 同机的 Cassandra 种子节点。连接器会先从种子节点拿拓扑元数据,然后每台 executor 各自直连数据节点。种子节点只起“引路”作用,不是数据访问中转站。

如果要提交到集群,建议把这些配置写在 Spark 提交脚本里,而不是硬编码到 Java / Scala 代码中,这样环境切换时不用改代码:

spark-submit \ --class com.example.CassandraEtlJob \ --master yarn \ --deploy-mode cluster \ --conf spark.cassandra.connection.host=172.16.1.10 \ --conf spark.cassandra.connection.port=9042 \ --conf spark.cassandra.auth.username=... \ --conf spark.cassandra.auth.password=... \ --jars spark-cassandra-connector-assembly_2.12-3.2.0.jar \ my-etl-job.jar

如果你发现连接后执行任务特别慢,先不要怀疑代码,先查连接器的executor 数量是否跟 Spark 并行度匹配。连接器默认每个 executor 会发起多个并发请求,但如果 executor 数太少,并行度再高也发挥不出来。

3. 实操:读写流水线的核心代码

3.1 从 Cassandra 读取数据并做清洗

我用 Scala 写 Spark 居多,Python 也能跑,但 Scala 在类型安全和调试上更顺手。下面这段是从一张用于埋点日志的宽表event_log里读取数据,并过滤掉非法字段、剔除近 7 天内重复记录的典型 ETL 代码,连接器通过隐式转换把 Cassandra 表变成 DataFrame 来用。

import org.apache.spark.sql.SparkSession import com.datastax.spark.connector._ import org.apache.spark.sql.functions._ val spark = SparkSession.builder() .appName("cassandra-spark-etl") .config("spark.cassandra.connection.host", "172.16.1.10,172.16.1.11") .config("spark.cassandra.connection.localDC", "DC1") .getOrCreate() val rawDf = spark.read .format("org.apache.spark.sql.cassandra") .options(Map( "table" -> "event_log", "keyspace" -> "analytics", "pushdown" -> "true" )) .load() .filter(col("ts").isNotNull) .filter(col("user_id").isNotNull) .filter(length(col("device_id")) > 0) .dropDuplicates("user_id", "device_id", "event_type", "ts")

这里有个关键参数pushdown,我一般都会显式设为true。它决定了连接器会不会把 filter、where 条件下的部分谓词下推到 Cassandra 层去执行,让 Cassandra 在返回数据前就帮你过滤掉一部分行。能省很多网络开销和序列化开销。

但注意,并不是所有过滤条件都能下推。只有闭合在单个分区键上的等值条件、主键范围条件、聚簇列条件这类 Cassandra 本身能高效索引的谓词才会被下推。像dropDuplicates这种跨分区的操作,肯定要到 Spark 内存里做。对它的理解不到位,你会误以为“只要写了 filter 就万事大吉”。

3.2 计算结果写回 Cassandra

ETL 清洗完,下一步往往是把结果写回 Cassandra,或者是把聚合后的报表写进新表。写入的核心写法如下:

val resultDf = analyticsDf .groupBy("user_id", "event_type", "day") .agg( count("*").as("event_count"), sum("value").as("total_value") ) resultDf.write .format("org.apache.spark.sql.cassandra") .mode("append") .options(Map( "table" -> "event_agg_daily", "keyspace" -> "analytics", "batch.size.rows" -> "500", "batch.size.bytes" -> "102400" )) .save()

写入的时候有几点必须提前确认。

第一,目标表的主键设计。Cassandra 的表必须提前建好,Spark 只是写入方,不会帮你建表。主键(partition key)设计直接决定写入是否能并发。比如event_agg_daily的 partition key 是user_id,写入时每个 Spark partition 里的数据可以根据user_id分布到不同的 Cassandra 节点,并行度很高。但如果你的表主键是day这种粒度很粗的字段,所有数据都堆到少量节点上,写入会严重热点化,跑起来像蜗牛。

第二,append 和 overwrite 的语义。Cassandra 没有传统数据库的覆盖写,overwrite模式在连接器里通常是先删除表再写入,这个操作很危险。我建议日常流水线一律用append,需要重跑历史时先进数据层的truncate,清理后再执行 append,避免生产环境误删数据。

第三个是batch.size.rows和batch.size.bytes,分别控制每个批量请求的行数和字节数。这两个参数看着不起眼,但直接影响写入吞吐。默认值偏保守,在压力测试后可以适当调大,比如把 rows 从默认的 1000 提升到 5000,bytes 提升到 128KB,吞吐能提升 30% 以上。但是别盲目调,批量太大对 Cassandra 协调节点压力很大,而且单批失败重试的成本也高。

3.3 增量与全量处理的代码组织

真实的业务流水线很少只跑一次全量,基本都是“全量初始化 + 每天增量”。我总结了一套固定套路。

全量初始化就用上面 3.1 的写法,不加额外过滤,直接全表扫。日常增量则依赖一个last_processed_offset表,记录每张表的处理水位。增量读取代码的关键点在于:把一个带 where 条件的 filter 压到 Cassandra 查询层,比如主键或者聚类列包含时间字段,可以直接用 CQL 下推过滤,让 Cassandra 只返回增量期间的数据,大幅缩短任务时间。

val incrementalDf = spark.read .format("org.apache.spark.sql.cassandra") .options(Map("table" -> "event_log", "keyspace" -> "analytics")) .load() .filter(col("day") >= lit(lastOffsetDay))

注意day必须是表聚类列,否则这种范围过滤 Cassandra 没法高效完成,连接器会把整表拉下来再在 Spark 里过滤,增量效果全废。所以建表时就要为后续增量场景设计好时间维度字段作为聚类列,这是我反复强调的一条落库设计经验。

水位表的更新跟处理逻辑要放在一个事务控制里,最稳妥的做法是先完成计算和回写,再更新水位。万一任务失败,水位不前进,下次重跑只重复处理失败的那一段,不做额外补偿。

4. 参数调优与性能实测

4.1 读取并行度与分区控制

很多人觉得 Spark 读 Cassandra “很慢”,其实大多数时候是并行度没调对。连接器默认的读取机制是:根据 Cassandra 每个 token range 划分 Spark partition,原则上分区数等于 Cassandra 节点数乘以 vnode 数。默认情况下每个分区的数据量可能差异很大,导致数据倾斜。

我实测过一个 12 节点集群,默认读取只用了 48 个 partition,每个分区拉的记录数有些不均衡。通过如下配置可以控制读取端的并行度:

spark.cassandra.input.split.sizeInMB 64 spark.cassandra.input.split.task.size 4096

split.sizeInMB是控制每个 Spark 分区大概承载多少 MB 的数据,调小可以让分区更细、并行度更高,但也不是越小越好。分区太多,Spark 调度开销和 Cassandra 协调请求都会翻倍。从我的经验来看,单分区控制在 4 到 8 个 Cassandra 分区键数据量比较合适,具体值要靠测试来定。64 MB 是一个很稳的起点。

如果发现 executor 的资源利用率不均衡,还可以配合repartition二次处理。比如读完 DataFrame 后执行.repartition(executorCores * 2),再做下游处理。这个方法通用有效,尤其是在 join 前的整理阶段。

4.2 写入吞吐与批量配置

写入调优比读取更微妙,因为 Cassandra 的写入性能跟批量大小、一致性级别、并发数高度耦合。

我维护的一条 3000 万条记录回写任务,最初用默认配置跑了 40 分钟,调优后压到了 11 分钟。改动只做了三件事:

  • spark.cassandra.output.concurrent.writes:默认是跟 executor 核数相关,我显式设成每个 executor 4 并发。
  • spark.cassandra.output.batch.grouping.key:默认是replica_set,意思是尽量把同副本的数据拼到一个 batch。如果业务上没有强事务需求,我建议改成none,这样连接器会按主键哈希均匀分布请求,吞吐更高。
  • batch.size.rows和batch.size.bytes:上面提过,按实际压力测试逐步上调。

还有一致性级别。写入用LOCAL_QUORUM是很多团队的标准配置,兼顾一致性和性能。如果你追求极致吞吐,可以在业务允许降级的情况下用LOCAL_ONE,但这时候读要容忍短时间的最终一致性。对于流水线这类场景,我通常保留LOCAL_QUORUM,因为下游报表对准确率的要求远高于那一点吞吐收益。

4.3 内存与 GC 调优经验

读 Cassandra 的任务,executor 内存主要耗费在两块:一块是连接器拉回来的 RDD / DataFrame 数据,一块是 Spark 做 shuffle 和聚合的缓冲。组合起来,内存参数建议这样设:

--executor-memory 8g --conf spark.executor.memoryOverhead=2g --conf spark.memory.storageFraction=0.2 --conf spark.sql.shuffle.partitions=400

storageFraction默认是 0.5,意思是 storage 内存和 shuffle 执行内存各占一半。像这种读源表做计算再回写的流水线,shuffle 占比高,storage 占比低,把 storage 调到 0.2,能显著减少 GC 压力。不过这个值也别压太狠,如果缓存用的多或者 join 后要复用 DataFrame,还是要留出空间。

GC 方面,Cassandra 连接器产生的对象量很大,尤其反序列化 CQL 行时会创建大量短生命周期对象。我通常给 executor 加-XX:+UseG1GC,并且设置-XX:MaxGCPauseMillis=200,比默认的 Parallel GC 更平稳。这个优化在堆内内存超过 8G 时效果尤其明显。

5. 常见问题与排查技巧实录

5.1 连接超时与集群不可用排查

最常见的问题就是任务启动后报All connection pools are busy或者Connection refused。

我排查这类问题的顺序是先看连接器日志,确认它实际连的是哪台 Cassandra 节点。很多时候你以为它连的是 seed 节点,但连接器会根据元数据去连所有数据节点,如果某些节点的rpc_address配置不对,Spark 端根本访问不到。

还有一种隐蔽场景:Cassandra 的 listen_address 和 rpc_address 配置。我在生产上遇到过 Spark 能连种子节点,但后续读数据时报 connection refused,后来发现是 Cassandra 节点间 broadcast 用的地址是内网 IP,而 Spark executor 在另一个网段,网络策略没放通。这个问题的排查坑就在于报错信息不直接说“IP 不通”,而是各种超时。处理方法是检查 Cassandra 配置文件中的rpc_address,保证 Spark 可以路由到该地址,并在防火墙放行 9042 端口。

连接池超时一般调大spark.cassandra.connection.pool.connection.pool.max.size到 8-16,解决高并发下连接池爆满的问题。

5.2 谓词下推失效分析

“我明明在 Spark 里加了 where 条件,为什么 Cassandra 还是慢成狗?”这是群里经常被问到的问题。

这里不怪大家,因为连接器对谓词下推的规则确实有点隐蔽。它支持的下推条件包括主键等值、聚簇列范围、IN子句等,凡是可以翻译成 CQL WHERE 的才会下推。如果 filter 里有函数运算,比如date_trunc("day", col("ts")) > ...,连接器不会下推,只能全表扫。

排查方法很简单:打开连接器 debug 日志或者把生成的查询打出来。用日志模式运行:

spark-submit --conf spark.cassandra.debug.level=DEBUG ...

日志里会打印Generated CQL query: SELECT ... WHERE ...,一眼就能看出哪些条件被下推了。我之前就是这样发现某个任务里隐藏的转换函数让下推失效,去掉之后从 2 小时直接降到 20 分钟。

5.3 TTL 与时间戳类型踩坑

Cassandra 表可以设置 TTL,数据到期自动删除。这个机制在写入时对 Spark 流水线有个隐藏坑:Spark DataFrame 的标准时间类型是 Timestamp,但 Cassandra 的timestamp类型存的是 UTC 时间。如果你的数据源是本地时区,直写的话时间会偏移 8 小时。

我踩过一次:调度系统按本地时间统计日活,结果 Cassandra 里存的是 UTC,第二天跑批时发现数据差 8 小时,线上报表连续错了两天。后来在 ETL 的字段转换阶段统一做to_utc_timestamp/from_utc_timestamp,才解决这个问题。

另一点是 TTL 只对写入时指定,Spark 连接器写回时如果不给每条记录指定 TTL,默认是无过期。如果你希望结果表数据保留一定周期,可以在写入时通过把 TTL 作为每行的列值写入,比如连接器支持writetime()和 TTL 函数:

import org.apache.spark.sql.cassandra._ analyticsDf.write .cassandraFormat("event_agg_daily", "analytics") .option("ttl", "86400") .save()

注意这里的 TTL 单位是秒。设了 TTL 之后,读取端可能出现数据“突然变少”,特别是集群节点间数据还没完全同步时,容易产生短时间的数据不一致。如果下游任务对数据完整性敏感,建议不要把 TTL 设得太短,或者让读取端容忍短窗口的不一致。

写在最后:几个值得记住的实战心得

上面这些经验,都是我从一个个线上问题里“喂”出来的。个人最深的体会是:Cassandra + Spark 集成本身并不难,难的是理解两个系统各自的脾气。Cassandra 擅长按主键索引点查、顺序扫描分布式 token range,但它的劣势也很明显——一切让 Cassandra 做全表扫描的尝试都会吃大亏。Spark 则相反,它天生为“暴力计算”而生,所以集成时一定要把“让 Cassandra 少干活、让 Spark 多干活”刻在脑子里的每一步设计里:能下推的谓词尽量让 Cassandra 过滤,只能在 Spark 做的复杂计算坚决不要尝试用 CQL 实现。

最后一个实用小技巧:正式上线前,一定要在准生产环境跑一次全量加增量验证,然后观察 Cassandra 节点的 CPU 和 GC 日志。如果某个热点表把节点 CPU 打满,你应该优先调整主键设计和连接器的分区策略,而不是盲目加 Spark 资源。类比的粗浅说法就是,水管漏水你要先修管道,而不是一味加泵。这套流水线搭好之后,后续扩展方向也很明确:可以接 Kafka 做实时流批一体,也可以把清洗后的数据直接落到 Iceberg 或 Hudi 做湖仓分层,落地空间挺大的。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询