拓冰建站拓冰建站
首页 / 资讯中心 / 正文

ClickHouse与Spark集成实践:混合大数据处理方案与踩坑指南

ClickHouse与Spark集成这件事我踩了将近两年的坑才把怎么混着用想明白。问个很现实的问题你的数仓里到底是用ClickHouse扛报表还是用Spark跑清洗可能很多人第一反应是都用但真正落地时就会发现——明明都是大数据处理界的明星工具真让它们协作起来却总有说不出的别扭。这篇博文就是围绕ClickHouse和Spark的集成把混合大数据处理方案讲透它们各自擅长什么、什么场景才需要把它们拼在一起、拼的时候有哪些能直接照抄的方案和代码以及我在生产环境里踩过的那些文档里查不到的坑。无论你是正在做技术选型的架构师还是已经在跑Spark和ClickHouse但被衔接问题折磨的工程师这篇文章都值得看完。1. ClickHouse与Spark的能力边界为什么选型摆错位置最容易翻车1.1 一个真实的翻车案例先说一个我亲历的案例这也是促使我真正研究这俩工具集成的起点。2023年我们团队负责一个电商数据分析平台最初架构很简单业务库通过Canal同步到Kafka然后一股脑灌进ClickHouse。白天还好一到晚上跑全量分位数计算、多表关联、漏斗分析ClickHouse的CPU就飙到90%以上查询动不动几十秒甚至超时。后来有个同事提出用Spark代替ClickHouse理由是Spark计算能力强。我们真把ETL和报表查询全迁到了Spark上结果更惨——明细查询本来秒级出数变成了十几秒起步报表页面经常转圈转半天。这个经历让我彻底明白了一个道理ClickHouse和Spark根本不是替代关系而是分属两个完全不同的生态位。选型摆错位置比不选更可怕。1.2 两个工具到底各自擅长什么ClickHouse本质上是一个列式存储的OLAP数据库它的DNA是查询——把千万行、上亿行数据放进内存里做聚合、过滤利用向量化执行引擎和稀疏索引在几毫秒到几秒内返回结果。它最舒服的场景是大宽表、明细查询、多维聚合、漏斗分析、用户画像、实时报表。Redis是键值的存储ClickHouse就是分析领域的内存加速器而且它有完整的列式压缩磁盘占用可以压到原始数据1/5甚至更低。但ClickHouse天生不擅长什么复杂的多表关联。虽然它支持JOIN但一旦关联的表大了、关联次数多了内存和临时文件的压力非常大性能会断崖式下跌。它也不适合做高并发点查毕竟不是为OLTP设计的。还有一个容易被忽略的点ClickHouse的SQL表达能力有限窗口函数、UDF虽然新版在补但和Spark SQL、PySpark生态相比还是太单薄。Spark则恰恰相反。Spark的核心是计算——它把数据分散到集群内存中通过RDD/DataFrame做任意复杂的变换能跑ETL、能算特征、能训练机器学习模型、能做流批一体处理。它的SQL语法完整度接近标准SQLUDF、窗口函数、复杂JOIN都是家常便饭。但Spark的问题也很突出调度开销大启动一个Task要几十毫秒到几百毫秒查询结果返回的路径长从Driver到Executor再到用户想做低延迟交互查询基本没戏。说白了Spark是算得动的批量计算器不是查得快的搜索引擎。而且它处理小文件、小数据量时效率极低1000万以下的数据量用Spark跑一次任务的时间往往比ClickHouse直接查还要慢。1.3 两个工具的正确分工经过那一次翻车案例我的准则变成了这样一句话数据清洗和计算交给Spark数据存储和高频查询交给ClickHouse。也就是让Spark把数据做熟再喂给ClickHouse去上菜。举几个具体的分工场景。场景一日志清洗。原始日志往往有脏数据、嵌套JSON、重复记录这种工作交给Spark用DataFrame做格式转换、字段提取、去重清洗完按天分区写出列存储格式然后导入ClickHouse。若是直接把原始日志灌进ClickHouse不仅字段数量膨胀、存储成本高而且后续每个查询都要带上清洗逻辑时间久了表结构都会失控。场景二多表关联加工。假设要做用户生命周期分析需要把订单表、用户表、商品表、退款表四张表关联。这四张表各有几千万到几亿行在ClickHouse里做一次四表大JOIN内存压力大得能让你彻夜难眠。正确做法是在Spark里做一轮关联加工成宽表输出用户ID行为标签金额等关键字段最后进ClickHouse。之后查询都变成单表聚合ClickHouse几乎不费吹灰之力。场景三离线统计与近实时查询的分层。日活、GMV这类指标每天由Spark跑批任务计算好结果存入ClickHouse的报表表而点击流这种需要实时查看的明细直接写入ClickHouse。前者是离线算好、实时查询后者是实时摄入、实时查询两条链路都合理但永远不要让Spark去承接用户的交互式查询请求也不要让ClickHouse承担超过其能力上限的重型计算。注意这里说的交给Spark不是指Spark分布式集群跑一天一夜而是要清楚它适合跑那些分批、批量的复杂计算任务。真正的优化不是删掉某一个工具而是让每个工具在自己的位置上干活。2. 按业务场景选型五条ClickHouse Spark集成链路当我明确了两个工具的分工后接下来就是怎么把两者连起来。网上资料大多是单点演示真正能指导生产选择的很少。我一共总结了五条集成链路每条都有明确的适用前提和代价你可以直接对照自己业务选型。2.1 轻量集成JDBC直连适合小流量场景这是最直觉、最简单的集成方式——Spark作为客户端通过JDBC驱动直接读取ClickHouse中的数据表或者把计算结果写入ClickHouse。前提是数据量不大单次拉取控制在百万行以内或者只是做维度表补全等小规模关联。具体来说在Spark的DataFrame API里配置jdbc数据源指定ClickHouse的连接地址和表名Spark会生成一个RDD把ClickHouse表当成一个分布式数据集来用。我们早期做的一个用户标签补全任务就是这么干的Spark读取订单明细每行数据需要关联用户的基础属性用户表一共200万行存放在ClickHouse里Spark在启动时通过JDBC拉取这张表广播出去然后在map阶段补全字段。整个任务跑完不到5分钟代码量很小也不用额外维护同步链路。JDBC直连的底细说穿了就是Spark把SQL推给ClickHouse执行然后拉结果。所以它适合小数据量的读取不适合大数据量的全量同步。2.2 离线批量Spark算完ClickHouse承接加速这是我目前最推崇的一条链路也是混合大数据处理方案里最能体现价值的部分。整体流程如下Spark从HDFS/Hive/数据湖中读取大规模原始数据进行复杂的清洗与关联加工加工好的结果数据以Parquet/ORC等列式格式写出到临时目录或对象存储ClickHouse通过表函数或clickhouse-client将数据批量导入自身的MergeTree表。这条链路的精髓在于Spark负责算ClickHouse负责查。计算量再大也是Spark承担的ClickHouse得到的永远是干净的、可索引的、低冗余的结果数据。查询时ClickHouse利用稀疏索引和列式存储秒级返回BI前端和即席查询都轻松应对。我们做过一个典型的离线报表平台Spark每天晚上从ODS层读取全天的订单、流量、行为明细经过近十步关联和清洗生成一张维度宽表大概每天5000万行然后导入ClickHouse。下次查询日活、GMV、转化率这些指标时ClickHouse响应时间基本在200毫秒以内。换成之前直接用Spark跑查询一次至少等半分钟业务方完全无法接受。2.3 准实时链路Kafka Spark Structured Streaming ClickHouse当数据需要秒级到分钟级的时效性时JDBC直连和离线批量都不够得引入消息队列和流计算。链路是业务数据通过Canal/Debezium捕获变更或直接由应用发送埋点日志到KafkaSpark Structured Streaming从Kafka消费在流上进行轻量的过滤、补全、聚合然后以微批的方式写入ClickHouse。需要注意ClickHouse官方也提供了Kafka引擎表可以直接用CREATE TABLE ... ENGINE Kafka消费Kafka数据不需要Spark也行。那什么时候要引入Spark呢答案是当Kafka里的原始格式需要复杂转换时——比如嵌套JSON展开、多流join、字段加密脱敏、去重。Kafka引擎表只能做简单的SQL转换稍微复杂一点就撑不住了。用Spark Streaming的好处是流任务和离线批任务可以共用同一套数据处理逻辑维护成本低。而且Spark对checkpoint、exactly-once语义支持比直接消费Kafka靠谱得多。2.4 加速层方案ClickHouse为数据湖/数仓加速这一条实际上是2.2的进阶变体。很多团队已经引入Iceberg、Hudi、Delta Lake这类数据湖技术用Spark读写湖表构建开放的湖仓架构。但湖表查询性能始终是个痛点尤其BI工具直接查Hive/Iceberg表往往需要跑MapReduce或Spark任务延迟根本压不下来。我的做法是把数据湖当全量存储把ClickHouse当热数据加速层。Spark负责将湖表的最新分区数据做增量计算输出到ClickHouse当用户要查明细或聚合时请求优先打到ClickHouse只有当ClickHouse里没有该粒度的数据时才回退到Spark查询湖表。这样既保留湖的开放性和低成本存储又获得ClickHouse的秒级查询体验。2.5 迁移场景怎么做还要单独提一类场景就是ClickHouse数据库整体迁移。热门搜索里这个词曝光很高说明踩过坑的人不少。ClickHouse迁移有两条路一是用官方Backup/Restore适合停机窗口内做物理备份恢复二是如果你需要跨版本、跨集群甚至异构环境迁移同时要做数据转换那就用Spark读旧ClickHouse全量数据写Parquet到对象存储再通过2.2的导入方式灌入新ClickHouse集群。优点是可以顺便做数据清洗、类型转换、分区分桶重排缺点是要考虑双跑期间的增量对齐。我们当时做跨机房迁移就是Backup/Restore处理全量、Kafka链路处理增量、最后用Spark补数校对三箭齐发才做到数据零丢失。2.6 五条链路的选型对照为了方便选择我把五条链路的适用场景和关键参数整理成了一个表。链路时效性数据量级建议典型场景核心风险JDBC直连分钟级单次百万行维度表补全、小规模关联并发连接打满CH、拉取慢离线批量导入T1单次千万~亿行夜间宽表加工、日报统计导入批次不当造成CH merge压力KafkaSpark Streaming秒~分钟级持续写入埋点日志清洗、实时指标流任务故障、延迟堆积湖仓CH加速层分钟~小时级大容量数据湖明细查询加速同步延迟导致数据不一致CH整体迁移一次性全量增量跨机房、跨版本迁移增量对齐、校验复杂3. Spark读写ClickHouse的落地细节与三个典型坑方案看着很多但真正动手写Spark和ClickHouse的连接代码时细节问题层出不穷。下面这部分是我在实际业务中踩过坑后总结出来的几个关键点建议照着做。3.1 版本与驱动注意包名变化这个大坑第一个坑来自ClickHouse JDBC驱动的包名。旧版本0.3.x及之前的驱动类名是ru.yandex.clickhouse.ClickHouseDriver很多老博客、老教程都在用这个。但新版本0.4.x之后官方把包名改成了com.clickhouse.jdbc.ClickHouseDriver。如果你引用的jar包升级了代码里还写着旧的类名会直接报ClassNotFoundException。我的建议是统一采用新版驱动。Maven坐标如下dependency groupIdcom.clickhouse/groupId artifactIdclickhouse-jdbc/artifactId version0.6.5/version classifierall/classifier /dependency注意classifierall/classifier很关键新版驱动如果不用这个分类器可能缺少一些HTTP客户端依赖。连接URL的基本格式是jdbc:clickhouse://ch-host:8123/default?connectTimeout30000socketTimeout600000ClickHouse的HTTP端口是8123如果你写9000端口那是native协议新版驱动也支持但建议默认走8123。3.2 读链路把过滤条件推给ClickHouse用Spark JDBC读ClickHouse时很多人直接写dbtable为表名然后让Spark全表扫描。如果你这么干大概率会遇到两个问题一是Spark会把整张表拉回来二是在ClickHouse端造成全表扫描。正确姿势是把过滤下沉到SQL里让ClickHouse充分利用稀疏索引。推荐写法是dbtable填一个子查询而不是直接填表名val df spark.read .format(jdbc) .option(url, jdbc:clickhouse://ch-host:8123/default) .option(dbtable, (SELECT event_date, uid, amount FROM app.orders WHERE event_date 2024-06-18) t) .option(user, readonly) .option(password, xxxxx) .option(driver, com.clickhouse.jdbc.ClickHouseDriver) .option(fetchSize, 50000) .load()原理是Spark JDBC DataSource支持谓词下推event_date 2024-06-18这类等值条件会被下推成SQL的WHERE子句由ClickHouse执行。但你要小心日期类型、Decimal类型的下推如果字段映射不对Spark无法把这个条件转成SQL字面量就会退化为全表扫描。所以最稳妥的办法是在子查询里就把分区条件写死别指望Spark的filter一定下推成功。另外读取ClickHouse时dbtable子查询里尽量只select需要的列。ClickHouse是列存少选一列IO就少一列这是基本素养。3.3 写链路别把JDBC当大批量导入工具这是我在生产环境踩过最深的坑之一。一开始我用Spark的write.format(jdbc)直接写ClickHouse结果每天跑完任务ClickHouse的system.parts表里冒出几万个partmerge线程告警查询性能从秒级变成分钟级。原因很简单MergeTree引擎每次INSERT都会生成一个新的data part后台线程再异步合并。如果Spark端一条一条insert或者批次太小就会产生海量小part。所以写ClickHouse有一条硬性准则宁可攒大包不要频繁小写。Spark JDBC写入虽然可以配置batchsize参数但底层对ClickHouse驱动而言跟原生的批量写入协议还是有差距。我的建议是超过千万级的数据量不要用write.format(jdbc)改用文件导入模式。3.4 文件导入模式Parquet clickhouse-client / s3表函数这是目前生产环境里我验证过最稳定的Spark写ClickHouse方案流程如下第一步Spark把结果写成Parquet列式文件到HDFS、本地临时目录或者S3对象存储。Parquet本身就是列存和ClickHouse的列式存储更匹配导入效率比CSV高很多还保留了数据类型信息。resultDF.write .mode(overwrite) .parquet(/tmp/spark-ch-export/20240618)第二步把Parquet导入ClickHouse。如果ClickHouse和Spark共用的同一个HDFS集群可以在ClickHouse节点上用clickhouse-client执行clickhouse-client --host ch-node1 --query INSERT INTO app.dws_result FORMAT Parquet /data/parquet/20240618/part-00000-xxx.snappy.parquet不过更现代化的做法是用ClickHouse的s3表函数让ClickHouse自己读取S3中的Parquet完成导入。Spark写S3ClickHouse从S3读取中间不用传文件。INSERT INTO app.dws_result SELECT * FROM s3(https://s3.amazonaws.com/bucket/path/20240618/*.parquet, AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY, Parquet)这种方式的好处是Spark和ClickHouse完全解耦S3作为中间缓冲导入速度由ClickHouse侧读取决定而且不占用Spark集群到ClickHouse的网络带宽。如果只是临时导入一次用clickhouse-client和文件形式最快如果每天定时批量导入走S3或者HDFS表函数更长久。第三步验证数据量和关键字段。导入完成后最好在Spark端和ClickHouse端对一下行数和SUM值避免数据重复或丢失。3.5 数据类型的映射坑Spark和ClickHouse的类型系统不完全一致字段映射经常出问题我简单列几个高频坑Spark的DecimalType映射到ClickHouse的Decimal64/Decimal128精度长度不一致时可能截断或报错。我建议导入前统一转成Decimal(38, 18)两边都留足空间。时间字段Spark的Timestamp映射到ClickHouse的DateTime或DateTime64。注意时区问题Spark默认UTCClickHouse默认服务器时区差8小时是常事。规范做法是统一采用DateTime64(3, UTC)展示层再做时区转换。字符串空值Parquet里的空字符串和null在ClickHouse里都变成空字符串除非字段设置为Nullable(String)。如果业务上要区分空字符串和NULL记得在导入前清洗一次避免语义错乱。ClickHouse的数组类型Array(Int32)在Spark里映射成ArrayTypeParquet支持但要注意嵌套层级不能太深太深会出现Spark和ClickHouse都不太能处理的状况。3.6 一个完整的Spark写ClickHouse示例以Spark 3.4 Scala 2.12 ClickHouse 24.x为例一个标准的离线导入任务长这样// 1. 读取Hive/数据湖 val sourceDF spark.sql( SELECT user_id, order_amount, order_time, channel FROM dwd_order_detail WHERE dt 2024-06-18 ) // 2. 在Spark里完成清洗加工 val resultDF sourceDF .filter(col(order_amount) 0) .withColumn(order_date, to_date(col(order_time))) .groupBy(user_id, order_date, channel) .agg( sum(order_amount).as(total_amount), count(*).as(order_cnt) ) // 3. 写出Parquet到对象存储 resultDF .repartition(10) .write .mode(overwrite) .parquet(s3a://data-lake/dws_daily_user/2024-06-18) // 4. 触发ClickHouse侧的s3表函数导入或者等待调度系统调用clickhouse-client这段代码看起来平平无奇但每一步都有讲究先过滤再聚合、写出前repartition控制文件数量、Parquet目录按日期分区。这些都是长期跑稳定任务的经验。4. 让两个系统和平相处的调优与保护策略Spark和ClickHouse一旦真正集成到同一套生产环境你还需要从治理层面保证它们不互相伤害。这个章节我会谈一些关键参数和实践心得。4.1 并发连接与连接数控制最典型的问题Spark读取ClickHouse时Spark可能生成几百个Task每个Task都会建立JDBC连接。如果并发打到ClickHouse上clickhouse-server的连接线程池会被打满其他业务查询直接卡死。这把火我点过。控制手段有两层。第一层Spark端限制读取并行度。Spark JDBC读取时分区数由partitionColumn、lowerBound、upperBound、numPartitions控制。对于ClickHouse这种OLAP库不需要那么多分区设置为8~16个就足够了并发太高只会徒增负担。.option(partitionColumn, user_id) .option(lowerBound, 1) .option(upperBound, 100000000) .option(numPartitions, 12)第二层ClickHouse端限制单用户并发查询数。给Spark用的账号设置一个较低的max_concurrent_queries_for_user或者通过chproxy这类网关做队列和限流防止突发流量压垮实例。4.2 批量写入与part管理写ClickHouse时批次的粒度直接决定merge的健康状况。ClickHouse官方建议单次INSERT至少1000行推荐1万行~100万行。如果一次INSERT数据量太大内存压力高太小part碎片多。Spark写入时控制好batchsize和写入并行度别让几十个Task同时往同一张表灌小批数据。另外每个MergeTree表建表时一定要设计好PARTITION BY和ORDER BY。比如按天分区用toYYYYMMDD(event_date)排序键选最常用的查询条件字段比如ORDER BY (event_date, user_id)。这样导入的每个part内部有序查询时能直接跳过大量数据也能减少后台merge压力。4.3 慢查询保护与内存限制ClickHouse一个让人爱恨交加的特性是查询太慢时它会死磕到底把所有CPU和内存吃光。这一点在Spark集成的场景里尤其危险——刚才提到Spark可能推来一个没下推过滤条件的查询ClickHouse会全表扫描然后整个查询把机器拖垮。所以生产环境必须给ClickHouse设置查询护栏。关键参数如下max_execution_time单位秒设置为300超过即中断查询防止SQL失控。max_memory_usage单查询最大内存默认是10GB你可以根据机器配置适当调低或保持。max_bytes_to_read单查询最多读取的数据量适合在BI账号上设置。readonly给Spark/BI的读账号设置readonly1禁止修改类的危险操作。我建议为Spark单独建一个低权限账号只允许SELECT并限制资源配额。这样哪怕Spark端误写一个全表扫描ClickHouse也能自保。4.4 任务调度与失败重试混合架构里Spark任务和ClickHouse导入任务是强依赖关系。Spark生成的数据文件没准备好就触发ClickHouse导入必然失败。所以调度策略要设计好。我们用的是Apache DolphinSchedulerSpark任务成功后再触发ClickHouse导入任务。每层任务都设置失败重试重试间隔2分钟、最多3次。还有一个容易被忽略的点导入任务完成之后必须加一个数据校验步骤对比Spark结果的行数和ClickHouse导入后的行数是否一致。别看它土很多线上数据事故都是栽在没有校验这一步上。5. 生产级混合大数据处理架构的长什么样讲了这么多最后落一个完整的架构图。这是我现在管理的线上平台正在跑的方案各模块都有明确职责基本能覆盖离线实时即席分析三类需求。5.1 数据分层与链路设计整体分四层ODS层原始数据落到HDFS/数据湖里用Hive/Iceberg管理保留全量、不做清洗成本最低、最安全DWD层Spark负责从ODS读取做清洗、去重、规范化、维度退化产出明细宽表依然存放在HDFS/数据湖DWS层Spark按业务主题做聚合产出汇总指标同时把需要高频查询的明细和聚合结果导入ClickHouseADS层BI报表、即席查询、用户画像服务直接读ClickHouse个别需要复杂关联的旁路任务继续走Spark。链路流转一般是两条离线链路ODS → Spark每日批处理 → DWD/DWS → Parquet → ClickHouse导入 → BI查询。实时链路业务数据/埋点 → Kafka → Spark Structured Streaming清洗 → 批量攒批 → ClickHouse导入 → 实时大屏/告警。5.2 各组件资源规划参考组件配置参考说明Spark集群4台 32C128G跑离线批处理和流任务动态资源池管理ClickHouse集群3分片1副本6节点16C64G按业务分库明细表用ReplicatedMergeTreeKafka3节点承接所有实时数据入口调度系统DolphinScheduler编排Spark批任务、CH导入、校验任务对象存储/HDFS10TB以上临时Parquet、全量数据存储资源上有一个建议Spark和ClickHouse不要混布在同一批物理机上。两者都是CPU和内存大户混在一起极易互相干扰。我们最早图省事混布过结果Spark跑大任务的时候ClickHouse查询P99直接从100毫秒涨到3秒后来拆开才好。5.3 性能收益与实际表现这套架构上线后我们做了几组对比效果是可以量化的。原来直接用Spark跑报表查询日活报表需要45秒现在Spark只负责计算和写入查询走ClickHouse响应时间在300毫秒以内。原来ClickHouse承担全部ETL加工时CPU经常90%以上现在ETL全部交给SparkClickHouse CPU稳定在20%以下查询并发承载能力明显提升。全量订单宽表约2.3亿行导入ClickHouse后单次查询按用户维度过滤加聚合响应时间稳定在400毫秒左右。另一个隐性收益是团队开发效率。数据开发只需要维护Spark的ETL逻辑和ClickHouse的建表语句责任边界清晰Spark管怎么算得对ClickHouse管怎么查得快出了问题排查范围也小。5.4 什么情况下其实不需要这套混合方案说句公道话混合大数据处理方案不是银弹。如果你的数据量还不到千万级查询压力也不大用ClickHouse单机版或者直接用MySQL索引就够用了引入Spark和ClickHouse两套系统纯粹是给自己找运维麻烦。另外如果你的查询主要是多层大表JOIN而且对灵活建模要求极高那ClickHouse未必合适。这种场景应该考虑Doris、StarRocks这类MPP分析型数据库它们对多表JOIN的支持比ClickHouse强不少而且本身可以融合导入链路未必需要Spark介入。反过来如果你的业务全是单表单聚类的报表用ClickHouse就足够了集成Spark反而是画蛇添足。我个人的体会是混合方案的真正价值在于计算与查询分离这个架构思想而不是非要把两个明星工具凑到一起。ClickHouse负责最后一公里的查询体验Spark负责前面几十公里的数据加工两者各司其职再用文件、消息队列、调度系统把它们衔接起来这就是目前我见过的最可靠、也最容易维护的大数据处理组合。如果你也在从纯离线数仓向实时化、服务化演进不妨沿着这条路线一步步搭起来。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门