0%

orcfile

https://cwiki.apache.org/confluence/display/Hive/LanguageManual+ORC

1 文件格式

ORC的全称是(Optimized Row Columnar),它并不是一个单纯的列式存储格式,仍然是首先根据行组分割整个表,在每一个行组内进行按列存储。ORC文件是自描述的,它的元数据使用Protocol Buffers序列化,并且文件中的数据尽可能的压缩以降低存储空间的消耗。

1.1 列存储的优势

  • 查询的时候不需要扫描全部的数据,而只需要读取每次查询涉及的列,这样可以将I/O消耗降低N倍,另外可以保存每一列的统计信息(min、max、sum等),实现部分的谓词下推。
  • 由于每一列的成员都是同构的,可以针对不同的数据类型使用更高效的数据压缩算法,进一步减小I/O。
  • 由于每一列的成员的同构性,可以使用更加适合CPU pipeline的编码方式,减小CPU的缓存失效。

1.2 文件结构

ORC文件也是以二进制方式存储的,ORC文件也是自解析的,它包含许多的元数据,这些元数据都是同构ProtoBuffer进行序列化的。

  • ORC文件:保存在文件系统上的普通二进制文件,一个ORC文件中可以包含多个stripe,每一个stripe包含多条记录,这些记录按照列进行独立存储,对应到Parquet中的row group的概念。
  • 文件级元数据:包括文件的描述信息PostScript、文件meta信息(包括整个文件的统计信息)、所有stripe的信息和文件schema信息。
  • stripe:一组行形成一个stripe,每次读取文件是以行组为单位的,一般为HDFS的块大小,保存了每一列的索引和数据。
  • stripe元数据:保存stripe的位置、每一个列的在该stripe的统计信息以及所有的stream类型和位置。
  • row group:索引的最小单位,一个stripe中包含多个row group,默认为10000个值组成。
  • stream:一个stream表示文件中一段有效的数据,包括索引和数据两类。索引stream保存每一个row group的位置和统计信息,数据stream包括多种类型的数据,具体需要哪几种是由该列类型和编码方式决定。

1.3 基于统计信息的过滤

三个层级的统计信息,分别为文件级别、stripe级别和row group级别的,都可以用来根据Search ARGuments(谓词下推条件)判断是否可以跳过某些数据,在统计信息中都包含成员数和是否有null值,并且对于不同类型的数据设置一些特定的统计信息。

(1)file level 在ORC文件的末尾会记录文件级别的统计信息,会记录整个文件中columns的统计信息。这些信息主要用于查询的优化,也可以为一些简单的聚合查询比如max, min, sum输出结果。 
(2)stripe level ORC文件会保存每个字段stripe级别的统计信息,ORC reader使用这些统计信息来确定对于一个查询语句来说,需要读入哪些stripe中的记录。比如说某个stripe的字段max(a)=10,min(a)=3,那么当where条件为a >10或者a <3时,那么这个stripe中的所有记录在查询语句执行时不会被读入。 
(3)row level 为了进一步的避免读入不必要的数据,在逻辑上将一个column的index以一个给定的值(默认为10000,可由参数配置)分割为多个index组。以10000条记录为一个组,对数据进行统计。Hive查询引擎会将where条件中的约束传递给ORC reader,这些reader根据组级别的统计信息,过滤掉不必要的数据。如果该值设置的太小,就会保存更多的统计信息,用户需要根据自己数据的特点权衡一个合理的值。

1.4 Stripe → Row Group → Column → Stream 的关系

1
2
3
4
5
6
7
8
9
10
11
12
Stripe
├── Row Group 0
│ ├── Column A
│ │ ├── Stream PRESENT
│ │ ├── Stream DATA
│ │ └── Stream LENGTH / DICTIONARY / ...
│ ├── Column B
│ │ ├── Stream PRESENT
│ │ └── Stream DATA
│ └── ...
├── Row Group 1
└── ...
  • Row Group 是索引与统计的最小过滤单位
  • 真正的数据存储是按 column + stream 来的
  • Row Group 并不单独存数据块,而是通过 ROW_INDEX 指向各个 stream 在文件中的偏移位置
  • 每个Stream只属于一个字段,但一个字段可以包含多个Stream

2 ORC file API

源码主要有core和mapreduce两个模块,其中MapReduce 模块就是Hadoop的InputFormat和OutFormat,这个模块下有两个包,mapred和MapReduce分别符合MapReduce的V1和V2的Input/Output Format.

四个类:OrcInputFormat、OrcOutputFormat、OrcMapreduceRecordReader和OrcMapreduceRecordWriter

2.1 OrcInputFormat与OrcOutputFormat

OrcInputFormat继承了FileInputFormat,getSplits用的就是FileInputFormat的实现,只是重写了createRecordReader方法。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
@Override
public RecordReader<NullWritable, V>
createRecordReader(InputSplit inputSplit,
TaskAttemptContext taskAttemptContext
) throws IOException, InterruptedException {
FileSplit split = (FileSplit) inputSplit;
Configuration conf = taskAttemptContext.getConfiguration();
Reader file = OrcFile.createReader(split.getPath(),
OrcFile.readerOptions(conf)
.maxLength(OrcConf.MAX_FILE_LENGTH.getLong(conf)));
return new OrcMapreduceRecordReader<>(file,
org.apache.orc.mapred.OrcInputFormat.buildOptions(conf,
file, split.getStart(), split.getLength()));
}

