一、数据存储格式Spark SQL支持多种文件存储格式主要分为行式存储和列式存储两大类。1. 行式存储数据按行组织同一行的所有字段值在物理上连续存放。以一张用户表为例ID, Name, Age, City行1| 1 | 张三 | 28 | 北京 | 行2| 2 | 李四 | 35 | 上海 | 行3| 3 | 王五 | 22 | 广州 |在磁盘上它会被存成[1,张三,28,北京][2,李四,35,上海][3,王五,22,广州]...主要特点写入非常快新增一行直接追加在末尾即可无需重组数据。整行读取效率高适合SELECT *或需要大部分字段的查询。压缩率较低一行内不同字段的数据类型、值分布差异大压缩算法难以发挥作用。列查询代价大如果只计算平均年龄SELECT AVG(Age)仍必须扫描所有行的全部字段I/O浪费严重。适用场景OLTP 在线事务处理频繁的增删改或者需要查询全部字段的情况。流式数据写入数据源逐条到达快速追加到文件如Kafka消息落盘成 Avro 文件。1.1CSV文本行式CSV是纯文本表格数据读写简单通用性极强但无内建类型和索引。字段用逗号/制表符等分隔不包含 Schema解析开销大压缩后不可拆分。适用于用户手工文件上传落表等情况。CREATE TABLE user_csv ( id INT, name STRING, age INT ) ROW FORMAT DELIMITED FIELDS TERMINATED BY , -- 指定按逗号分隔 STORED AS TEXTFILE -- 使用gzip压缩数据量不大的情况也可以删除该行不使用压缩 TBLPROPERTIES (compressiongzip);1.2 JSON半结构化文本行式JSON是一种轻量化的数据交换格式广泛应用于Web应用和分布式系统之间的数据交互。相比CSV来说支持复杂嵌套结构。Spark 2.2 可直接使用JSONFILECREATE TABLE user_json (id INT, name STRING, age INT) STORED AS JSONFILE TBLPROPERTIES (compressiongzip);或老版本用 Hive SerDeCREATE TABLE user_json (id INT, name STRING, age INT) ROW FORMAT SERDE org.apache.hive.hcatalog.data.JsonSerDe STORED AS TEXTFILE;1.3 Avro二进制行式Apache Avro 是一个基于 Schema 的、与语言无关的二进制数据序列化系统。它不像 JSON 那样是纯文本也不像 Parquet 那样是列式分析格式。Avro 的核心使命是让数据在分布式系统的不同组件之间高效、安全地流动同时允许数据结构随时间自由演化。适用于流处理、Schema 频繁演化的数据落地比如基于Kafka的数据管道。核心设计Schema 与数据分离Schema 用 JSON 定义描述数据的字段名、类型、默认值等人类可读。数据以紧凑二进制存储序列化后的记录不携带字段名和结构信息只是按 Schema 顺序紧密排列值。这带来了极小的体积和极快的解析速度。读取必须拥有 Schema写入时的 Schema写 Schema与读取时的 Schema读 Schema可以不同Avro 会自动按照字段名进行匹配和转换——这正是 Schema 演化的基石。Schema 演化向后与向前兼容双 Schema 解析Avro 数据文件中存储的是 Writer Schema 编码的二进制数据。读取时若提供不同的 Reader SchemaAvro 解析器会根据字段名映射和默认值规则自动转换数据而非直接报错。CREATE TABLE user_avro ( id INT, name STRING, age INT ) STORED AS AVRO TBLPROPERTIES (avro.compresssnappy);1.4SequenceFile二进制键值对行式Hadoop的键值对二进制格式常用于Hadoop 中间数据传递Spark SQL较少直接使用适用于MapReduce 中间结果、数据合并、小文件归档的情况。2. 列式存储数据按列组织同一列的所有值在物理上连续存放。同样用用户表举例磁盘布局会变成ID列 [1,2,3,...] Name列 [张三,李四,王五,...] Age列 [28,35,22,...] City列 [北京,上海,广州,...]主要特点分析查询极快列裁剪只读需要的列极大减少I/O。压缩率极高同一列数据类型一致值域往往相近非常适合游程编码Run-Length Encoding, RLE、字典编码Dictionary Coding如LZW等存储空间通常只有行式的1/3~1/10。谓词下推与向量化列存文件内部通常包含列块统计信息min/max等可以快速跳过不满足条件的整块数据向量化引擎能一次处理一批列值CPU效率高。写入较慢需要将行数据拆分成列并缓冲成列块后才能写入内存和计算开销较大。单行更新昂贵如果要修改一行中的某个字段可能需要重写整个列块不适合频繁随机更新。适用场景OLAP 在线分析查询大宽表上只对少数列做聚合、分组、过滤。数据仓库 / 数据湖海量历史数据存储追求高压缩比和扫描效率。2.1 ParquetApache Parquet 是一种开源的列式存储文件格式专为大数据处理和分析场景设计。一个Parquet文件的内容由Header、Data Block和Footer三部分组成。在文件的首尾各有一个内容为PAR1的Magic Number用于标识这个文件为Parquet文件。Header部分就是开头的Magic Number。Data Block是具体存放数据的区域由多个Row Group组成具体概念内容如下Row Group行组数据集的水平切片。Parquet 将整张表按行数水平划分为N个 Row Group。每个 Row Group 包含该组内所有列的数据是一个独立的并行处理单元。默认大小为128 MB与HDFS Block一致由于一个Map Task一般处理一个Block的数据这样设置可以增大任务执行并行度。Row Group 0: 包含第 1 ~ 200,000 行的所有列数据 Row Group 1: 包含第 200,001 ~ 400,000 行的所有列数据 ... Row Group 4: 包含第 800,001 ~ 1,000,000 行的所有列数据这样划分的好处每个 Row Group 可被不同的 Spark Task 并行读取Row Group 级别的统计信息Footer 中能让我们直接跳过整个不相关的 Row Group。Column Chunk列块一个 Row Group 内某一列的所有数据被组织成多个Page。Row Group 0 ├── Column Chunk for id (200,000 个整数) ├── Column Chunk for name (200,000 个字符串) ├── Column Chunk for age (200,000 个整数) └── Column Chunk for city (200,000 个字符串)每个 Column Chunk 在 Footer 中都会记录自己的元数据统计信息min_value,max_value,null_count等编码与压缩信息用了什么编码如 Dictionary RLE、压缩算法Snappy物理偏移量该 Column Chunk 的起始位置和大小这些统计正是谓词下推的依据。例如city列的 Column Chunk 统计显示min_value Baoding,max_value Shanghai如果查询条件WHERE city Zhengzhou这个 Row Group 的 city 值范围完全不包含 ‘Zhengzhou’整个 Row Group 就可以被跳过。Page页最小的 I/O 和编码单位。包含三种类型数据页 (Data Page)存储列的实际值经过编码和压缩。字典页 (Dictionary Page)如果该 Column Chunk 使用了字典编码会有一个字典页记录所有唯一值及其映射后面的数据页则只存储整数索引。索引页Page 级别的偏移量索引可选。Column Chunk: city (Row Group 0) ├── Dictionary Page: [0-Beijing, 1-Shanghai, ..., 19-Hangzhou] ├── Data Page 0: (行 1~20,000 的 city 索引) → [0,1,5,0,0,...] (RLE 压缩) ├── Data Page 1: (行 20,001~40,000 的索引) └── Data Page 2: (行 40,001~60,000 的索引) ...如果查询需要city BeijingParquet 读取器会先加载字典页将过滤条件转化为索引0然后在每个 Data Page 内查找是否包含0。Page 级别的索引Parquet Page Index2.0 版本还能记录每个 Page 内的min/max索引帮助直接跳过不包含索引0的 Page进一步提升效率。Footer部分用来存储整个文件的元数据由File Metadata、Footer Length和Magic Number三部分组成。FileMetaData记录文件元数据信息包括Schema和每个Row Group的Metadata。每个Row Group的Metadata又由各个Column的Metadata组成每个Column Metadata包含了其Encoding、Offset、Statistic信息等等。FileMetaData { version : int (格式版本如 1 或 2) schema : listSchemaElement (表结构的完整描述) num_rows : long (文件总行数) row_groups : listRowGroup (每个 Row Group 的元数据) key_value_metadata : optional listKeyValue (自定义属性) created_by : optional string (生成该文件的库与版本如 parquet-mr version 1.12.2) column_orders : optional listColumnOrder (列的排序规则) }Footer Length4 字节有符号整数用于标识Footer部分的大小帮助找到Footer的起始指针位置Magic Number再次 “PAR1”2.2 ORCOptimized Row Columnar优化行列式与 Parquet 相比ORC 的 Header/Footer 架构更显“重索引”——它把部分索引从中心元数据下放到 Stripe 内部牺牲了一点简单性换取了在 Hive 数仓场景下更锐利的点查询和范围过滤性能。Header只有 3 字节的Magic NumberORC作用与 Parquet 的PAR1完全相同格式识别和完整性校验。文件主体由多个Stripe组成每个 Stripe 都是独立的并行读取单元。每个 Stripe 由三部分组成Index Data、Row Data、Stripe Footer。Index Data轻量级行级索引存储该 Stripe 内列的统计信息如 min/max、布隆过滤器等用于快速过滤和跳过不必要的数据。这是 ORC 与 Parquet最大差异所在。Index Data 区为每列保存了 Row Index及可选 Bloom Filter Index。Bloom Filter Index如果建表时开启会为每个行组额外创建一个布隆过滤器快速判断一个值是否“肯定不存在”。例如WHERE city Zhengzhou布隆过滤器能瞬间回答“该行组不含 Zhengzhou”直接跳过。以city列在 Stripe 0 中的 200,000 行为例步长 10,000会生成 20 个索引条目每 10,000 行一组Row Index for city: entry 0 (行 1~10000): min Anshan, max Changsha, position offset_0 entry 1 (行 10001~20000): min Beijing, max Fuzhou, position offset_1 ... entry 19 (行 190001~200000): min Shanghai, max Shanghai, position offset_19Row Data列数据的流式存储实际存储该 Stripe 中所有行的列数据按列划分为多个 Stream流。Parquet先将整个文件数据水平切分为多个Row Group行组每个 Row Group 包含该组内所有行。在 Row Group 内部数据按列独立存储每一列的数据成为一个Column ChunkColumn Chunk 再切分成Page最小 I/O 单元。ORC同样先将文件水平切分为多个Stripe每个 Stripe 包含该组内所有行。在 Stripe 内部也是按列独立存储每列的数据由一组Stream构成如 DATA 流、PRESENT 流等Stream 就是 ORC 的最小物理单元。常见流类型PRESENT stream布尔位图标记该列哪些行是 NULL。DATA stream列的实际值根据类型和编码不同整数用 varint字符串用字典 ID 或直接字节等。LENGTH stream对于字符串类型记录每个值的字节长度。DICTIONARY_DATA stream如果使用字典编码存储字典内容。SECONDARY stream辅助信息如 Decimal 的小数部分。Stream: PRESENT (200,000 bit几乎全 1) Stream: DICTIONARY (包含 25 个不同城市名字符串) Stream: DATA (200,000 个字典索引整数用 varint 编码) Stream: LENGTH (空因为字典索引是定长 varint)Stripe Footer每个 Stripe 末尾有一个小 Footer包含该 Stripe 内每列的编码方式、流位置每个流的物理偏移与大小以及 Stripe 级别的统计信息。它与 File Footer 中的统计可能重复但这是为了在读取 Stripe 时无需跨文件查找。File Footer集中存放 Schema 和每个 Stripe 的列统计信息相当于 Parquet 的FileMetaData但 ORC 还在这里放布隆过滤器索引。Footer { headerLength : uint64 (Header 长度实际为 3) contentLength : uint64 (所有 Stripe 总字节数用于跳过数据) stripes : listStripeInformation (每个 Stripe 的物理位置) types : listType (完整的 Schema 树) metadata : listUserMetadataItem (自定义键值对) numberOfRows : uint64 (文件总行数) statistics : listColumnStatistics (每个 Stripe 每列的统计) rowIndexStride : uint32 (行索引步长默认 10000) }stripes每个StripeInformation记录了一个 Stripe 的offsetStripe 起始偏移indexLength索引区长度dataLength数据区长度footerLengthStripe Footer 长度numberOfRows该 Stripe 的行数statistics一个扁平的列表按 Stripe 顺序 × 列顺序排列。每个ColumnStatistics包含minValue、maxValue、hasNull、numberOfValues等。这就是谓词下推的依据。types记录完整 Schema支持 Struct、List、Map、Union 等复杂类型。rowIndexStride默认 10000表示 Stripe 内每 10000 行生成一组 Row Index 统计。Postscript保存 File Footer 的长度footerLength、压缩算法compression、Magic NumberORC等反序列化所需信息是读取 Footer 的钥匙。File Footer文件级统计记录整个文件每列的 min/max 等汇总信息用于在读取时直接跳过整个文件。Stripe FooterStripe 级统计记录该 Stripe 内每列的 min/max 等汇总信息用于跳过整个 Stripe。Index DataStripe 内部的行组级统计默认每 1 万行一组记录每组的 min/max 和布隆过滤器用于跳过 Stripe 内部分行组过滤粒度最细。2.3 ORC 与 Parquet 对比Parquet像一个仓库先把货物按批次Row Group堆放每批里再按品类Column Chunk装箱Page。仓库门后只有一张总清单Footer记录各批各箱的位置和统计。ORC也像仓库但每个批次Stripe内部还多放了一个小本子Index Data详细记录了每 10000 件货物的品类范围和位置方便快速翻找不用全箱打开。部分ParquetORCHeader4 字节PAR13 字节ORC元数据位置Footer 集中存放File Footer 集中存放但多一层 PostscriptPostscript无有存储 Footer 长度和压缩设置保证 Footer 可压缩Stripe/RowGroup 索引仅 ColumnChunk 级统计在 Footer 中无行内索引Footer 中有 Stripe 级统计Stripe 内部还有 10k 级 Row Index 和布隆过滤器编码与压缩每个 Page 独立编码压缩列流编码后整 Stripe 统一压缩SchemaFooter 中包含完整 SchemaFooter 中包含完整 Schema自定义属性key_value_metadatametadata列表Parquet 追求极简 Footer一次 Footer 读取即可规划全表扫描剩下全靠列存本身的 Page 跳过2.0 后才加入 Page Index。ORC 构建多层索引File FooterStripe 级→ Stripe 内 Index Data10k 行级→ 布隆过滤器过滤粒度更精细适合点查和复杂过滤但元数据层级稍多。二、压缩格式Snappy基于LZ77使用固定大小的哈希表快速查找重复字符串输出“长度-偏移量”对不做熵编码直接字节流输出速度极快但压缩比一般。Gzip采用DEFLATE算法先通过LZ77消除重复再对字面量、长度和距离分别进行Huffman 编码静态或动态属于经典的 LZ77 熵编码组合压缩比较高但速度较慢。Zstd同样基于LZ77变种但改用有限状态熵编码器FSE基于 ANS替代 Huffman支持更大的搜索窗口和更精细的建模通过多级序列化可在压缩比和速度之间灵活调整接近 Gzip 的压缩比却拥有 LZ4 级别的解压速度。LZ4极简的LZ77实现将匹配长度和字面量长度打包成一个令牌字节紧跟偏移量和字面量数据无熵编码逻辑非常简单以极低的 CPU 开销换取极高的压缩和解压速度。Bzip2完全不同于 LZ 系先对数据块进行Burrows-Wheeler 变换BWT让相似字符聚集再使用游程编码RLE压缩重复序列最后经过Huffman 编码输出。压缩比很高但 BWT 和排序导致压缩/解压极慢且内存消耗大。LZO同样是LZ77的优化实现使用哈希表匹配重复并对长度和偏移量进行紧凑编码支持重叠匹配以提升压缩比设计重点是解压速度极快比压缩快得多压缩时需要额外的块级处理是遗留 Hadoop 生态中的常用选择。压缩格式压缩比压缩速度解压速度原生可拆分典型适用场景Snappy中低极快极快否通用平衡之选列存内部默认Zstd高可调快极快可新版 Hadoop新一代平衡列存理想选择LZ4低极快极快否Shuffle 中间数据、低延迟场景Gzip / Zlib高较慢中等否归档、文本压缩、极致压缩比Bzip2很高极慢慢是极低存储成本的归档极少用LZO中等快快需索引遗留 Hadoop 生态需额外安装注Snappy/LZ4 为速度牺牲体积Zstd 通过调整等级可同时逼近高压缩比和高速度成为现代首选Gzip/Bzip2 严重偏向压缩比CPU 代价高。1. Avro内部结构文件由多个Data Block组成每个 Block 包含多条记录可独立压缩。支持的压缩Snappy、DeflateGzip、Zstandardzstd、LZ4 等。推荐Snappy写入/读取极快是流处理落地Kafka - Avro的默认选择。Zstd若需要更小体积且可接受轻微 CPU 增加是更好的替代。2. CSV / JSON (文本行式)自身无内部分块若需压缩只能对整个文件应用 Gzip、Bzip2 等。问题直接压缩为.gz将导致文件不可分割严重影响并行性。变通办法若要压缩并保持可并行读可使用Bzip2原生可分块或使用容器格式如 Avro二次存储。仅用于小规模数据交换或归档时可用Gzip。推荐尽量避免在生产中大规模使用压缩的 CSV/JSON 作为分析数据源。若要使用优先选择 Bzip2牺牲速度换可分割性或将数据先转为列式格式再压缩。3. SequenceFileHadoop 原生行式键值对格式支持块压缩NONE/RECORD/BLOCK 三种级别。常用压缩Snappy、Gzip、LZO 等BLOCK 级别压缩可保持可分割性。Spark SQL 中较少直接建表但若使用可通过spark.hadoop.mapred.output.compression.codec等参数控制。4. Parquet压缩粒度每个 Page 独立编码后再对整个 Page 应用通用压缩。支持的压缩Snappy默认、Gzip、Zstd、LZ4、None。推荐Snappy速度最快Spark 默认适合大多数分析场景。Zstd比 Snappy 高约 20-30% 的压缩率解压速度依然很快是新一代最佳平衡。如果存储成本敏感强烈推荐。Gzip牺牲读写速度换取更高压缩比适合长期冷数据归档。5. ORC压缩粒度Stripe 内部的列流Stream在 Stripe 级别统一压缩。支持的压缩Zlib默认、Snappy、LZO、LZ4、Zstd、None。特点ORC 自身编码整数 Varint、字符串字典等已经非常紧凑再结合压缩最终体积通常比 Parquet 略小。推荐Snappy平衡之选Spark 中常用覆盖默认的 Zlib。Zstd若引擎和 Hive 版本支持逐渐成为最佳实践兼顾速度和体积。ZlibHive 传统默认压缩比高但 CPU 消耗大Spark 下通常建议改为 Snappy。列存中同一列的数据类型相同、取值范围接近编码后连续相同值极多此时再施加 Snappy/Zstd 这类 LZ77 系算法可获得极高压缩比。而行式混合字段编码困难直接压缩效率低。因此列存与快速压缩Snappy/Zstd的组合既保证了查询性能又实现了优异的存储密度。三、补充Parquet 在SQL查询执行时的全流程假设我们要在 Spark SQL 中执行SELECT AVG(age) FROM users WHERE city Beijing;1. 读取 Footer获取全局元数据Spark 首先读取文件尾部 Footer得到Schema信息users表的全部字段信息以及类型所有 Row Group 的列表以及每个 Row Group 内每个 Column Chunk 的统计信息min/max/null2. 基于统计信息跳过 Row Group谓词下推针对city列Footer 显示每个 Row Group 的统计Row Group 0:citymin‘Baoding’, max‘Shanghai’ → 包含 ‘Beijing’需要读取Row Group 1:citymin‘Chengdu’, max‘Wuhan’ → 不包含 ‘Beijing’整个跳过Row Group 2: min‘Anshan’, max‘Beijing’ → 包含 ‘Beijing’需要读取Row Group 3: min‘Nanjing’, max‘Zhengzhou’ → 不包含跳过Row Group 4: min‘Beijing’, max‘Shenzhen’ → 包含需要读取最终只需读取 Row Group 0, 2, 4I/O 立即减少约 40%。3. 列裁剪只读需要的列在每个需要读取的 Row Group 内Spark 只读取涉及的两列city用于过滤age用于聚合id和name的 Column Chunk 完全不被访问再次大幅减少 I/O。4. 在 Row Group 内部通过 Page 精确读取以 Row Group 0 的city列为例先读取字典页将 ‘Beijing’ 转换成索引 0。扫描city的各个 Data Page可配合 Page Index 跳过不含索引 0 的页找出所有city Beijing的行号。根据这些行号到age列的 Column Chunk 中读取对应行的年龄值。因为age也是列式存储且同步分页Parquet 可以只读取包含这些目标行的agePage。5. 向量化计算读取出的age值以批处理向量化的方式送入 CPU计算平均值最终返回结果。整个过程Parquet 将全表扫描优化成了只读部分 Row Group 只读两列 只读满足过滤条件的少数 Page性能提升可达数十甚至上百倍。四、补充Parquet的Row Group、Column Chunk 的元数据信息RowGroup { total_byte_size : long (该 Row Group 总字节数) num_rows : long (该 Row Group 包含的行数) columns : listColumnChunk (每列的 Chunk 元数据) file_offset : optional long (Row Group 在文件中的起始偏移Parquet 2.0) total_compressed_size : optional long sorting_columns : optional listSortingColumn }ColumnChunk { file_path : string (如果文件是集合中的一个通常为 null) file_offset : long (该 ColumnChunk 在文件中的起始偏移量) meta_data { type : Type (列的数据类型) encodings : listEncoding (使用的编码如 PLAIN_DICTIONARY, RLE) path_in_schema: liststring (列在 Schema 树中的路径如 [users, name]) codec : CompressionCodec (压缩算法SNAPPY, GZIP 等) num_values : long (该 Chunk 中值的数量) total_uncompressed_size : long total_compressed_size : long key_value_metadata : optional listKeyValue /** 统计信息用于谓词下推 **/ statistics : Statistics { max : binary (最大值编码后的形式) min : binary (最小值) null_count: long distinct_count: long (可选唯一值大约数量) max_value : binary (Parquet 2.0 增强统计) min_value : binary } } offset_index_offset : optional long (Page 索引位置2.0) offset_index_length : optional long column_index_offset : optional long (列索引位置2.0) column_index_length : optional long }五、补充ORC 在SQL查询执行时的全流程假设我们要在 Spark SQL 中执行SELECT AVG(age) FROM users WHERE city Beijing;1. 文件打开读取 Postscript 与 File Footer从末尾读 1 字节 → Postscript 长度 → 读取 Postscript → 获取 Footer 长度 → 读取 File Footer。解析 Schema确认city和age列。获取statistics每个 Stripe 每列的 min/max。2. Stripe 级过滤利用 File Footer 统计根据 Footer 中的统计Stripe 0:citymin“Anshan”, max“Shanghai” → 包含 “Beijing”需读取。Stripe 1:citymin“Chengdu”, max“Zhengzhou” → 不包含 “Beijing”整个 Stripe 跳过。通过StripeInformation中的offset和dataLength直接跳过 Stripe 1 的物理区域。3. 列裁剪只读取 Stripe 0 的city和age列。在 Row Data 中仅解压并处理这两列对应的流忽略id和name。4. Stripe 内行级过滤Index Data在 Stripe 0 内读取city列的 Index Data20 个行组索引。对每个行组检查其min/max若“Beijing”不在范围内整个行组10,000 行跳过。假设只匹配到其中 4 个行组行 10001~20000行 50001~60000 等。如果 Bloom Filter 已开启ORC 还会用布隆过滤器快速确认“Beijing”是否可能存在进一步加速。5. 读取目标行组的数据流根据 Index Data 中匹配行组的position直接 seek 到city的 DATA stream 的对应区段解压缩读取字典 ID过滤出Beijing对应的行号。同时利用这些行号到age的 DATA stream 中读取相应位置的年龄值。6. 计算平均值向量化处理读取出的age值返回结果。整个过程ORC 用三级过滤Stripe → 10k 行组 → 布隆过滤器将扫描数据量压缩到极小。相比 Parquet 仅凭 Row Group 级统计ORC 在点查和谓词多时 I/O 更少但代价是索引构建和存储开销稍大。六、Spark Catalyst RBO / CBO 优化与列存文件格式的优化的区别Catalyst RBO 并不直接读取 Parquet/ORC 文件内部的统计信息如 min/max来做优化。列裁剪、谓词下推等优化是由 RBO 逻辑规则驱动的但规则本身只变换逻辑计划并不触碰文件格式。真正利用文件统计信息跳过数据的是物理执行阶段的数据源。1. RBO 阶段规则改写不碰文件统计Catalyst 的 RBO 是一组逻辑规则如PushDownPredicate、ColumnPruning它们基于关系代数等价变换完全不依赖数据文件的内容或统计信息。列裁剪Column Pruning规则做什么分析查询只需要哪些列将Project节点下推使得Scan节点仅输出需要的列。与格式的关系规则只是标记“只需列 A、C”至于怎么从文件中只读取这两列是后续数据源的事。谓词下推Predicate Pushdown规则做什么将Filter条件尽可能下推到靠近Scan的位置甚至下推到Scan内部如分区条件。与格式的关系规则只是把过滤条件表达式传给Scan节点。最终物理计划中FileSourceScanExec会把这些条件分为两类partitionFilters分区过滤dataFilters普通数据过滤此时并没有看文件里 min/max 跳过 Row Group 的统计那是在执行时发生的。2. 执行阶段数据源利用统计信息进行过滤当物理计划生成并开始执行时Parquet/ORC 的统计信息才真正发挥作用但这已不再是 RBO而是数据源本身的能力。FileSourceScanExec在构建扫描任务时会通过 Hadoop 的InputFormat或 Parquet/ORC 原生接口将dataFilters表达式序列化并传递给文件读取器。文件读取器做的事Parquet读取 Footer对每个 Row Group 检查其 Column Chunk 的min/max如果过滤条件不满足整个 Row Group 跳过。ORC读取 File Footer 和每个 Stripe 的 Index Data利用 Stripe 级统计和行级 Row Index及 Bloom Filter精细跳过不满足条件的行组。列裁剪同样在物理读取时只读取需要的列的 Column Chunk/Stream。这些操作是运行时数据源过滤Data Source Filter Pushdown不是 Catalyst RBO 规则。3. CBO 与统计信息的关系Catalyst 的CBO基于代价的优化会利用表的统计信息如行数、列的不同值数量来估算代价并选择更优的物理计划如 Join 顺序、Broadcast 阈值等。但这些统计信息不是直接从 Parquet/ORC 文件内嵌的 Row Group 统计获取的而是通过ANALYZE TABLE命令预先计算并存储在 Metastore 中的全局统计。文件内嵌的 min/max 主要用于执行时跳过数据不参与 CBO 代价估计。4.ANALYZE TABLE命令Hive 的ANALYZE TABLE命令用于收集表的统计信息供优化器做基于成本的优化CBO。存放于Hive Metastore元数据库。收集的统计信息包括表/分区级别行数、文件数、数据总大小列级别每列的 distinct 值、NULL 数、最大/最小值等常见用法-- 收集整个表的统计信息 ANALYZE TABLE my_table COMPUTE STATISTICS; -- 收集指定分区的统计信息 ANALYZE TABLE my_table PARTITION(dt2024-01-01) COMPUTE STATISTICS; -- 收集所有分区的统计信息 ANALYZE TABLE my_table PARTITION(dt) COMPUTE STATISTICS; -- 收集列级统计信息 ANALYZE TABLE my_table COMPUTE STATISTICS FOR COLUMNS col1, col2;作用优化器根据这些统计信息估算数据量、选择更优的 JOIN 策略如是否走 MapJoin、决定 reducer 数量等从而提升查询性能。对于 ORC/Parquet 这类列式存储Hive 的ANALYZE TABLE会优先利用文件内部的统计信息如行数、min/max、null 数等来快速计算表或分区的统计信息避免全表扫描。但并不完全基于文件内部统计信息因为有些统计如 distinct 值、列平均长度文件内部可能没有直接提供Hive 可能需要额外计算或使用近似算法如 HLL估算。不同点对比维度ORC/Parquet 文件内部统计信息HiveANALYZE TABLE统计信息存储位置数据文件内部Footer、Stripe/Row Group 索引Hive Metastore元数据库粒度非常细文件级、Stripe/Row Group 级、甚至每万行索引较粗表/分区级列级也较粗不细分到块内容主要是每列的min/max用于谓词下推过滤数据块行数、数据大小、distinct 值、null 数、平均长度等用于成本估算生成方式数据写入时由存储格式自动生成无需额外操作需要手动执行ANALYZE TABLE或开启自动收集使用时机查询执行时在读取数据前根据 min/max 跳过不满足条件的数据块查询规划时优化器CBO根据统计估算选择 Join 策略、Reducer 数等