尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Flink/Spark中IN查询为何慢?五大根因与优化方案

Flink/Spark中IN查询为何慢?五大根因与优化方案 你是不是也遇到过这种场景本地MySQL跑一个WHERE id IN (1,2,3...)毫秒级返回同样的SQL搬到Flink或者Spark上直接几分钟起步运气差点直接OOM调度平台上一排红色报警。很多大数据工程师第一反应是“框架不行”但实际真正的原因多数情况下是我们对分布式引擎的执行机制理解得不够透。这篇文章不打算讲那种网上满天飞的“IN改成JOIN就好啦”的片汤话而是老老实实把IN查询在Flink/Spark上为什么慢这件事拆开揉碎。从执行计划、数据分发、状态管理、数据倾斜几个层面讲清楚根因再给出一套我自己在多个项目里验证过的优化思路和实操步骤。无论你是刚接触实时数仓的新手还是已经被慢查询折磨了几轮的资深开发这篇文章都值得收藏下次再遇到这类问题可以直接按图索骥。1. 先从根子上说清楚IN查询在分布式引擎里到底干了什么1.1 你以为的IN查询和实际执行的IN查询不是一回事写SQL的人心里通常有个隐含假设IN查询就是“拿一个值去集合里比对一下命中就返回”。这个假设在单机数据库里基本成立因为有索引、有优化器、有统计信息。但在Flink和Spark这种分布式计算引擎里事情完全变了。举个最简单的例子你在Spark里写SELECT * FROM user_orders WHERE user_id IN (1001, 1002, 1003, ...)这张user_orders表如果存储在Hive/HDFS上数据是分散在几十个甚至上百个文件里的。Spark执行这个查询时会先启动一个全表扫描任务把每个文件都读一遍然后在内存里逐行判断user_id是否在IN列表里。整个过程中Spark根本没有“索引”可用它只能靠全量扫描加上逐行过滤来完成任务。这就暴露了第一个核心差异单机数据库有索引可以把随机IO变成少量精确IO分布式引擎没有全局索引只能全量IO加过滤。Flink SQL如果跑的是批模式Batch逻辑和Spark类似。如果跑的是流模式Streaming问题更复杂因为流式数据是无穷无尽的你要判断“当前这条数据的user_id是否在给定的列表里”这个列表本身也要作为状态保存下来还会牵扯到状态过期、状态恢复等一系列问题。1.2 分布式引擎的“分布式”三个字既是能力也是枷锁很多工程师忽略了一个关键点在分布式系统里数据分布决定了执行策略。一张表的数据被分成多个分区分散在不同节点上每个节点只能看到自己本地那部分数据。IN查询要想得到全局正确的结果必须经历这么几步每个节点扫描本地数据逐行判断是否命中IN列表。如果有join或者后续聚合还需要把数据按照某个key重新分区shuffle。把所有节点的结果汇总到Driver/JobManager再返回给用户。每一步都有成本全量扫描是IO成本shuffle是网络传输成本汇总阶段是单点压力成本。而IN查询在分布式引擎里恰恰没有针对“判断一个值是否在给定集合中”这个动作做任何特殊优化。我之前看过一个生产案例Spark SQL跑一个IN查询IN列表里其实只有50个ID但目标表有几十亿行。整个任务跑了40分钟其中绝大多数时间都耗在了全表扫描上。后来改成先用过滤条件把表裁剪到几百万行再执行IN查询整个任务跑完不到2分钟。这中间差的就是对分布式引擎执行机制的理解。2. IN查询慢到离谱的五个常见根因2.1 数据倾斜一个分区拖垮整个作业数据倾斜在分布式计算里就像木桶效应——整个作业的耗时取决于最慢的那个任务。而IN查询引发的数据倾斜往往藏得很深。举个例子你查WHERE province IN (北京, 上海, 新疆)看似是三个值但如果你的数据是按省份分区的而北京和上海的数据量占了整个表的80%这两个分区的任务就要处理海量数据其他分区的任务可能几秒就干完了然后整个作业就卡在最慢的那两个任务上。更隐蔽的是shuffle阶段的数据倾斜。如果IN查询后面还跟着GROUP BY province或者其他聚合操作按省分组时北京、上海这两个key的数据量天然就大会导致某个Reduce任务处理的数据量是其他任务的几十倍直接把这个任务拖垮。怎么判断是不是倾斜打开Spark UI或者Flink Web UI看各Task的耗时分布和Shuffle Read/Write量。如果出现“绝大部分Task秒级完成少数几个Task十几分钟都跑不完”基本可以断定是倾斜。2.2 谓词下推失效数据到内存之后才想起来过滤谓词下推Predicate Pushdown是优化器的基本功理论上应该把WHERE条件尽量下推到数据源让引擎只读取满足条件的部分数据。但IN查询在实际执行时谓词下推经常失效。主要原因有两个文件格式不支持如果底层数据是纯文本或者Parquet但没有开启统计信息引擎没法提前判断文件里是否包含满足条件的数据只能全量读进来再过滤。优化器保守策略当IN列表过长或者查询过于复杂时某些优化器会放弃下推直接全量扫描。这在Spark的CBO基于成本的优化还没完全生效的老版本里特别常见。我在实际项目中踩过一个坑有一张以order_date做分区字段的表SQL里同时写了WHERE order_date 2024-06-01 AND user_id IN (...)。理论上Spark应该只读取6月1日这一个分区的数据但实际执行计划里显示扫描了全表。后来排查发现是因为表没有做分区发现Partition DiscoverySpark根本没识别出分区字段导致IN查询变成了全表扫描。2.3 流式处理中的状态膨胀与重放问题FlinkSQL如果跑在流模式IN查询的语义需要重新审视。流式数据是无限的你要判断一条数据是否满足某个条件就需要一个状态存储来记住“历史数据”或者“维表数据”。如果这个状态无限增长Flink的State Backend会承受巨大压力。假设你要做实时订单和用户维表的关联SQL写成SELECT * FROM orders WHERE user_id IN (SELECT user_id FROM users WHERE vip_level 1)在流模式下Flink要维护两张表的状态。users表的每次更新都会触发状态变更orders表的数据每进来一条都要去状态里查询。如果users表的数据量几十万且更新频繁状态会非常大RocksDB的读写压力也会飙升整个作业的吞吐量直接下降一个数量级。另外流式任务如果发生故障需要从Checkpoint恢复。Checkpoint里包含了状态数据状态越大恢复时间越长这也是很多生产作业“越跑越慢”的原因之一——不是查询本身变慢了而是状态管理带来的开销在持续累积。2.4 广播变量缺失每个任务都在重复拉取数据IN查询本质上是在“判断一个值是否属于一个集合”。如果这个集合很小比如几千个ID最合理的做法是把这个集合广播到每个计算节点上让每个Task在本地内存里直接判断。但在很多默认配置下Spark和Flink并不会自动做这件事。Spark的spark.sql.autoBroadcastJoinThreshold默认是10MB只有当IN对应的子查询结果小于这个阈值且优化器识别出可以转成Broadcast Hash Join时才会走广播路径。一旦超过阈值优化器会退化成Sort Merge Join或者直接走嵌套循环性能急剧下降。更隐蔽的一个问题是很多人在Spark里写WHERE id IN (SELECT id FROM dim_table)虽然子查询的维度表很小但如果dim_table没有统计信息优化器不知道它有多小干脆就不走了广播路径导致每个Reduce任务都要重新读取或者shuffle一份维度表数据网络IO和CPU全部被打满。2.5 小文件问题扫描本身成了最大成本这一点在Hive/Spark场景特别突出。很多团队的数据管道没有做好小文件合并Compact一个几GB的表被拆成了几千个甚至几万个小于1MB的小文件。Spark读取文件时每个文件至少启动一个Task来读取小文件越多Task数量越多每个Task都有固定的调度开销、JVM启动开销和元数据拉取开销。当IN查询触发全表扫描时这些开销就被放大了。我遇到过最离谱的一个案例一个数据量只有2GB的表因为小文件太多Spark启动了4000多个Task实际计算时间不到30秒但整个作业跑了12分钟绝大多数时间都耗在Task调度和文件打开关闭上。这种情况下不解决小文件问题优化SQL写法根本治标不治本。3. 一套可以直接抄的优化方案3.1 能用Join就别用IN让优化器去干活这不是一句空洞的口号而是有底层逻辑的。在分布式引擎里Join的优化器支持程度和运行时执行策略远比IN查询要成熟。Spark对Join有Brodcast Hash Join、Sort Merge Join、Shuffle Hash Join等多种策略Flink也支持多种Join优化。你把IN改成Join实际上是把决定权交给了优化器。以Spark为例下面这两段SQL是等价的-- 写法一IN SELECT * FROM orders WHERE user_id IN (SELECT user_id FROM vip_users) -- 写法二LEFT SEMI JOIN SELECT o.* FROM orders o LEFT SEMI JOIN vip_users v ON o.user_id v.user_id改成SEMI JOIN之后优化器至少有三个额外手段可用如果vip_users足够小自动走Broadcast Hash Join每个Task本地判断。如果足够大走Sort Merge Join通过排序加合并的方式做匹配避免了IN查询那种“逐行去查集合”的低效模式。可以结合spark.sql.autoBroadcastJoinThreshold参数来控制广播阈值让优化器更智能地选择执行策略。但要注意LEFT SEMI JOIN和IN还有一个关键区别——IN会自动去重而LEFT SEMI JOIN不会。如果子查询结果里有重复的user_id可能影响最终结果。我建议在子查询里显式加上DISTINCT语义更清晰也便于优化器估算数据量。3.2 广播小维表让每个Task在本地判断如果你的IN列表本质上是对应一个“小维表”的查询那最好的方式就是让这个维表广播到所有节点。在Spark里可以这样操作// 广播维度数据 val vipUserIds spark.sql(SELECT DISTINCT user_id FROM vip_users) .collect() .map(_.getLong(0)) .toSet // 广播到所有Executor val bcVipUserIds spark.sparkContext.broadcast(vipUserIds) // 在主表上过滤 import spark.implicits._ val result userOrders .filter(row bcVipUserIds.value.contains(row.getAs[Long](user_id)))这个方案的本质是把集合判断从分布式计算变成了本地内存判断。数据在读取出来之后内存里一个contains操作就完成了判断完全避免了shuffle。Flink里头也可以做类似的事情用BroadcastStream把维表广播到所有算子实例然后通过KeyedBroadcastProcessFunction或者BroadcastProcessFunction实现本地判断。3.3 Bloom Filter用极小的代价过滤掉99%的无效数据Bloom Filter是一个很老但很实用的数据结构。它的核心思想是用多个哈希函数把集合映射到一个bit数组上判断一个元素“是否不存在”非常准如果Bloom Filter说“不存在”那一定不存在如果它说“可能存在”才需要去真正验证。在IN查询场景你可以在主表扫描前加一道Bloom Filter过滤把那些绝对不可能命中的数据直接过滤掉只有可能命中的数据才进入完整的判断流程。Spark里通常这样用import org.apache.spark.util.sketch.BloomFilter // 从IN列表构建Bloom Filter val expectedNumItems 10000 val fpp 0.01 // 1%的误判率 val bf BloomFilter.create(expectedNumItems, fpp) inList.foreach(id bf.putLong(id)) // 数据过滤 val filtered userOrders .filter(row bf.mightContainLong(row.getAs[Long](user_id)))尤其在大表上如果主表有几十亿行而IN列表只有几百个ID用Bloom Filter可以把超过95%的数据在扫描阶段就过滤掉极大减少后续shuffle和计算的负载。这里有一个我踩过的坑Bloom Filter的误判率参数不要设成0否则会占用大量内存性能反而下降。一般设置在0.01到0.05之间就足够了多filter出来的一点点数据后续判断成本很低。3.4 拆分IN列表避免单个任务计算量过大当IN列表特别长比如几万个ID时优化器没准会把它拆成多个子任务。但更多情况下过长的IN列表会导致两个问题如果走的是逐行判断CPU的哈希计算压力会很大。如果走的是Join那么一个巨大的集合参与Joinshuffle的数据量会成倍增加。我的实操建议是把IN列表拆成多个批次每个批次几百或几千个ID分批查询后再合并结果。每一批的数据都远小于阈值优化器更容易选择Broadcast策略。// 伪代码分批查询 ListListLong batches partitionList(inList, 500); ListRow result new ArrayList(); for (ListLong batch : batches) { String inStr batch.stream() .map(String::valueOf) .collect(Collectors.joining(,)); DatasetRow batchResult spark.sql( SELECT * FROM user_orders WHERE user_id IN ( inStr ) ); result.addAll(batchResult.collectAsList()); }这个方案在JDBC数据源上效果特别明显因为数据库本身有索引小批量的IN查询可以利用索引快速返回结果。3.5 参数调优让引擎自带的优化机制生效很多IN查询慢不是SQL本身的问题而是平台的默认参数配置没有跟上数据规模。几个最常被忽视的参数引擎参数作用建议值Sparkspark.sql.autoBroadcastJoinThreshold广播阈值根据维表大小调整为50MB~200MBSparkspark.sql.shuffle.partitionsShuffle分区数根据数据量调整为200~2000Sparkspark.sql.adaptive.enabled动态执行开启Flinktable.exec.state.ttl状态TTL根据业务需求设置如1小时Flinktable.optimizer.join-reorder-enabledJoin重排开启Flinktable.exec.shuffle-adaptive.enabled动态Shuffle开启这里特别点名spark.sql.adaptive.enabled也就是Spark 3.0引入的动态分区裁剪和动态调整Shuffle分区数功能。开启之后Spark可以根据运行时的shuffle数据量自动调整Reduce任务数量还能自动处理数据倾斜。我把这个参数打开之后多个任务直接提速了3倍以上。Flink SQL如果用的是流模式table.exec.state.ttl一定要设置合理值。如果状态不设置TTL会无限增长最终压垮整个作业。设置为业务需要的窗口大小或者维表更新周期即可。4. 实操过程与核心环节实现4.1 Spark SQL案例从IN到Join的改造实录说一个完整案例。有张订单表orders存储在Hive按天分区有张用户维表usersMySQL里的表通过JDBC读取。业务需求是查出所有VIP用户的订单。原始SQL长这样SELECT * FROM orders WHERE order_date 2024-06-01 AND user_id IN ( SELECT user_id FROM users WHERE vip_level 1 )这个SQL在Spark跑出来的执行计划是先全量读取6月1日的订单分区同时全量读取users表然后做Sort Merge Join。订单表当天有3000万行users表有50万行任务跑完用了14分钟Shuffle数据量接近5GB。优化后的做法分三步第一步先把用户在JDBC端过滤掉广播到Executor上val vipUserIds spark.read .jdbc(mysqlUrl, users, props) .filter(vip_level 1) .select(user_id) .distinct() .collect() .map(_.getLong(0)) .toSet val bcVip spark.sparkContext.broadcast(vipUserIds)第二步通过广播变量过滤订单表spark.sql(SELECT * FROM orders WHERE order_date 2024-06-01) .filter(row bcVip.value.contains(row.getAs[Long](user_id)))第三步加一个Bloom Filter做预过滤进一步减少判断次数val bf BloomFilter.create(1000000, 0.02) vipUserIds.foreach(bf.putLong) spark.sql(SELECT * FROM orders WHERE order_date 2024-06-01) .filter(row bf.mightContainLong(row.getAs[Long](user_id))) .filter(row bcVip.value.contains(row.getAs[Long](user_id)))这个方案跑下来时间从14分钟降到了1分40秒左右。核心收益来自于全表数据从逐行比对集合变成了先在本地做Bloom Filter粗筛再在本地做Set精确判断完全没有触发Shuffle。4.2 Flink SQL案例用Lookup Join替换流式INFlink流式场景下IN查询的常用替代方案是Lookup Join。假设你要实时处理订单流需要关联MySQL里的用户维度判断用户是否为VIP再输出。不推荐的写法是这样的——直接在每一条数据上去查MySQLSELECT o.order_id, o.user_id, o.amount FROM orders o WHERE o.user_id IN ( SELECT user_id FROM users WHERE vip_level 1 )这个写法在流模式里需要维护一个巨大的状态来同步MySQL里的用户表性能极差。推荐的做法是用Flink的CREATE TABLE声明维表然后走Lookup JoinCREATE TABLE users ( user_id BIGINT, vip_level INT, PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/test, table-name users, lookup.cache.max-rows 500000, lookup.cache.ttl 3600s, lookup.max-retries 3 ); CREATE TABLE orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECONDS ) WITH ( connector kafka, topic orders, properties.bootstrap.servers localhost:9092, format json ); SELECT o.order_id, o.user_id, o.amount FROM orders o JOIN users FOR SYSTEM_TIME AS OF o.ts u ON o.user_id u.user_id WHERE u.vip_level 1关键点在于FOR SYSTEM_TIME AS OF语法它表示“按照订单数据的时间戳去查询当时的维表快照”。通过设置lookup.cache.max-rows和lookup.cache.ttl让Flink在本地维护一个维表缓存查询不需要每次都走网络到MySQL。这样既避免了全量状态的膨胀又大幅降低了维表查询的延迟。我上线过类似的实时计算任务用这个方案把维表查询的QPS从每秒几百提升到了每秒几万而且完全没有给MySQL造成压力。4.3 验证与调优怎么判断优化有没有生效优化做完别急着庆祝。要用数据验证执行计划是否达到了预期效果。看执行计划Spark里用EXPLAIN SELECT o.* FROM orders o LEFT SEMI JOIN vip_users v ON o.user_id v.user_id重点关注Physical Plan里有没有出现BroadcastHashJoin或者BroadcastExchange。如果出现的还是SortMergeJoin说明广播阈值没生效需要调整参数或者强制广播。看运行时指标Spark UI里SQL选项卡看每个Stage的Shuffle Read/Write大小有没有降下来。任务耗时分布如果最长和最短任务差距在3倍以上可能有数据倾斜。Executor的GC时间如果GC时间占比过高说明Executor内存不足考虑增加内存或者减少并发数。Flink侧在Web UI里观察每个算子实例的numRecordsIn/Out看看数据是否倾斜。State的大小和RocksDB的读写延迟。背压水位BackPressure如果某些算子一直是高水位说明下游处理能力不足。另外我强烈建议开启慢查询日志直接在平台层面捕获那些执行时间超过阈值的SQL。有了基线数据才能准确评估优化前后的提升而不是凭感觉“好像快了一点”。5. 常见问题与排查技巧实录5.1 “timer执行查询报空指针”这类异常怎么排查我在Flink任务里不止一次遇到类似“timer执行查询是报空指针”的异常。这类问题的典型场景是Flink SQL里用了TUMBLE窗口或者INTERVAL JOIN内部会注册Timer。当Timer触发时它要去查关联的状态或者外部的维表但此时对应的状态已经被清理或者KeyBy的上下文信息缺失就会抛空指针。排查思路三个方向先看完整的异常堆栈定位是哪个算子抛的空指针。检查table.exec.state.ttl是不是设置得太短窗口还没结束状态就过期被清了。检查KeyBy逻辑确保同一个Key的数据进入同一个算子实例否则状态和Timer不匹配。如果是维表查询导致的空指针大概率是序列化问题。检查维表连接器的返回类型是否和SQL声明一致特别是DECIMAL、BIGINT这类容易在反序列化时出错的类型。5.2 IN列表太长导致SQL解析失败有些业务场景需要传入上万个ID直接把SQL拼出来可能超过解析器限制。Spark跑的时候可能报Cannot parse the input SQLFlink也有可能报语法错误。经验上推荐两种方式方案一把ID列表写入临时表然后用子查询方式。-- 先创建一个临时表 CREATE TEMP VIEW in_list AS SELECT explode(array(1001, 1002, 1003, ...)) AS user_id; -- 再用IN子查询 SELECT * FROM orders WHERE user_id IN (SELECT user_id FROM in_list)方案二用Hive的VALUES子句。SELECT * FROM orders WHERE user_id IN (VALUES 1001, 1002, 1003, ...)这两种方式既能绕开语法限制更重要的是让优化器拿到一个真实的“数据集”来估算大小更容易触发广播优化。5.3 优化后反而更慢了先检查数据分布有时候改成Join后性能不升反降。我遇到过两个典型场景。第一个是广播变量太大。如果维表真实数据超过广播阈值直接把几GB的数据广播到每个Executor上反而导致Executor内存溢出GC频繁。解决方法是重新评估阈值或者改用Sort Merge Join。第二个是子查询去重后数据膨胀。LEFT SEMI JOIN如果子查询不加DISTINCT可能会因为重复key导致shuffle数据量翻倍。我当时在一个用户标签场景里发现子查询里的user_id有大量重复join之后shuffle数据直接翻了三倍加上DISTINCT之后才恢复正常。所以在优化之后一定要对比执行计划中的shuffle数据量和任务耗时不要只看最终结果的对不对。6. 写在最后的经验从我个人的实操经验来看处理大数据引擎上的IN查询性能问题最忌一上来就改代码。先定位瓶颈在哪个阶段是扫描、shuffle、状态还是任务调度再去选对应的优化手段才能一次性打准。如果只是偶尔跑一次的低频SQL用3.1里面那种改Join的方式就够了简单直接。如果是高频调度任务那就要认真考虑广播变量加Bloom Filter的组合方案把每个环节的计算量都压到最低。还有一个小技巧在Spark SQL里尽量给表加上统计信息也就是执行一下ANALYZE TABLE。有了准确的统计信息优化器才能做出更合理的执行策略选择尤其是在判断“小表”和“广播阈值”的时候效果立竿见影。我自己在多个项目里实践下来仅仅补上统计信息这一步就解决了不少IN查询莫名走得极慢的问题。希望这篇内容能帮你少踩几个坑。下次再遇到IN查询在大数据引擎上慢到离谱先别急着骂框架试着从执行计划和数据分布入手大概率能找到真正的症结。
返回列表