Apache Arrow C++ Skyhook 扫描示例实战:将过滤与投影下推到 Ceph 集群
2026/9/24 16:45:45 网站建设 项目流程
  • 数据工程
  • 大数据
  • 序列化
  • 数据分析

【免费下载链接】arrow

Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing

项目地址:https://gitcode.com/gh_mirrors/arrow13/arrow
点击查看免费下载

本文基于 Apache Arrow 仓库中的官方示例 dataset_skyhook_scan_example.cc 与其配套文档 docs/source/cpp/examples/dataset_skyhook_scan_example.rst,完整讲解如何在 Ubuntu 20.04+ 环境下编译带 Skyhook 支持的 Arrow,部署单节点 Ceph 集群,并把 Parquet 扫描的过滤(Filter)与投影(Projection)算子下推到 Ceph OSD 端执行。读完本文,你将掌握 Skyhook 客户端库、RADOS 连接参数、OSD 端对象类方法(cls)的完整调用链路,以及一条从零开始可复现的端到端运行路径。

Skyhook 是什么:仓库内的布局与架构

Skyhook 是一种将 Arrow Dataset 的扫描算子(过滤、投影)下推到存储节点执行的方案:数据以 Parquet 或 Arrow IPC 格式的 RADOS 对象形式存放在 Ceph 中,扫描时客户端不把整份文件读回,而是把"扫描请求"(包含过滤表达式、投影 schema、分区表达式等)序列化后发给 Ceph OSD,由 OSD 上加载的libcls_skyhook.so对象类直接执行 Arrow 扫描并只返回结果表,从而显著减少集群与客户端之间的数据传输量。

在本仓库中,Skyhook 的完整实现位于 cpp/src/skyhook,按职责分为三个子目录:

子目录作用关键文件
client客户端侧:实现SkyhookFileFormat,作为arrow::dataset::FileFormat的一种自定义格式file_skyhook.h、file_skyhook.cc
clsOSD 侧:Ceph 对象类(object class)实现,注册scan_op方法cls_skyhook.cc
protocol客户端与 OSD 之间共享的请求/响应结构及 flatbuffers 序列化逻辑skyhook_protocol.h、skyhook_protocol.cc

构建层面,cpp/src/skyhook/CMakeLists.txt 通过find_package(librados REQUIRED)引入 Ceph 的 RADOS 库,产出两个目标:arrow_skyhook(客户端库)和cls_skyhook(需要拷贝到 OSD 的libcls_skyhook.so);同时编译skyhook-cls-testskyhook-protocol-test两个测试。而是否启用 Skyhook 由顶层 cpp/CMakeLists.txt 中的ARROW_SKYHOOK选项控制。

示例程序整体流程

示例的入口main位于 dataset_skyhook_scan_example.cc:只接受一个命令行参数(数据集根路径的 URI),若没有参数则直接返回成功(这是为了兼容 CI 环境);有参数时调用Main执行如下五步流水线:

  1. InstantiateSkyhookFormat():创建并初始化skyhook::SkyhookFileFormat,建立到 Ceph RADOS 的连接;
  2. fs::FileSystemFromUri(dataset_root, &path):把file:///mnt/cephfs/nyc之类的 URI 解析为FileSystem与路径;
  3. GetDatasetFromPath():按路径是目录还是文件分别构建Dataset
  4. GetScannerFromDataset():构造Scanner,应用投影列、过滤表达式与并行线程配置;
  5. scanner->ToTable():执行扫描并输出结果表行数Table size: N

从源码结构看,这个流程复用了 Arrow Dataset 标准的"文件系统 → 工厂 → 数据集 → 扫描器"链路,唯一的特殊之处在于第一步:文件格式不是ParquetFileFormat,而是把扫描请求下推到 OSD 的SkyhookFileFormat

示例代码核心要点

1. 配置结构体:投影、过滤与并行

