ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Arrow C++ Skyhook 扫描示例实战:将过滤与投影下推到 Ceph 集群

Apache Arrow C++ Skyhook 扫描示例实战:将过滤与投影下推到 Ceph 集群 数据工程大数据序列化数据分析【免费下载链接】arrowApache 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.ccclsOSD 侧Ceph 对象类object class实现注册scan_op方法cls_skyhook.ccprotocol客户端与 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-test与skyhook-protocol-test两个测试。而是否启用 Skyhook 由顶层 cpp/CMakeLists.txt 中的ARROW_SKYHOOK选项控制。示例程序整体流程示例的入口main位于 dataset_skyhook_scan_example.cc只接受一个命令行参数数据集根路径的 URI若没有参数则直接返回成功这是为了兼容 CI 环境有参数时调用Main执行如下五步流水线InstantiateSkyhookFormat()创建并初始化skyhook::SkyhookFileFormat建立到 Ceph RADOS 的连接fs::FileSystemFromUri(dataset_root, path)把file:///mnt/cephfs/nyc之类的 URI 解析为FileSystem与路径GetDatasetFromPath()按路径是目录还是文件分别构建DatasetGetScannerFromDataset()构造Scanner应用投影列、过滤表达式与并行线程配置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::vectorstd::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.ccuse_threads控制Scanner::ToTable是否并行消费默认trueprojected_columns {total_amount}只请求total_amount一列避免反序列化不需要的列filter payment_type 1构造一个列引用与字面量的比较表达式来自arrow/compute/expression.h。这里payment_type同时也是下文 Hive 分区目录的字段因此该过滤条件在理想情况下可以直接利用分区信息裁剪。2. 数据集构建与 Hive 分区GetDatasetFromDirectoryL58-L83演示了从目录构建数据集的完整写法fs::FileSelector s; s.base_dir dir; s.recursive true; ds::FileSystemFactoryOptions options; options.partitioning std::make_sharedds::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_typeint32与VendorIDint32对应数据集目录nyc/payment_type1/VendorID1/...的目录结构factory-Inspect()推断所有文件公共 schemafactory-Finish()产出数据集。GetDatasetFromFileL85-L100是单文件版本而GetDatasetFromPathL102-L110先通过fs-GetFileInfo(path)判断路径是目录还是文件再分派到上述两个函数。3. RADOS 连接参数InstantiateSkyhookFormatL128-L157构造RadosConnCtx五个参数与源码 file_skyhook.h 中RadosConnCtx结构体一一对应参数示例值含义ceph_config_path/etc/ceph/ceph.confCeph 集群配置文件路径包含集群级配置与连接信息ceph_data_poolcephfs_data存放待扫描对象的数据池data poolceph_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();GetScannerFromDatasetL112-L126展示了标准的ScannerBuilder用法Project指定投影列、Filter指定过滤表达式、UseThreads控制并行度最后Finish()得到Scanner。Main中执行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-devRADOS C 客户端开发库供find_package(librados)使用rados-objclass-devCeph 对象类object class开发头文件编译cls_skyhook需要ceph、ceph-common、ceph-osd、ceph-mon、ceph-mgr、ceph-mds、rbd-mirror组成一个完整 Ceph 集群所需的核心组件monitor、manager、OSD、MDS、镜像等ceph-fuse用户态挂载 CephFS 到/mnt/cephfs的 FUSE 客户端rapidjson-dev、libboost-all-devSkyhook/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_SKYHOOKON \ -DARROW_PARQUETON \ -DARROW_WITH_SNAPPYON \ -DARROW_BUILD_EXAMPLESON \ -DARROW_DATASETON \ -DARROW_CSVON \ -DARROW_WITH_LZ4ON \ .. make -j install cp release/libcls_skyhook.so /usr/lib/x86_64-linux-gnu/rados-classes/各 CMake 选项的作用与依据选项作用依据ARROW_SKYHOOKON启用 Skyhook 构建使 cpp/CMakeLists.txt 进入src/skyhook子目录构建配置ARROW_PARQUETON构建 Parquet 支持Skyhook 扫描的数据格式与示例数据集均为 Parquet编译依赖ARROW_WITH_SNAPPYON、ARROW_WITH_LZ4ONParquet 压缩编解码器支持SNAPPY/LZ4编译依赖ARROW_BUILD_EXAMPLESON构建 examples/arrow 下的示例程序其中 CMakeLists.txt 把示例链接到arrow_skyhook并依赖parquet目标构建配置ARROW_DATASETON构建 Dataset APIarrow/dataset/*SkyhookFileFormat继承自arrow::dataset::FileFormat构建配置ARROW_CSVON构建 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含自动生成的fsid、auth client required none免认证、osd pool default size 1单副本、osd objectstore memstore内存存储、osd class load list *加载全部对象类等配置然后依次启动ceph-mon、ceph-osd、ceph-mgr并创建cephfs_data/cephfs_metadata两个池及 CephFS 文件系统。仓库中 ci/scripts/integration_skyhook.sh 提供了等价的可复现部署逻辑它生成同一套ceph.conf、启动单 OSDosd objectstore memstore、创建cephfs_data与cephfs_metadata池、通过ceph fs new cephfs cephfs_metadata cephfs_data建立 CephFS、把libcls_skyhook*拷贝到rados-classes/目录再执行ceph-fuse /mnt/cephfs挂载文件系统——这一步为后续把数据集放入 CephFS 并让示例通过本地路径扫描做好了准备。生成示例数据集第四步生成数据并放入 CephFSpip 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构造包含total_amount与fare_amount两列、共 500 行的 pandas DataFrame值随机生成写入skyhook_test_data.parquet按payment_type取值 1~4与VendorID取值 1~2两维分区生成目录结构nyc/payment_type{p}/VendorID{v}/{p}.{v}.parquet共 8 个 Parquet 文件。这个目录结构与示例代码中HivePartitioning声明的payment_type、VendorID两个分区字段完全对应。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是数据集根目录的 URIFileSystemFromUri会将其解析为本地文件系统路径正常运行会输出Table size: 行数若 RADOS 连接或 OSD 扫描失败会向stderr打印Status错误信息并以非零码退出。注意示例的main对无参数调用如 CI 冒烟测试会直接返回成功因此只有在传入正确 URI 时才会真正执行 Skyhook 下推扫描。底层原理扫描请求如何下推到 OSD客户端侧构造请求并调用对象类方法SkyhookFileFormat::ScanBatchesAsyncfile_skyhook.cc是下推的核心把字符串格式名映射为SkyhookFileType::type枚举PARQUET/IPC非法格式返回Status::Invalid通过SkyhookDirectObjectAccess::Stat对文件执行 POSIXstat取得文件大小组装skyhook::ScanRequest包含filter_expression过滤表达式、partition_expression分区表达式、projection_schema投影 schema、dataset_schema数据集 schema、file_size、file_format用skyhook::SerializeScanRequest把请求序列化进ceph::bufferlistflatbuffers 格式见 skyhook_protocol.h调用SkyhookDirectObjectAccess::Exec(st.st_ino, scan_op, request, result)执行对象类方法scan_opDeserializeTable把返回的 bufferlist 反序列化为RecordBatch向量并包装成RecordBatchGenerator实现中特别注明关闭线程解压以避免嵌套线程问题相关注释引用了历史问题 ARROW-12597。SkyhookDirectObjectAccessskyhook_protocol.h解决了一个关键映射问题CephFS 中一个文件在 RADOS 底层由多个 stripe 对象组成对象 ID 形如[hex(inode)].[stripe_index 的 8 位二进制]。Skyhook 保证每个文件只有一个 stripestripe index 恒为 0因此ConvertInodeToOID把 inode 转为十六进制后拼接.00000000即可得到唯一的 RADOS 对象 ID再通过librados::execAPI 在该对象上执行对象类方法。OSD 侧对象类扫描与回传libcls_skyhook.so的实现位于 cls_skyhook.cc。它通过CLS_VER(1, 0)、CLS_NAME(skyhook)声明版本与类名在__cls_initL263-L267中注册类skyhook和方法scan_op只读方法CLS_METHOD_RD。scan_opL211-L261的执行流程DeserializeScanRequest反序列化客户端请求失败返回错误码SCAN_REQ_DESER_ERR_CODE按req.file_format分派PARQUET走ScanParquetObject、IPC走ScanIpcObjectDoScanL153-L173是真正的扫描逻辑用RandomAccessObject包装 RADOS 对象实现arrow::io::RandomAccessFile接口内部用cls_cxx_read按位置读取对象字节见 L42-L144构造FileSource与FileFragment再用ScannerBuilder套用客户端的过滤、投影与线程配置执行scanner-ToTable()得到结果表SerializeTable把结果表序列化回 bufferlist写回out返回 0 表示成功失败则记录错误并返回SCAN_ERR_CODE/SCAN_RES_SER_ERR_CODE等错误码。简言之客户端只发送我要哪些列、过滤条件是什么、文件在哪的描述性请求而真正的 Parquet 解码、过滤与投影全部在数据所在的 OSD 上完成回传的只有满足条件的紧凑结果表。验证与测试仓库为 Skyhook 提供了两层验证手段单元测试构建时生成的skyhook-cls-test与skyhook-protocol-test分别覆盖 cls_skyhook_test.cc 与 skyhook_protocol_test.cc前者验证对象类扫描逻辑后者验证请求/响应的序列化往返一致性端到端集成测试ci/scripts/integration_skyhook.sh 在ARROW_SKYHOOKON时完整走一遍启动单节点 Cephmemstore 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 大小的大文件不适用该直接映射格式限制客户端侧格式串仅支持parquet与ipc其余取值会在ScanBatchesAsync中返回Invalid错误写入MakeWriter尚未实现返回NotImplemented因此 Skyhook 当前定位为扫描路径的加速方案。通过本文的 5 个步骤与源码级剖析你可以在本地完整复现Arrow Skyhook Ceph的计算下推扫描并以此为基础进一步阅读 skyhook 协议定义 与 对象类实现理解存储侧计算的具体实现细节。赞分享数据工程大数据序列化数据分析【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址https://gitcode.com/gh_mirrors/arrow13/arrow点击查看免费下载相关推荐Apache Arrow C 数据集 Skyhook 扫描示例将过滤与投影下推到 Ceph 集群Apache Arrow C 数据集 Skyhook 扫描示例将过滤与投影下推到 Ceph 集群 Apache Arrow 仓库中的 dataset_sk数据工程数据分析大数据Apache Arrow C Datasets 实战指南从读取、过滤、投影到分区读写Apache Arrow C Datasets 实战指南从读取、过滤、投影到分区读写 本文以 docs/source/cpp/examples/datas数据工程数据分析大数据Apache Arrow Gandiva 表达式、投影器Projector与过滤器FilterC 实战指南Apache Arrow Gandiva 表达式、投影器Projector与过滤器FilterC 实战指南 本篇技术指南以 Apache Arrow数据工程数据分析大数据创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表