这里面用到了org.apache.orc.Reader和org.apache.orc.OrcFile. 模板里面的V通常会是OrcStruct,Orc中一个表的schema可以表示为一个OrcStruct。
同样滴,OrcOutputFormat当中也是重写了craeteRecordWriter方法:

1
2
3
4
5
6
7
8
9
10
@Override
public RecordWriter<NullWritable, V>
getRecordWriter(TaskAttemptContext taskAttemptContext
) throws IOException {
Configuration conf = taskAttemptContext.getConfiguration();
Path filename = getDefaultWorkFile(taskAttemptContext, EXTENSION);
Writer writer = OrcFile.createWriter(filename,
org.apache.orc.mapred.OrcOutputFormat.buildOptions(conf));
return new OrcMapreduceRecordWriter<V>(writer);
}

这里面用到了org.apache.orc.Writer和org.apache.orc.OrcFile.
以上用到的这三个类都是core 模块里的,之后再去读。

2.2 OrcMapreduceRecordReader

这个类里面其实复用了org.apache.orc.mapred包下面的一些Writable类,这些Writable类用来在MapReduce中支持Orc中的一些特殊数据类型,比如Map、List、Struct等等,可以将这些数据类型从DataInput中反序列化出来,也可以将数据序列化到DataOutput中去。

从构造方法可以看到:

1
2
3
4
5
6
7
8
9
10
11
12
public OrcMapreduceRecordReader(Reader fileReader,
Reader.Options options) throws IOException {
this.batchReader = fileReader.rows(options);
if (options.getSchema() == null) {
schema = fileReader.getSchema();
} else {
schema = options.getSchema();
}
this.batch = schema.createRowBatch();
rowInBatch = 0;
this.row = (V) OrcStruct.createValue(schema);
}

其中batchReader是org.apache.orc.RecordReader类型的,这个RecordReader在core中。batch是真正有数据的地方。batch的类型是org.apache.hadoop.hive.ql.exec.vector.VectorizedRowBatch,从java doc可以看出来它是干嘛用的:

1
2
3
4
5
6
7
/**
* A VectorizedRowBatch is a set of rows, organized with each column
* as a vector. It is the unit of query execution, organized to minimize
* the cost per row and achieve high cycles-per-instruction.
* The major fields are public by design to allow fast and convenient
* access by the vectorized query execution code.
*/

而schema是org.apache.orc.TypeDescription类型的,而TypeDescription其实是一个数据类型的描述,它有一个category,可以是INT、FLOAT等等,也可以是STRUCT。这里schema的category通常会是STRUCT。STRUCT和OrcStruct是对应的,是一个复合类型,其中可以包含其他类型的属性,用来表示一个表的schema。

row就是batchReader从batch中读出来的一行记录。

OrcMapreduceRecordReader中另外一个重要的方法是nextKeyValue:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
@Override
public boolean nextKeyValue() throws IOException, InterruptedException {
if (!ensureBatch()) {
return false;
}
if (schema.getCategory() == TypeDescription.Category.STRUCT) {
OrcStruct result = (OrcStruct) row;
List<TypeDescription> children = schema.getChildren();
int numberOfChildren = children.size();
for(int i=0; i < numberOfChildren; ++i) {
result.setFieldValue(i, OrcMapredRecordReader.nextValue(batch.cols[i], rowInBatch,
children.get(i), result.getFieldValue(i)));
}
} else {
OrcMapredRecordReader.nextValue(batch.cols[0], rowInBatch, schema, row);
}
rowInBatch += 1;
return true;
}

可以看出来,这个方法中,判断当schema是STRUCT的category时,将它看做一个表的schema,取出里面包含的children,即各个字段的TypeDescription,然后读取各个字段的值,存入row中。这其中复用了OrcMapredRecordReader.nextValue()方法.

2.3 OrcMapreduceRecordWriter

OrcMapreduceRecordWriter和OrcMapreduceRecordReader类似,只不过其中的batchReader换成了org.apache.orc
.Writer类型的writer。构造方法如下:

1
2
3
4
5
6
public OrcMapreduceRecordWriter(Writer writer) {
this.writer = writer;
schema = writer.getSchema();
this.batch = schema.createRowBatch();
isTopStruct = schema.getCategory() == TypeDescription.Category.STRUCT;
}

这个类的主要方法就是write方法:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
@Override
public void write(NullWritable nullWritable, V v) throws IOException {
// if the batch is full, write it out.
if (batch.size == batch.getMaxSize()) {
writer.addRowBatch(batch);
batch.reset();
}

// add the new row
int row = batch.size++;
// skip over the OrcKey or OrcValue
if (v instanceof OrcKey) {
v = (V)((OrcKey) v).key;
} else if (v instanceof OrcValue) {
v = (V)((OrcValue) v).value;
}
if (isTopStruct) {
for(int f=0; f < schema.getChildren().size(); ++f) {
OrcMapredRecordWriter.setColumn(schema.getChildren().get(f),
batch.cols[f], row, ((OrcStruct) v).getFieldValue(f));
}
} else {
OrcMapredRecordWriter.setColumn(schema, batch.cols[0], row, v);
}
}

可以看出来,write先把V类型(通常是OrcStruct)的记录写入batch,如果batch写满了,就批量写入writer。

2.4 代办

  • Flink消费orc文件,能否按照file+strip作为offset?

2.5 【参考文献】

  1. ORC源码阅读(1) - mapreduce 模块