struct Configuration { // Indicates if the Scanner::ToTable should consume in parallel. bool use_threads = true; // Indicates to the Scan operator which columns are requested. This // optimization avoid deserializing unneeded columns. std::vector<std::string> projected_columns = {"total_amount"}; // Indicates the filter by which rows will be filtered. This optimization can // make use of partition information and/or file metadata if possible. cp::Expression filter = cp::greater(cp::field_ref("payment_type"), cp::literal(1)); ds::InspectOptions inspect_options{}; ds::FinishOptions finish_options{}; } kConf;

三个关键配置(见 dataset_skyhook_scan_example.cc):

  • use_threads:控制Scanner::ToTable是否并行消费,默认true
  • projected_columns = {"total_amount"}:只请求total_amount一列,避免反序列化不需要的列;
  • filter = payment_type > 1:构造一个列引用与字面量的比较表达式,来自arrow/compute/expression.h。这里payment_type同时也是下文 Hive 分区目录的字段,因此该过滤条件在理想情况下可以直接利用分区信息裁剪。

2. 数据集构建与 Hive 分区

GetDatasetFromDirectory(L58-L83)演示了从目录构建数据集的完整写法:

fs::FileSelector s; s.base_dir = dir; s.recursive = true; ds::FileSystemFactoryOptions options; options.partitioning = std::make_shared<ds::HivePartitioning>( arrow::schema({arrow::field("payment_type", arrow::int32()), arrow::field("VendorID", arrow::int32())})); ARROW_ASSIGN_OR_RAISE(auto factory, ds::FileSystemDatasetFactory::Make(fs, s, format, options)); ARROW_ASSIGN_OR_RAISE(auto schema, factory->Inspect(kConf.inspect_options)); ARROW_ASSIGN_OR_RAISE(auto dataset, factory->Finish(kConf.finish_options));

要点:使用FileSelector递归发现目录下所有文件;用HivePartitioning声明分区字段payment_type(int32)与VendorID(int32),对应数据集目录nyc/payment_type=1/VendorID=1/...的目录结构;factory->Inspect()推断所有文件公共 schema;factory->Finish()产出数据集。GetDatasetFromFile(L85-L100)是单文件版本,而GetDatasetFromPath(L102-L110)先通过fs->GetFileInfo(path)判断路径是目录还是文件,再分派到上述两个函数。

3. RADOS 连接参数

InstantiateSkyhookFormat(L128-L157)构造RadosConnCtx,五个参数与源码 file_skyhook.h 中RadosConnCtx结构体一一对应:

参数示例值含义
ceph_config_path/etc/ceph/ceph.confCeph 集群配置文件路径,包含集群级配置与连接信息
ceph_data_poolcephfs_data存放待扫描对象的数据池(data pool)
ceph_user_nameclient.admin访问集群的 Ceph 用户名
ceph_cluster_nameceph集群名,多站点多集群架构下用于标识当前会话所属集群
ceph_cls_nameskyhook对象类名,对应 OSD 加载的libcls_skyhook.so中注册的类

随后skyhook::SkyhookFileFormat::Make(rados_ctx, "parquet")创建格式实例(第二个参数声明底层文件格式,当前支持"parquet""ipc")。Make内部会调用Init()建立到 RADOS 集群的连接并实例化SkyhookDirectObjectAccess(见 file_skyhook.cc)。

4. Scanner 构建与执行

ARROW_ASSIGN_OR_RAISE(auto scanner_builder, dataset->NewScan()); if (!columns.empty()) { ARROW_RETURN_NOT_OK(scanner_builder->Project(columns)); } ARROW_RETURN_NOT_OK(scanner_builder->Filter(filter)); ARROW_RETURN_NOT_OK(scanner_builder->UseThreads(use_threads)); return scanner_builder->Finish();

GetScannerFromDataset(L112-L126)展示了标准的ScannerBuilder用法:Project指定投影列、Filter指定过滤表达式、UseThreads控制并行度,最后Finish()得到ScannerMain中执行scanner->ToTable()后打印结果表行数:

ARROW_ASSIGN_OR_RAISE(auto table, scanner->ToTable()); std::cout << "Table size: " << table->num_rows() << "\n";

环境准备:安装 Ceph 与 Skyhook 依赖

官方文档要求 Ubuntu 20.04 或更高版本。第一步安装系统级依赖:

apt update apt install -y cmake \ libradospp-dev \ rados-objclass-dev \ ceph \ ceph-common \ ceph-osd \ ceph-mon \ ceph-mgr \ ceph-mds \ rbd-mirror \ ceph-fuse \ rapidjson-dev \ libboost-all-dev \ python3-pip

各包的角色:

