ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Cassandra与Spark集成实战:构建高效大数据处理流水线

Cassandra与Spark集成实战:构建高效大数据处理流水线 做数据后端这几年我一直觉得 Cassandra 和 Spark 是小团队最容易低估、但实际收益极高的一对组合。前者天生为海量写入而生分布式的数据模型让它在大规模并发插入场景下几乎没有对手后者则是离线计算和实时流处理的事实标准。把这两个一起纳入一条大数据处理流水线意味着从数据落库到分析之间省掉了“导出 CSV 再导入数仓”这种中间搬运步骤架构上会清爽很多。这篇文章我想把一套可落地的集成方案完整拆开来讲内容覆盖为什么要选这对组合、环境搭建和版本兼容怎么避坑、Spark 读写 Cassandra 的配置与 API 用法、性能调优的关键参数包括你大概率会被线上状况逼疯的 Spark 内存问题、以及如何把批处理和调度串成一条稳定的流水线。适合已经有一点大数据基础、准备把 Cassandra 接入业务分析链路的人参考也适合第一次在集群环境里集成这两个组件时想少走弯路的读者。1. 为什么要把流水线搭在 Cassandra 和 Spark 这对组合上1.1 两个组件天然互补省掉中间搬运Cassandra 的强项不是做复杂关联计算而是以分布式、多副本、可线性扩展的方式承接海量写入流量。它的数据模型是面向列族的宽表一张表可以轻松保存上亿行数据并且在正常运维下很少出现性能骤降。但它的弱项也同样明显要是你想跑一段复杂的聚合分析比如按天统计不同渠道的用户活跃情况、对时间窗口做滑动计算直接用 CQL 写 JOIN 和嵌套子查询非常痛苦而且会侵占线上业务的 coordinator 资源搞不好把在线查询一起拖垮。Spark 的角色正好补上这块短板。它把计算任务切分成多个 stage在 executor 上分布式执行对 Cassandra 里几亿行的全表扫描、过滤、聚合、join都是它的常规操作。所以最常见的架构形态是线上请求的写入继续走 Cassandra离线或近线的分析任务则由 Spark 直接读取 Cassandra 数据完成计算结果再写回 Cassandra 或其他存储。这就直接把“数据仓库搬运层”给删掉了。1.2 连接器就是两者之间的转换桥真正把 Cassandra 和 Spark 连接起来的是 spark-cassandra-connector。它的核心工作是两件事把 Cassandra 表映射成 Spark 里的 DataFrame 或 RDD让 Spark 可以用熟悉的 DataFrame API 做计算以及把 Spark 处理完的结果以批量写入的方式落回 Cassandra。它内部实现了分区切分、一致性控制、批量提交、谓词下推这些逻辑不需要你自己手写底层协议。很多人看到“集成”两个字以为要自己封装 CQL 解析器或者写 JNI 接口但实际上连接器已经把脏活累活做完了。你需要掌握的只是几个配置项和 API。不过在动手之前版本兼容这一点如果不处理好后面会被坑得很惨。1.3 典型业务场景描一遍我接触过的项目中这套组合用得相对多的场景是这些用户行为分析客户端埋点数据直接进 CassandraSpark 每天凌晨跑一次聚合产出用户留存、活跃、路径漏斗等指标写回结果表。IoT 时序数据设备上报的指标按时间维度落库Spark 按小时窗口做温度、能耗等指标统计异常设备再通过流式任务告警。推荐特征计算线上推荐系统需要的候选集特征由 Spark 读取 Cassandra 里的历史行为数据离线生成再同步给在线服务。事件驱动报表业务系统产生的事件通过消息队列进入 CassandraSpark Structured Streaming 消费这批数据做近实时统计刷新大屏和报表。这些场景的共同点是写入量大、数据规模大、计算逻辑复杂。而 Cassandra Spark 的组合不需要额外引入 Hadoop 生态那一整套东西就能把这些问题处理完。2. 环境准备版本兼容、集群搭建与连接器选型2.1 版本兼容矩阵如果你的项目是同步进行的这里给你一份可以拿来参考的组合方式组件推荐版本范围说明Cassandra3.11.x / 4.0.x / 4.1.x4.0 以上版本建议开启全量一致性读性能更稳定Spark3.2.x / 3.3.x / 3.4.x3.3 之后的 SerializedSafe 支持对 Cassandra 更友好connectorscalaJava 8 / 11如果 Cassandra 是旧版本建议 Java 8spark-cassandra-connector与 Spark 版本配套的 3.x注意连接器与 Spark 大版本必须严格匹配连接器版本的命名通常像这样spark-cassandra-connector_2.12:3.3.0前面的2.12是 Scala 版本后面的3.3.0是对应 Spark 的大版本。如果你的 Spark 是 3.3.x就找 3.3.0 的连接器Spark 是 3.4.x就找 3.4.0 的连接器。乱配的后果往往是运行时直接抛NoSuchMethodError或NoClassDefFoundError排查起来非常浪费时间。2.2 Spark 集群搭建与 Cassandra 网络互通Spark 集群我用得比较多的是 Standalone 模式理由很简单不引入 YARN 和 Hadoop组件少、排查问题少。搭建的时候有几点容易被忽略master 和 worker 节点需要能够互相访问并且 worker 节点必须能直连 Cassandra 的rpc_address通常是 9042 端口。很多问题都出在 worker 能访问 Cassandra 而 master 不能或者反过来结果任务一跑就报连接超时。每个 worker 的可用内存不要一次性全部分配给 executor要预留一部分给系统缓存和 Spark 自身开销。比如机器 64GB 内存spark.executor.memory建议控制在 32GB 到 40GB 之间剩下留给 OS page cache。如果你的 Cassandra 是多数据中心部署每个 Spark worker 最好放在同一个数据中心内通过spark.cassandra.connection.local_dc指定本地数据中心避免跨机房读数据造成大量延迟。启动 Spark 集群的命令大致如下# 在 master 节点启动 $SPARK_HOME/sbin/start-master.sh -h 192.168.10.10 -p 7077 # 在 worker 节点启动 $SPARK_HOME/sbin/start-worker.sh spark://192.168.10.10:7077 -m 40G -c 8-m指定该 worker 可使用的内存总量-c指定空闲核心数。到这里基础环境就绪下一步把 Cassandra 集群信息接进 Spark 会话。2.3 连接器依赖的引入方式最常见的做法是在spark-shell或spark-submit时用--packages引入连接器依赖省去手动下载 jar 的麻烦spark-shell \ --packages com.datastax.spark:spark-cassandra-connector_2.12:3.3.0 \ --conf spark.cassandra.connection.host192.168.10.20,192.168.10.21 \ --conf spark.cassandra.connection.local_dcdc1如果你走的是代码提交模式在build.sbt中加依赖就行libraryDependencies com.datastax.spark %% spark-cassandra-connector % 3.3.0这里必须提醒一句连接器对 Cassandra 原生传输协议版本有要求Cassandra 4.x 默认使用 native protocol v5较老的连接器可能只支持 v4。遇到握手失败时先看 Cassandra 的native_transport_max_protocol_version配置再确认连接器版本很多“连不上”其实是协议版本不匹配。3. 核心读写从 SparkSession 配置到 DataFrame API3.1 SparkSession 初始化与连接配置无论读还是写第一步都是在 SparkSession 里把 Cassandra 节点的连接信息放进去。我常用的配置是这样的from pyspark.sql import SparkSession spark SparkSession \ .builder \ .appName(CassandraSparkPipeline) \ .master(spark://192.168.10.10:7077) \ .config(spark.cassandra.connection.host, 192.168.10.20,192.168.10.21) \ .config(spark.cassandra.connection.local_dc, dc1) \ .config(spark.cassandra.auth.username, spark_user) \ .config(spark.cassandra.auth.password, ******) \ .getOrCreate()有几个配置项理解清楚之后排查问题会简单很多spark.cassandra.connection.host填 Cassandra 任意几个节点的 IP 即可连接器会自动获取整个集群拓扑。但注意这里填的节点必须对 Spark worker 网络可达。spark.cassandra.connection.local_dc建议必填。不填的时候连接器默认请求当前连接节点的 DC一旦 Spark 和 Cassandra 不在同一个机房额外的跨 DC 查询会严重影响性能。spark.cassandra.auth.username/password生产环境必须用专门服务账号不要给超级权限。3.2 从 Cassandra 读取数据读取的核心是使用format(org.apache.spark.sql.cassandra)加上options指定 keyspace 和 table 名。以 Scala 为例val rawEventsDF spark .read .format(org.apache.spark.sql.cassandra) .options(Map( keyspace - user_events, table - raw_events )) .load() rawEventsDF.printSchema() rawEventsDF.createOrReplaceTempView(raw_events)这里有个很关键的概念这个 DataFrame 不是把整张表一次性拉到 Spark 内存里。连接器会根据 Cassandra 表的分区分布把表切分成多个 split再按分区推送到 executor 上执行。你执行count()时它其实触发的是每个分区的并行计数而不是先全量拷贝。所以你在写查询时尽量把过滤条件写在 Spark 的 API 层让连接器能把where下推到 Cassandra 端执行。例如val filtered spark .sql(SELECT user_id, event_time, event_type FROM raw_events WHERE event_type click AND event_time 2024-06-01)只要过滤字段是 Cassandra 表的主键或索引列连接器就会在下推阶段生成对应的 CQL WHERE 条件而是不是全表扫描性能差别是几个数量级。3.3 写入 Cassandra 与 JSON 字段处理写入的 API 同样很直接resultDF .write .format(org.apache.spark.sql.cassandra) .options(Map( keyspace - user_events, table - event_stats )) .mode(append) .save()写入时需要注意的强制约束是目标表必须有明确的主键设计。Cassandra 的表如果没有主键Spark 写入会直接报错。这个主键在 Cassandra 中就是PRIMARY KEY定义的分区键和聚簇键写操作本质上是对这些键的 upsert。另外实际项目里经常遇到 Cassandra 表存的是原始 JSON 字符串而 Spark 计算时需要把 JSON 字段展开成结构化列。我用的是 Spark 内置的from_json函数from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, LongType schema StructType([ StructField(channel, StringType()), StructField(duration, LongType()), StructField(page, StringType()), ]) parsed rawEventsDF \ .select( user_id, event_time, F.from_json(F.col(payload), schema).alias(data) ) \ .select( user_id, event_time, F.col(data.channel).alias(channel), F.col(data.duration).alias(duration) )这种做法的好处是不需要提前把 JSON 拆成独立字段再入库写入端保持灵活分析端按需解析恰好发挥了两边各自的优势。3.4 只读指定列和谓词下推的经验连接器还有一个容易被忽视的能力是“只读需要的列”。Cassandra 宽表经常有几十个字段如果你只需要其中三列在.load()后面的 select 只做投影连接器仍然会把你需要的列映射成 CQL 选择列表。我的建议是在 DataFrame 上尽早投影、尽早过滤让下推生效减少从 Cassandra 到 Spark 的网络传输量。实际测试里如果我们把一条全字段扫描的任务改成“只读 5 个必要字段 主键过滤”同样规模的数据在 Cassandra 端耗时可以降 40% 以上。这不是 Spark 的功劳而是连接器的下推做得足够好。前提是你愿意花十分钟检查一下df.queryExecution().executedPlan()里的扫描节点确认下推真的触发了而不是默默做了全表拉取。4. 性能调优分区、一致性、内存与写入批量参数4.1 读取阶段的分区与并发控制Spark 读 Cassandra 时任务被拆成多少个分片直接决定并发度。控制这一行为的关键参数是spark.cassandra.input.split.size_in_mb连接器会尽量把表的数据按这个大小分成若干 split 交给不同的 executor 处理。把值调小分片变多、并发变高适合节点资源比较多的时候调大分片变少、并发降低适合避免过度抢占线上 Cassandra 节点资源。一个常见实践是如果 Cassandra 和 Spark 各占独立资源池split.size_in_mb可以压到 32 或 16尽量把整个集群的 CPU 跑满如果是共享集群建议保持 64 及以上防止瞬时读放大影响在线服务。另外读取的一致性级别建议默认设为LOCAL_QUORUMspark.cassandra.input.consistency.levelLOCAL_QUORUM这能保证读到的是多数副本的数据同时不会跨数据中心等待响应兼顾一致性和延迟。如果没有特殊的一致性要求不建议使用ONE因为一旦某个副本出问题会读到缺失数据。4.2 写入阶段的批量参数写入 Cassandra 时连接器也是走批量 upsert 的方式。性能相关的参数有两个值得重点调参数默认值调优思路spark.cassandra.output.batch.size.in.rowsauto每批写入行数过大容易引发 Cassandra 端超时spark.cassandra.output.batch.size.in.bytes16KB每批数据大小两个参数同时满足才成批spark.cassandra.output.consistency.levelLOCAL_QUORUM写入所需的副本确认数我还习惯在写入量大的任务里配合spark.cassandra.output.concurrent.writes控制并行写线程数。如果 Spark 写 Cassandra 的同时还在跑其他任务并发写太高会让 coordinator 节点 CPU 打满把它限制在 4 到 8 之间通常会比较稳妥。这里的底层逻辑是Cassandra 的批量写入不是越多越快每个 batch 都要经过 coordinator 的分发确认batch 过大会造成超时重试反而更慢。4.3 Spark 内存参数——最容易踩的坑“Spark 内存”这个关键词背后藏着一系列线上事故。我在多个项目里见过类似的现象Cassandra 集群很健康Spark 任务却莫名其妙 executor 挂掉日志里清一色的OutOfMemoryError或Container killed by YARN。根因几乎都是 Spark 侧内存参数配置不当。我的建议是从这几个维度排查和调整spark.executor.memory单个 executor 可用的堆内存。如果每个任务处理的数据量大、shuffle 量高就把它调高。但注意executor 数量太多时总内存会被 Cassandra 多副本读取放大适得其反。spark.executor.cores单个 executor 的并发能力。一般是 4 到 8 个核心。如果核心太多而内存没跟上容易把堆内存耗尽。spark.sql.shuffle.partitions默认 200但在读写 Cassandra 场景里如果数据量不大200 个分区会产生大量小任务拉低效率。根据你的数据量适当调低比如 48 到 96 之间。spark.memory.offHeap.enabled和spark.memory.offHeap.size如果开启了 off-heap记得给堆外内存留足够余量否则序列化和 JDBC 等操作会 OOM。一个经典调法是把总 executor 内存先设成 32GBspark.memory.fraction0.4Spark 旧版本默认 0.6试试性能再根据 GC 日志微调。我实际线上的稳定配置是spark.executor.memory32G spark.executor.cores6 spark.sql.shuffle.partitions96 spark.memory.fraction0.4 spark.memory.storageFraction0.5这只是基线具体还是要以自己的数据和集群规格为准。记住一点内存调整不是一劳永逸遇到新任务类型先跑小型采样数据验证再上全量。4.4 避免读取热点分区导致的计算倾斜Cassandra 数据模型里同一分区键的数据会被路由到同一个节点。如果某个 key 的访问频率远高于其他 keySpark 读它就容易形成单点热点明明是分布式计算却只有一两个 executor 在忙。应对措施通常在数据建模时就要考虑设计分区键时不要用单一高基数字段而是把“高基数 低基数”组合成复合分区键让数据尽量分散。分析任务中如果对某个热点 key 做过滤可以先用 CQL 层面做预聚合或者加 salting 字段。Spark 侧可以给 DataFrame 加一层自定义分区器把热点 key 打散到多个 executor 上但这要配合业务代码一起改成本较高。我在处理用户日活统计分析时就踩过类似问题某些头部用户一天产生的行为事件是普通用户的上万倍扫这两三个 key 时单 executor 要跑几十分钟其他 executor 早就空转结束。后来把分区键从单一的user_id改成(user_id, event_date)数据终于均匀分散任务整体时间从 1 小时降到 7 分钟。5. 流水线落地调度、增量处理与监控5.1 从 Kafka 到 Cassandra 再到 Spark 的完整链路集成方案最终要落成一条可持续运行的流水线而不是手动在 spark-shell 里敲几次命令。我常用的流水线链路是这样的业务数据实时或准实时写入 Kafka或者直接写入 Cassandra 业务表。如果数据先到 Kafka用 Spark Structured Streaming 消费并写入 Cassandra 原始表这一步承担了消息落库的工作。批处理任务定时读 Cassandra 里的原始表做清洗、聚合写回结果表。结果表再被在线服务、报表平台或 BI 工具直接查询。这条链路的好处是每层职责单一Cassandra 是数据底座Streaming 保证新鲜度批处理保证口径稳定。你不需要为了“实时”而实时很多业务只需要小时级新鲜度批处理就够用。5.2 用 DolphinScheduler 编排批处理任务在调度层DolphinScheduler 是一个很适合这套组合的选型。它的 shell 任务可以直接执行spark-submit同时提供依赖管理、失败重试、告警等机制比在 crontab 里裸跑 spark-submit 要安心得多。一个典型的调度工作流如下任务 A校验 Cassandra 集群连通性和待处理的 keyspace 是否存在。任务 B执行 Spark 批处理任务使用spark-submit提交主类读取原始表并产出聚合表。任务 C校验聚合表行数与采样数据失败则触发重跑。任务 D清理临时目录、写入成功标记。DolphinScheduler 的任务配置里命令可以写成这样spark-submit \ --master spark://192.168.10.10:7077 \ --deploy-mode cluster \ --class com.example.AggregateJob \ /data/jars/cassandra-spark-pipeline-1.0.0.jar \ --keyspace user_events \ --srcTable raw_events \ --dstTable event_stats调度里的依赖关系会让任务 B 只有在任务 A 成功后才会启动失败时按配置重试还支持任务超时和日志滚动。这些是裸 crontab 不容易做到的。5.3 增量处理与数据一致性保障所有批处理任务都会面临同一个问题增量怎么定义。Cassandra 表没有内置的自增 ID最稳妥的做法是使用时间戳字段作为增量切分点。比如event_time是按小时写入的那批处理任务就处理[上次成功时间, 当前时间)区间的数据。实现时利用 Spark 的谓词下推把时间窗口条件拼进 Cassandra 读取的过滤条件里。跨任务的数据一致性我的做法是“幂等写”。每次聚合结果写回 Cassandra 时主键包含统计维度 统计时间窗口例如主键字段说明stat_date聚合日期channel_dim渠道维度event_type事件类型这样即使某个时间窗口的任务因为失败重跑了几次最终的写入结果也是一致的不会产生重复累计。这比在外部维护一套 offset 或 last_processed_id 要简单得多也天然抗失败。5.4 监控的日常操作Spark 历史服务器和 Cassandra 的监控面板是必要工具。我习惯记录三个维度的指标任务运行时长正常情况下批处理在固定时间窗口内完成一旦环比上升 30% 以上就要看是不是数据量增长、split 倾斜或 Cassandra 节点延迟。executor GC 时间如果 GC 时间超过任务运行时间的 20%多半是内存配置不合理需要调整 executor 内存或减少 shuffle 分区数。Cassandra 节点协调延迟如果 Spark 读写期间 Cassandra 的 coordinator CPU 飙高就需要降低连接器并发参数。监控不一定要上复杂的平台先建一个基础面板把这三个指标的曲线拉出来已经能应对 90% 的故障预警。6. 踩坑实录连接失败、数据倾斜和写入异常排查6.1 “连接超时”但 Cassandra 明明活着这是最常见的一个误判。现象是 spark-shell 能进去但一旦执行读取就报类似Exception while executing SELECT ...或All host(s) tried for query failed。排查链路通常是先在 Spark worker 节点上用cqlsh连 Cassandra验证 9042 端口是否通。检查 Cassandra 的rpc_address是否设置为 worker 可达的 IP而不是只绑定在内网 IP 上。检查连接器版本是否与服务端 native protocol 版本匹配尤其是 Cassandra 4.x 配旧版连接器时。检查是否误把 CQL 端口9042和 Thrift 端口9160搞混Thrift 在 Cassandra 4.0 之后默认禁用。很多时候 Cassandra 进程看起来正常但那只是因为nodetool status显示 UN实际上 Spark 连接器连的是另一个逻辑数据中心或另一个网段网络策略没放行而已。6.2 写入报错 “Key may not be empty” 或主键字段缺失这个报错几乎都是因为目标表的分区键字段在 DataFrame 里是 null 或空字符串。Cassandra 对主键有强约束空值直接拒绝写入。我的排查思路是先打印 DataFrame schema再看目标表 CQL 定义的主键顺序确认 DataFrame 列名与表结构完全一致。特别是用 Scala 写时字段名大小写不能错Cassandra 默认大小写不敏感但 CQL 里加了引号的 xxxxx 会和 Spark 列名对不上。还有一种隐蔽情况DataFrame 里包含多个同名但不同大小写的列写入映射时也可能导致找不到主键。用.select()先把主键列名统一再做 upsert能避开这个问题。6.3 数据倾斜导致某个 executor 长期不结束我在前面性能一节提过热点分区的问题这里补充一个排查方法。当任务卡住不结束去看 Spark UI 的 stage 页面如果某个 task 的 duration 明显高于其他 task 几个数量级基本可以断定是数据倾斜。具体思路是这样先确认倾斜的表的分区键是什么再用 CQL 手工查该分区下有多少行数据确认是不是那几个 key 撑爆了单个 task。如果无法改表结构临时解决方案是给查询条件加一个高基数的过滤字段比如日期让连接器把数据拆成更小的分片。但如果倾斜是业务固有的比如头部用户效应最终方案还是要做数据模型改造把单一热点 key 拆成多个加盐 key。这个改造成本不低所以建议在表设计初期就把这个风险评估进去。6.4 写放大与磁盘占满的隐性风险Cassandra 的写路径是 append-only写入快但后台 compaction 会把小 SSTable 合并成大 SSTable这个过程中磁盘 IO 和空间消耗都会显著增长。Spark 批量写入往往会在短时间内制造大量小 SSTable如果 compaction 跟不上数据文件会暂时膨胀极端情况可能把磁盘打满。我自己遇到过一次比较惊险的情况Spark 任务写入的量并不算大但因为每条数据的主键几乎都不相同Cassandra 端一直在做大量 compaction 操作导致写入延迟从 5ms 飙到 400ms最终任务超时重试。后来调整了写入批大小并降低了写入并发同时给目标表设置了合适的 compaction strategy比如SizeTieredCompactionStrategy情况才逐步稳定下来。所以在做性能测试的时候不要只看 Spark 任务本身的时间还要盯着 Cassandra 节点的磁盘占用和 compaction 指标。有时候“Spark 写完了”不等于“系统稳定了”它只是把压力推迟到了后台。6.5 重试逻辑一定要配合幂等设计最后一个教训和代码本身无关但非常重要。批处理任务一旦失败重跑如果写入逻辑没有幂等设计很容易把指标加重复。比如统计用户点击量如果每次重跑都把全量数据累加一遍最终结果必然翻倍。我推荐的做法前面也提过聚合结果表的主键一定要包含统计窗口和维度字段写入模式统一用 append这样重跑只会覆盖相同主键的数据不会产生重复行。这个设计在 Cassandra 这种基于主键 upsert 的存储里极其自然也几乎是零成本实现。如果你做的是多表 join 后写明细表那幂等的粒度要落到明细行的唯一键上同样把唯一键作为 Cassandra 的主键。设计阶段多花一点时间后续运维会少掉无数个失眠的夜晚。
返回列表