资讯动态

Apache Arrow IPC 列式序列化与进程间通信格式深度解析

发布时间:2026/9/23 21:47:53 来源:尧图企业网站定制
数据工程大数据序列化数据分析【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址https://gitcode.com/gh_mirrors/arrow13/arrow点击查看免费下载本指南以 Apache Arrow 官方列式规范文档中的「Serialization and Interprocess Communication (IPC)」章节即 docs/source/format/IPC.rst 所指向的正文位于 Columnar.rst 的 format-ipc 锚点为骨架结合仓库内的 FlatBuffers 协议定义与 C 参考实现系统讲解 Arrow 的二进制消息封装、流式格式Streaming Format、文件格式File Format与字典消息协议。读完本文你将掌握 Arrow IPC 消息的逐字节布局、Schema/RecordBatch/DictionaryBatch 三类消息的内部结构、缓冲区扁平化规则以及如何在应用层通过custom_metadata扩展协议语义。说明本文所述规范在仓库中的正式定义文件为 format/Message.fbs、format/Schema.fbs 与 format/File.fbsC 参考实现位于 cpp/src/arrow/ipc。文中所有二进制布局与约定均以这些文件与文档为准。基本概念Record Batch 与 SchemaArrow 列式格式中序列化数据的基本单元是record batch记录批。从语义上讲record batch 是有序的数组集合其中每个数组称为一个字段field所有字段长度相同但数据类型可以不同。record batch 中字段的名称与类型共同构成该批的schema模式。IPC 协议的核心目标是定义一套把 record batch 序列化为二进制载荷流、并在不进行内存复制的前提下从这些载荷重建 record batch 的机制。为此Arrow 定义了单向的二进制消息流其中只包含三种核心消息类型Schema模式消息RecordBatch记录批消息DictionaryBatch字典消息在此基础上规范进一步定义了所谓的封装 IPC 消息encapsulated IPC message格式一个序列化的 FlatBuffer 元数据外加一个可选的 message body。在描述如何序列化各组成消息类型之前必须先定义这种封装格式。封装消息格式Encapsulated Message Format对于简单的流式与基于文件的序列化场景Arrow 定义了「封装」消息格式用于进程间通信。此类消息只需检查消息元数据就能反序列化为内存中的 Arrow 数组对象无需复制或搬移任何实际数据零拷贝重建的前提。封装后的二进制消息布局如下continuation: 0xFFFFFFFF metadata_size: int32 metadata_flatbuffer: bytes padding message body各部分含义为32 位 continuation 指示符continuation indicator固定取值0xFFFFFFFF表示后面是一条有效消息。这一组件自 Arrow 0.15.0 引入部分目的是解决 FlatBuffers 的 8 字节对齐要求。32 位小端little-endian长度前缀表示后面 metadata 的字节大小。消息元数据metadata_flatbuffer使用 Message.fbs 中定义的Message类型序列化。填充字节padding将元数据补齐到 8 字节边界。消息体message body其长度必须是 8 的倍数。整条序列化消息必须为 8 的倍数这样才能保证消息可以在不同流之间重新定位否则元数据与消息体之间的填充量将无法确定non-deterministic。关于metadata_size的边界它包含了Message本身及其填充的总大小。metadata_flatbuffer中的MessageFlatBuffer 值内部包含一个版本号version一个具体的消息值Schema、RecordBatch或DictionaryBatch之一消息体的大小body length一个custom_metadata字段用于携带应用自定义元数据Message的 FlatBuffer 定义format/Message.fbs与文档一致table Message { version: org.apache.arrow.flatbuf.MetadataVersion; header: MessageHeader; bodyLength: long; custom_metadata: [ KeyValue ]; }其中header是一个联合类型union在MessageHeader中除三种核心消息外还包含Tensor与SparseTensorformat/Message.fbs。规范建议为获得最大兼容性最好使用RecordBatch传输数据Arrow 实现无需实现全部消息类型。读取时的一般流程是先从输入流解析并校验Message元数据以获得 body 大小再读取 body。这一点在 C 参考实现中得到印证cpp/src/arrow/ipc/message.cc 中通过「跳过 8 字节32 位 continuation 指示符 32 位小端长度前缀」来定位并校验元数据而 continuation 令牌在 cpp/src/arrow/ipc/metadata_internal.h 中定义为kIpcContinuationToken -1其二进制表示正是0xFFFFFFFF。Schema 消息FlatBuffers 文件 format/Schema.fbs 包含所有内建数据类型以及代表给定 record batch 模式的Schema元数据类型的定义。schema 由有序的字段序列组成每个字段包含名称与类型。序列化的Schema不包含任何数据缓冲区只包含类型元数据。FieldFlatBuffer 类型保存单个数组的元数据包括字段名称字段数据类型字段在语义上是否可空nullable。虽然这对数组的物理布局没有影响但许多系统会区分可空与不可空字段Arrow 允许保留该元数据以实现忠实的 schema 往返schema round trip子字段childField集合用于嵌套类型dictionary属性表示字段是否被字典编码若被编码则分配一个字典「id」用于将后续字典 IPC 消息与对应字段匹配从 format/Schema.fbs 可以看到Field表定义中确实包含name、nullable、dictionaryDictionaryEncoding以及custom_metadata等字段Schema表同样提供custom_metadataformat/Schema.fbs。此外Arrow 在 schema 级和 field 级都提供custom_metadata属性允许各系统插入自己的应用定义元数据以定制行为。RecordBatch 消息RecordBatch消息包含与 schema 决定的物理内存布局相对应的实际数据缓冲区。该消息的元数据提供每个缓冲区的位置与大小从而允许仅通过指针算术即可重建 Array 数据结构无需内存复制。record batch 的序列化形式包含两部分数据头data header即 Message.fbs 中定义的RecordBatch类型body一段平坦的内存缓冲区序列端到端写入并带有保证最小 8 字节对齐的填充数据头中包含record batch 中每个扁平化字段flattened field的长度与空值计数null countbody 中每个组成Buffer的内存偏移量与长度FieldNode结构体format/Message.fbs定义字段节点元数据length表示嵌套树中该层 Arrow 数组的值槽数量null_count表示观测到的空值数量。值得注意的是null_count 0的字段可以选择不写出物理有效性位图而将该位图缓冲区的长度设为 0。RecordBatch表中的nodes[FieldNode]与buffers[Buffer]分别对应预序遍历扁平化的逻辑 schema 与缓冲区树。字段与缓冲区的扁平化字段与缓冲区通过对 record batch 字段进行前序深度优先遍历来扁平化。文档给出了如下示例 schemacol1: Structa: Int32, b: Listitem: Int64, c: Float64 col2: Utf8其扁平化结果为FieldNode 0: Struct namecol1 FieldNode 1: Int32 namea FieldNode 2: List nameb FieldNode 3: Int64 nameitem FieldNode 4: Float64 namec FieldNode 5: Utf8 namecol2对应产生的缓冲区序列为buffer 0: field 0 validity buffer 1: field 1 validity buffer 2: field 1 values buffer 3: field 2 validity buffer 4: field 2 offsets buffer 5: field 3 validity buffer 6: field 3 values buffer 7: field 4 validity buffer 8: field 4 values buffer 9: field 5 validity buffer 10: field 5 offsets buffer 11: field 5 data这里可以看到规律每个字段都可能有 validity有效性位图定长值类型有 values 缓冲区List/Utf8等变长类型有 offsets偏移量与 data/values 缓冲区Struct只有 validity。BufferFlatBuffer 值描述一段内存的位置与大小通常相对于上述「封装消息格式」来解释。规范特别提醒Buffer的size字段不要求包含填充字节。由于该元数据可用于在库之间传递内存指针地址建议将size设为实际内存大小而非填充后的大小。变长缓冲区Variadic Buffers自Arrow Columnar Format 1.4起规范引入了变长缓冲区支持。Utf8View、BinaryView等类型使用可变数量的缓冲区表示对于预排序扁平化逻辑 schema 中的每个此类FieldRecordBatch中会有一个variadicBufferCounts条目指示该字段在当前 record batch 中占用的变长缓冲区数量。示例 schemacol1: Structa: Int32, b: BinaryView, c: Float64 col2: Utf8View该 schema 有两个含变长缓冲区的字段因此每个 record batch 的variadicBufferCounts都有两个条目。当variadicBufferCounts [3, 2]时扁平化缓冲区为buffer 0: col1 validity buffer 1: col1.a validity buffer 2: col1.a values buffer 3: col1.b validity buffer 4: col1.b views buffer 5: col1.b data buffer 6: col1.b data buffer 7: col1.b data buffer 8: col1.c validity buffer 9: col1.c values buffer 10: col2 validity buffer 11: col2 views buffer 12: col2 data buffer 13: col2 data即col1.bBinaryView拥有 1 个 validity、1 个 views 和 3 个 data 缓冲区col2Utf8View拥有 1 个 validity、1 个 views 和 2 个 data 缓冲区与variadicBufferCounts [3, 2]一一对应。该字段在 Message.fbs 中定义为variadicBufferCounts: [long]并且「仅当 schema 中不含任何可变缓冲区字段如 BinaryView/Utf8View时才可省略」。字节序EndiannessArrow 格式默认采用小端little-endian。序列化的Schema元数据中包含一个 endianness 字段表示 RecordBatch 的字节序通常即生成该 RecordBatch 的系统的字节序。其主要用途是在字节序相同的系统之间交换 RecordBatch。规范指出最初在读取字节序与底层系统不匹配的 Schema 时会返回错误参考实现聚焦于小端并提供相应测试未来可能通过字节交换提供自动转换。IPC 流式格式Streaming FormatArrow 提供用于 record batch 的流式协议Streaming Format它呈现为一系列封装消息的序列每条消息都遵循上述封装格式。schema 位于流的开头其后所有 record batch 共用同一 schema。如果 schema 中有字段被字典编码则会包含一个或多个DictionaryBatch消息。DictionaryBatch与RecordBatch消息可以交错排列但任何字典键dictionary key在RecordBatch中使用之前必须先在其之前的DictionaryBatch中定义。流的总体结构SCHEMA DICTIONARY 0 ... DICTIONARY k - 1 RECORD BATCH 0 ... DICTIONARY x DELTA ... DICTIONARY y DELTA ... RECORD BATCH n - 1 EOS [optional]: 0xFFFFFFFF 0x00000000边界情况说明当 record batch 包含完全为空的字典编码数组时存在交错字典与 record batch 的边缘情况——此时编码列对应的字典可能出现在第一个 record batch 之后。流读取器实现读取流时在每条消息之后可以读取接下来的 8 字节以判断流是否继续以及后续消息元数据的大小读取完消息 flatbuffer 后再读取消息体。流写入器可以通过两种方式发出流结束信号EOS写入 8 字节4 字节 continuation 指示符0xFFFFFFFF后跟 0 元数据长度0x00000000或直接关闭流接口规范推荐使用.arrows文件扩展名表示流式格式尽管在许多场景下这些流根本不会被存储为文件。流式格式不支持随机访问只支持顺序读取。IPC 文件格式File FormatArrow 定义了一种支持随机访问的「文件格式」它是流式格式的扩展magic number ARROW1 empty padding bytes [to 8 byte boundary] STREAMING FORMAT with EOS FOOTER FOOTER SIZE: int32 magic number ARROW1文件以魔数字符串ARROW1加填充开始和结束。该魔数在 C 实现中定义于 cpp/src/arrow/ipc/metadata_internal.hkArrowMagicBytes ARROW1并由 cpp/src/arrow/ipc/feather.cc 用于文件类型识别。魔数之后的内容与流式格式完全相同。在文件末尾写入一个footer页脚其中包含 schema 的冗余副本schema 本就是流式格式的一部分以及文件中每个数据块的内存偏移量与大小。正是 footer 使得对文件中任意 record batch 的随机访问成为可能。footer 的精确细节定义于 format/File.fbs。从 format/File.fbs 可见Footer表包含version、schema、dictionaries[Block]、recordBatches[Block]与custom_metadataBlock结构体由offset指向 RecordBlock 起始位置注意是消息头之后的位置、metaDataLength元数据长度与bodyLength数据长度已对齐可与元数据之间存在空隙组成。文件格式与流式格式在字典使用上有两点差异在文件格式中不要求字典键在RecordBatch使用前先通过DictionaryBatch定义——只要键在文件中的某处已定义即可。每个字典 ID 最多只能有一个非 delta字典批即不支持字典替换delta 字典按其在文件 footer 中出现的顺序应用。文件格式推荐使用.arrow扩展名。需要特别说明的是用该格式创建的文件有时被称为「Feather V2」并使用.feather扩展名。FeatherV1是 Arrow 项目早期的一个概念验证用于 Pythonpandas与 R 的语言无关的快速数据帧存储Feather V2 的名称与扩展名即源于此。字典消息Dictionary Messages字典在流式与文件格式中作为一系列单字段 record batch写入。因此一串 record batch 的完整语义 schema 由 schema 与所有字典共同构成。字典类型位于 schema 中所以必须先读取 schema 以确定字典类型然后才能正确解释字典。DictionaryBatch的 FlatBuffer 定义为table DictionaryBatch { id: long; data: RecordBatch; isDelta: boolean false; }定义见 format/Message.fbs。字典id可以在 schema 中被引用一次或多次因此同一个字典可以服务于多个字段。isDelta标志允许对已有字典进行扩展以用于后续 record batch 的物化isDelta为 true表示该批字典向量应与之前同id的任何批次向量拼接追加。isDelta为 false默认表示该字典替换同 ID 的现有字典。文档给出了一个非常直观的 delta 示例。对于一个编码列字符串序列[A, B, C, B, D, C, E, A]使用 delta 字典批可以编码为SCHEMA DICTIONARY 0 (0) A (1) B (2) C RECORD BATCH 0 0 1 2 1 DICTIONARY 0 DELTA (3) D (4) E RECORD BATCH 1 3 2 4 0 EOS这里RECORD BATCH 0用索引0,1,2,1表示A,B,C,BDICTIONARY 0 DELTA追加D索引 3与E索引 4RECORD BATCH 1用3,2,4,0表示D,C,E,A。整个过程中字典通过追加不断增长索引无需重排。对应的非 delta替换编码方式则是在第二个DICTIONARY 0中完整给出新字典SCHEMA DICTIONARY 0 (0) A (1) B (2) C RECORD BATCH 0 0 1 2 1 DICTIONARY 0 (0) A (1) C (2) D (3) E RECORD BATCH 1 2 1 3 0 EOS此时字典被整体替换但注意文件格式不支持字典替换同一字典 ID 只能有一个非 delta 字典批delta 是文件格式中更稳妥的增量方式。自定义应用元数据Custom Application MetadataArrow 在三个层级提供custom_metadata字段作为开发者在 Arrow 协议消息中传递应用特定元数据的机制Field级Schema级Message级命名空间规则冒号符号:用作命名空间分隔符且一个 key 中可以多次使用。ARROW模式是保留命名空间专供 Arrow 内部在custom_metadata字段中使用例如ARROW:extension:name。扩展类型Extension Types用户自定义的「扩展」类型可以通过在Field元数据结构的custom_metadata中设置特定KeyValue对来定义。扩展键如下ARROW:extension:name标识自定义数据类型的字符串名称。规范建议使用「命名空间」式前缀以最大程度减少同一应用中多个 Arrow 读写器之间的冲突例如使用myorg.name_of_type而非简单的name_of_type。ARROW:extension:metadata用于重建自定义类型的ExtensionType的序列化表示。保留规则以arrow.开头的扩展名是保留给规范化的canonical扩展类型使用的第三方扩展类型不应使用此前缀。扩展元数据可以注解任意内建 Arrow 逻辑类型。其设计意图是不支持某扩展类型的实现仍然可以处理底层数据。例如16 字节的 UUID 值可以嵌入FixedSizeBinary(16)中不支持该扩展类型的实现仍能处理底层二进制值并在后续 Arrow 协议消息中原样传递custom_metadata。扩展类型可以使用也可以不使用ARROW:extension:metadata字段。文档给出了四类示例uuid表示为FixedSizeBinary(16)元数据为空latitude-longitude表示为structlatitude: double, longitude: double元数据为空tensor多维数组存储为Binary值带有指示每个值的数据类型与形状的序列化元数据例如 JSON{type: int8, shape: [4, 5]}表示 4x5 的单元张量trading-time表示为Timestamp带有的序列化元数据指示数据对应的市场交易日历实现指南子集实现与扩展性规范在结尾给出了面向执行引擎或框架、UDF 执行器、存储引擎等的实现约束实现规范子集如果只产生produce而不消费consumeArrow 向量可以实现向量规范与对应元数据的任意子集。如果同时消费与产生向量存在一个需支持的最小向量子集。产生任意子集向量及其元数据总是可以的但消费向量时至少应将不支持的输入向量转换为受支持的子集例如将Timestamp.millis转换为timestamp.micros或将int32转换为int64。扩展性执行引擎实现者也可以在内部使用自定义向量扩展其内存表示前提是这些自定义向量绝不对外暴露。在将数据发送给期待 Arrow 数据的其他系统之前这些自定义向量必须转换为 Arrow 规范中存在的类型。结语从规范到代码本文所述的 IPC 协议并非停留在纸面仓库中的 format/Message.fbs、format/Schema.fbs、format/File.fbs 是协议的权威定义cpp/src/arrow/ipc 目录下的message.cc、metadata_internal.h、feather.cc等文件则提供了完整的参考实现包括 continuation 令牌、ARROW1魔数、消息读写与校验逻辑。若要深入理解流式读写与随机访问文件的落地细节可以从 cpp/src/arrow/ipc/reader.cc 与 cpp/src/arrow/ipc/writer.cc 入手结合本文的封装消息布局与 File/Stream 结构逐段对照阅读。赞分享数据工程大数据序列化数据分析【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址https://gitcode.com/gh_mirrors/arrow13/arrow点击查看免费下载相关推荐Apache Arrow IPC 序列化与进程间通信格式完全指南从封装消息到流式与文件格式Apache Arrow IPC 序列化与进程间通信格式完全指南从封装消息到流式与文件格式 Apache Arrow 列式格式的原始序列化单元是 record大数据数据分析数据工程序列化Apache Arrow Columnar Format 深度解析列式内存布局与 IPC 序列化规范Version 1.5Apache Arrow Columnar Format 深度解析列式内存布局与 IPC 序列化规范Version 1.5 Apache Arrow 的大数据数据分析数据工程序列化Apache Arrow pyarrow 序列化与 IPC 编程指南Streaming 与 File 格式全解Apache Arrow pyarrow 序列化与 IPC 编程指南Streaming 与 File 格式全解 本文是 Apache Arrow Python数据工程大数据序列化数据分析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价