ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Spark SQL中数据存储格式与压缩格式

Spark SQL中数据存储格式与压缩格式

一、数据存储格式

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 ('compression'='gzip');
1.2 JSON半结构化文本行式

JSON是一种轻量化的数据交换格式,广泛应用于Web应用和分布式系统之间的数据交互。相比CSV来说,支持复杂嵌套结构。

Spark 2.2+ 可直接使用JSONFILE

CREATE TABLE user_json (id INT, name STRING, age INT) STORED AS JSONFILE TBLPROPERTIES ('compression'='gzip');

或老版本用 Hive SerDe:

CREATE 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 Schema,Avro 解析器会根据字段名映射和默认值规则自动转换数据,而非直接报错。

CREATE TABLE user_avro ( id INT, name STRING, age INT ) STORED AS AVRO TBLPROPERTIES ('avro.compress'='snappy');
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 Parquet

Apache 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 = 'Beijing',Parquet 读取器会先加载字典页,将过滤条件转化为索引0,然后在每个 Data Page 内查找是否包含0。Page 级别的索引(Parquet Page Index,2.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 : list<SchemaElement> (表结构的完整描述) num_rows : long (文件总行数) row_groups : list<RowGroup> (每个 Row Group 的元数据) key_value_metadata : optional list<KeyValue> (自定义属性) created_by : optional string (生成该文件的库与版本,如 "parquet-mr version 1.12.2") column_orders : optional list<ColumnOrder> (列的排序规则) }
  • Footer Length:4 字节有符号整数,用于标识Footer部分的大小,帮助找到Footer的起始指针位置

  • Magic Number:再次 “PAR1”

2.2 ORC(Optimized Row Columnar)优化行列式

与 Parquet 相比,ORC 的 Header/Footer 架构更显“重索引”——它把部分索引从中心元数据下放到 Stripe 内部,牺牲了一点简单性,换取了在 Hive 数仓场景下更锐利的点查询和范围过滤性能。

Header:只有 3 字节的Magic NumberORC,作用与 Parquet 的PAR1完全相同:格式识别和完整性校验。

文件主体:由多个Stripe组成,每个 Stripe 都是独立的并行读取单元。每个 Stripe 由三部分组成:Index DataRow DataStripe 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_19
  • Row Data(列数据的流式存储):实际存储该 Stripe 中所有行的列数据,按列划分为多个 Stream(流)。

  • Parquet:先将整个文件数据水平切分为多个Row Group(行组),每个 Row Group 包含该组内所有行。在 Row Group 内部,数据按列独立存储,每一列的数据成为一个Column Chunk,Column 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 : list<StripeInformation> (每个 Stripe 的物理位置) types : list<Type> (完整的 Schema 树) metadata : list<UserMetadataItem> (自定义键值对) numberOfRows : uint64 (文件总行数) statistics : list<ColumnStatistics> (每个 Stripe 每列的统计) rowIndexStride : uint32 (行索引步长,默认 10000) }
  • stripes:每个StripeInformation记录了一个 Stripe 的:

    • offset:Stripe 起始偏移

    • indexLength:索引区长度

    • dataLength:数据区长度

    • footerLength:Stripe Footer 长度

    • numberOfRows:该 Stripe 的行数

  • statistics:一个扁平的列表,按 Stripe 顺序 × 列顺序排列。每个ColumnStatistics包含minValuemaxValuehasNullnumberOfValues等。这就是谓词下推的依据

  • 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 件货物的品类范围和位置,方便快速翻找,不用全箱打开。

部分ParquetORC
Header4 字节PAR13 字节ORC
元数据位置Footer 集中存放File Footer 集中存放,但多一层 Postscript
Postscript有,存储 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 Footer(Stripe 级)→ Stripe 内 Index Data(10k 行级)→ 布隆过滤器,过滤粒度更精细,适合点查和复杂过滤,但元数据层级稍多。

二、压缩格式

  • 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、Deflate(Gzip)、Zstandard(zstd)、LZ4 等。

  • 推荐

    • Snappy:写入/读取极快,是流处理落地(Kafka -> Avro)的默认选择。

    • Zstd:若需要更小体积且可接受轻微 CPU 增加,是更好的替代。

2. CSV / JSON (文本行式)

  • 自身无内部分块,若需压缩,只能对整个文件应用 Gzip、Bzip2 等。

  • 问题:直接压缩为.gz将导致文件不可分割,严重影响并行性。

  • 变通办法

    • 若要压缩并保持可并行读,可使用Bzip2(原生可分块)或使用容器格式(如 Avro)二次存储。

    • 仅用于小规模数据交换或归档时,可用Gzip

  • 推荐尽量避免在生产中大规模使用压缩的 CSV/JSON 作为分析数据源。若要使用,优先选择 Bzip2(牺牲速度换可分割性),或将数据先转为列式格式再压缩。

3. SequenceFile

  • Hadoop 原生行式键值对格式,支持块压缩(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 版本支持,逐渐成为最佳实践,兼顾速度和体积。

    • Zlib:Hive 传统默认,压缩比高但 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/null)

2. 基于统计信息跳过 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, 4,I/O 立即减少约 40%。

3. 列裁剪:只读需要的列

在每个需要读取的 Row Group 内,Spark 只读取涉及的两列:

  • city(用于过滤)

  • age(用于聚合)

idname的 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 : list<ColumnChunk> (每列的 Chunk 元数据) file_offset : optional long (Row Group 在文件中的起始偏移,Parquet 2.0+) total_compressed_size : optional long sorting_columns : optional list<SortingColumn> }
ColumnChunk { file_path : string (如果文件是集合中的一个,通常为 null) file_offset : long (该 ColumnChunk 在文件中的起始偏移量) meta_data { type : Type (列的数据类型) encodings : list<Encoding> (使用的编码,如 PLAIN_DICTIONARY, RLE) path_in_schema: list<string> (列在 Schema 树中的路径,如 ["users", "name"]) codec : CompressionCodec (压缩算法,SNAPPY, GZIP 等) num_values : long (该 Chunk 中值的数量) total_uncompressed_size : long total_compressed_size : long key_value_metadata : optional list<KeyValue> /** 统计信息,用于谓词下推 **/ 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,确认cityage列。

  • 获取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中的offsetdataLength直接跳过 Stripe 1 的物理区域。

3. 列裁剪

只读取 Stripe 0 的cityage列。在 Row Data 中,仅解压并处理这两列对应的流,忽略idname

4. Stripe 内行级过滤(Index Data)

在 Stripe 0 内,读取city列的 Index Data(20 个行组索引)。

  • 对每个行组检查其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 是一组逻辑规则(如PushDownPredicateColumnPruning),它们基于关系代数等价变换,完全不依赖数据文件的内容或统计信息

列裁剪(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(dt='2024-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 数等
返回列表