
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.13.2.x2.123.2.03.3.x2.123.2.0 以上2.4.x2.112.4.3注意连接器的 groupId 是com.datastax.sparkartifactId 是spark-cassandra-connector-assembly也就是带所有依赖的打包版本。如果你用 Maven可以这样引入dependency groupIdcom.datastax.spark/groupId artifactIdspark-cassandra-connector-assembly_2.12/artifactId version3.2.0/version /dependency这里有个容易踩的坑连接器内部带的guava、netty等依赖经常和 Spark 自带的版本冲突。表现为启动时一堆 NoSuchMethodError 或 ClassNotFound非常迷惑。强烈建议在提交任务时使用--conf spark.driver.userClassPathFirsttrue和--conf spark.executor.userClassPathFirsttrue让用户 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.host172.16.1.10 \ --conf spark.cassandra.connection.port9042 \ --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 计算结果写回 CassandraETL 清洗完下一步往往是把结果写回 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 提升到 5000bytes 提升到 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 4096split.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.memoryOverhead2g --conf spark.memory.storageFraction0.2 --conf spark.sql.shuffle.partitions400storageFraction默认是 0.5意思是 storage 内存和 shuffle 执行内存各占一半。像这种读源表做计算再回写的流水线shuffle 占比高storage 占比低把 storage 调到 0.2能显著减少 GC 压力。不过这个值也别压太狠如果缓存用的多或者 join 后要复用 DataFrame还是要留出空间。GC 方面Cassandra 连接器产生的对象量很大尤其反序列化 CQL 行时会创建大量短生命周期对象。我通常给 executor 加-XX:UseG1GC并且设置-XX:MaxGCPauseMillis200比默认的 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.levelDEBUG ...日志里会打印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 做湖仓分层落地空间挺大的。