  • cmake:构建 Arrow/Skyhook 的构建工具;
  • libradospp-dev:RADOS C++ 客户端开发库,供find_package(librados)使用;
  • rados-objclass-dev:Ceph 对象类(object class)开发头文件,编译cls_skyhook需要;
  • cephceph-commonceph-osdceph-monceph-mgrceph-mdsrbd-mirror:组成一个完整 Ceph 集群所需的核心组件(monitor、manager、OSD、MDS、镜像等);
  • ceph-fuse:用户态挂载 CephFS 到/mnt/cephfs的 FUSE 客户端;
  • rapidjson-devlibboost-all-dev:Skyhook/Arrow 构建过程中的 JSON 与 Boost 依赖;
  • python3-pip:用于安装 pandas/pyarrow 以生成示例数据集。

编译构建带 Skyhook 的 Arrow

第二步克隆源码并启用 Skyhook 相关编译选项:

git clone https://gitcode.com/gh_mirrors/arrow13/arrow cd arrow/ mkdir -p cpp/release cd cpp/release cmake -DARROW_SKYHOOK=ON \ -DARROW_PARQUET=ON \ -DARROW_WITH_SNAPPY=ON \ -DARROW_BUILD_EXAMPLES=ON \ -DARROW_DATASET=ON \ -DARROW_CSV=ON \ -DARROW_WITH_LZ4=ON \ .. make -j install cp release/libcls_skyhook.so /usr/lib/x86_64-linux-gnu/rados-classes/

各 CMake 选项的作用与依据:

选项作用依据
ARROW_SKYHOOK=ON启用 Skyhook 构建,使 cpp/CMakeLists.txt 进入src/skyhook子目录构建配置
ARROW_PARQUET=ON构建 Parquet 支持,Skyhook 扫描的数据格式与示例数据集均为 Parquet编译依赖
ARROW_WITH_SNAPPY=ONARROW_WITH_LZ4=ONParquet 压缩编解码器支持(SNAPPY/LZ4)编译依赖
ARROW_BUILD_EXAMPLES=ON构建 examples/arrow 下的示例程序,其中 CMakeLists.txt 把示例链接到arrow_skyhook并依赖parquet目标构建配置
ARROW_DATASET=ON构建 Dataset API(arrow/dataset/*),SkyhookFileFormat继承自arrow::dataset::FileFormat构建配置
ARROW_CSV=ON构建 CSV 支持(示例依赖的 Arrow Dataset 组件之一)构建配置

make -j install完成后,libcls_skyhook.so位于构建目录release/下。把该动态库拷贝到/usr/lib/x86_64-linux-gnu/rados-classes/后,Ceph OSD 启动时会按配置的 class list 加载它(参见下文osd class load list)。

部署单节点 Ceph 集群

第三步使用 skyhookdm 项目提供的micro-osd.sh脚本,在单机上拉起一个含单个内存 OSD 的最小 Ceph 集群:

./micro-osd.sh /tmp/skyhook

脚本会以/tmp/skyhook为数据目录,生成一套ceph.conf(含自动生成的fsidauth client required = none免认证、osd pool default size = 1单副本、osd objectstore = memstore内存存储、osd class load list = *加载全部对象类等配置),然后依次启动ceph-monceph-osdceph-mgr,并创建cephfs_data/cephfs_metadata两个池及 CephFS 文件系统。

仓库中 ci/scripts/integration_skyhook.sh 提供了等价的可复现部署逻辑:它生成同一套ceph.conf、启动单 OSD(osd objectstore = memstore)、创建cephfs_datacephfs_metadata池、通过ceph fs new cephfs cephfs_metadata cephfs_data建立 CephFS、把libcls_skyhook*拷贝到rados-classes/目录,再执行ceph-fuse /mnt/cephfs挂载文件系统——这一步为后续把数据集放入 CephFS 并让示例通过本地路径扫描做好了准备。

生成示例数据集

第四步生成数据并放入 CephFS:

pip install pandas pyarrow python3 ../../ci/scripts/generate_dataset.py cp -r nyc /mnt/cephfs/

注意../../ci/scripts/generate_dataset.py是相对cpp/release构建目录的路径;从仓库根目录看,脚本位于 ci/scripts/generate_dataset.py。该脚本的行为(见 generate_dataset.py):

  1. 构造包含total_amountfare_amount两列、共 500 行的 pandas DataFrame,值随机生成;
  2. 写入skyhook_test_data.parquet
  3. payment_type(取值 1~4)与VendorID(取值 1~2)两维分区,生成目录结构nyc/payment_type={p}/VendorID={v}/{p}.{v}.parquet,共 8 个 Parquet 文件。

这个目录结构与示例代码中HivePartitioning声明的payment_typeVendorID两个分区字段完全对应。cp -r nyc /mnt/cephfs/把整个数据集拷贝进 CephFS(数据实际以 RADOS 对象形式落在cephfs_data池中),因此示例运行时才能通过file:///mnt/cephfs/nyc定位到这些对象。

运行示例

第五步执行扫描:

LD_LIBRARY_PATH=/usr/local/lib release/dataset-skyhook-scan-example file:///mnt/cephfs/nyc

说明:

  • release/dataset-skyhook-scan-example是上一步make产出的示例可执行文件(对应 cpp/examples/arrow/CMakeLists.txt 中的add_arrow_example(dataset_skyhook_scan_example ...));
  • LD_LIBRARY_PATH=/usr/local/lib用于让程序找到安装到/usr/local/lib的 Arrow 动态库;
  • 参数file:///mnt/cephfs/nyc是数据集根目录的 URI,FileSystemFromUri会将其解析为本地文件系统路径;
  • 正常运行会输出Table size: <行数>;若 RADOS 连接或 OSD 扫描失败,会向stderr打印Status错误信息并以非零码退出。

注意:示例的main对无参数调用(如 CI 冒烟测试)会直接返回成功,因此只有在传入正确 URI 时才会真正执行 Skyhook 下推扫描。

底层原理:扫描请求如何下推到 OSD

客户端侧:构造请求并调用对象类方法

SkyhookFileFormat::ScanBatchesAsync(file_skyhook.cc)是下推的核心:

  1. 把字符串格式名映射为SkyhookFileType::type枚举(PARQUET/IPC),非法格式返回Status::Invalid
  2. 通过SkyhookDirectObjectAccess::Stat对文件执行 POSIXstat,取得文件大小;
  3. 组装skyhook::ScanRequest:包含filter_expression(过滤表达式)、partition_expression(分区表达式)、projection_schema(投影 schema)、dataset_schema(数据集 schema)、file_sizefile_format
  4. skyhook::SerializeScanRequest把请求序列化进ceph::bufferlist(flatbuffers 格式,见 skyhook_protocol.h);
  5. 调用SkyhookDirectObjectAccess::Exec(st.st_ino, "scan_op", request, result)执行对象类方法scan_op
  6. DeserializeTable把返回的 bufferlist 反序列化为RecordBatch向量,并包装成RecordBatchGenerator(实现中特别注明关闭线程解压以避免嵌套线程问题,相关注释引用了历史问题 ARROW-12597)。

SkyhookDirectObjectAccess(skyhook_protocol.h)解决了一个关键映射问题:CephFS 中一个文件在 RADOS 底层由多个 stripe 对象组成,对象 ID 形如[hex(inode)].[stripe_index 的 8 位二进制]。Skyhook 保证每个文件只有一个 stripe(stripe index 恒为 0),因此ConvertInodeToOID把 inode 转为十六进制后拼接.00000000即可得到唯一的 RADOS 对象 ID,再通过librados::execAPI 在该对象上执行对象类方法。

OSD 侧:对象类扫描与回传

libcls_skyhook.so的实现位于 cls_skyhook.cc。它通过CLS_VER(1, 0)CLS_NAME(skyhook)声明版本与类名,在__cls_init(L263-L267)中注册类skyhook和方法scan_op(只读方法CLS_METHOD_RD)。

scan_op(L211-L261)的执行流程:

  1. DeserializeScanRequest反序列化客户端请求,失败返回错误码SCAN_REQ_DESER_ERR_CODE
  2. req.file_format分派:PARQUETScanParquetObjectIPCScanIpcObject
  3. DoScan(L153-L173)是真正的扫描逻辑:用RandomAccessObject包装 RADOS 对象(实现arrow::io::RandomAccessFile接口,内部用cls_cxx_read按位置读取对象字节,见 L42-L144),构造FileSourceFileFragment,再用ScannerBuilder套用客户端的过滤、投影与线程配置,执行scanner->ToTable()得到结果表;
  4. SerializeTable把结果表序列化回 bufferlist,写回out,返回 0 表示成功;失败则记录错误并返回SCAN_ERR_CODE/SCAN_RES_SER_ERR_CODE等错误码。

简言之,客户端只发送"我要哪些列、过滤条件是什么、文件在哪"的描述性请求,而真正的 Parquet 解码、过滤与投影全部在数据所在的 OSD 上完成,回传的只有满足条件的紧凑结果表。

验证与测试

仓库为 Skyhook 提供了两层验证手段:

  1. 单元测试:构建时生成的skyhook-cls-testskyhook-protocol-test,分别覆盖 cls_skyhook_test.cc 与 skyhook_protocol_test.cc,前者验证对象类扫描逻辑,后者验证请求/响应的序列化往返一致性;
  2. 端到端集成测试:ci/scripts/integration_skyhook.sh 在ARROW_SKYHOOK=ON时完整走一遍"启动单节点 Ceph(memstore OSD)→ 创建 CephFS → 拷贝 cls 库 →ceph-fuse挂载 → 生成并拷贝 nyc 数据集 → 运行两个 skyhook 测试",与本文的 5 个步骤一一对应,是验证环境是否就绪的最佳参考。

前提与注意事项

  • 系统版本:官方文档明确要求 Ubuntu 20.04 或更高版本;
  • 必选依赖:RADOS 开发库(libradospp-dev)、对象类开发包(rados-objclass-dev)缺一不可,编译期会因find_package(librados REQUIRED)失败而报错;
  • 集群形态:示例针对单节点内存 OSD 的最小集群编写,ceph.conf使用免认证与单副本配置;生产多节点、多副本集群需相应调整 RADOS 参数与对象类加载配置;
  • 单 stripe 假设SkyhookDirectObjectAccess依赖"每个文件对应单个 RADOS 对象"的前提,超过 stripe 大小的大文件不适用该直接映射;
  • 格式限制:客户端侧格式串仅支持parquetipc,其余取值会在ScanBatchesAsync中返回Invalid错误;写入(MakeWriter)尚未实现,返回NotImplemented,因此 Skyhook 当前定位为扫描路径的加速方案。

通过本文的 5 个步骤与源码级剖析,你可以在本地完整复现"Arrow + Skyhook + Ceph"的计算下推扫描,并以此为基础进一步阅读 skyhook 协议定义 与 对象类实现,理解存储侧计算的具体实现细节。

  • 数据工程
  • 大数据
  • 序列化
  • 数据分析

【免费下载链接】arrow

Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing

项目地址:https://gitcode.com/gh_mirrors/arrow13/arrow
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询