ARTICLE DETAIL

资讯详情

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

cuDF libcudf 列聚合(Column Aggregation)API 全解:聚合工厂、归约扫描、GroupBy 与滚动窗口

cuDF libcudf 列聚合(Column Aggregation)API 全解:聚合工厂、归约扫描、GroupBy 与滚动窗口 数据分析数据工程机器学习【免费下载链接】cudfcuDF - GPU DataFrame Library项目地址https://gitcode.com/gh_mirrors/cu/cudf点击查看免费下载cuDF 是 NVIDIA 开源的 GPU DataFrame 库其底层 C 库 libcudf 提供了一套统一的列聚合Column Aggregation抽象供归约reduction、分组聚合groupby、滚动窗口rolling window等场景复用。本文以 libcudf API 文档 column_aggregation.rst 为纲结合 aggregation.hpp、reduction.hpp、groupby.hpp 与 rolling.hpp 的源码级注释系统讲解聚合类型枚举、全部聚合工厂函数的参数语义、输出类型规则与典型调用方式帮助你直接在 libcudf 层面编写 GPU 聚合代码。一、聚合体系概览一个抽象基类贯穿所有聚合 APIlibcudf 中所有想对一列数据做什么的意图都被建模为cudf::aggregation对象。在 aggregation.hpp 中aggregation是抽象基类它通过Kind枚举标识具体的聚合操作并支持is_equal()相等比较、do_hash()哈希与clone()克隆。class aggregation { public: enum Kind : int32_t { /* 见下文完整枚举 */ }; Kind kind; /// 要执行的聚合操作 virtual ~aggregation() default; [[nodiscard]] bool is_valid() const; [[nodiscard]] virtual bool is_equal(aggregation const other) const; [[nodiscard]] virtual size_t do_hash() const; [[nodiscard]] virtual std::unique_ptraggregation clone() const 0; };为区分哪种 API 能接受哪种聚合该头文件定义了若干虚拟派生类作为用途标记派生类用途reduce_aggregation整列归约cudf::reducesegmented_reduce_aggregation分段归约cudf::segmented_reducescan_aggregation前缀扫描cudf::scangroupby_aggregation分组聚合groupby::aggregategroupby_scan_aggregation分组扫描groupby::scanrolling_aggregation滚动窗口聚合cudf::rolling_window这种设计使得同一套聚合语义如 SUM、MEAN、RANK可以跨 API 复用而各 API 通过模板工厂参数Base约束可用的聚合类型。注意并非所有聚合都支持所有 APIaggregation.hpp 文件头注释明确提醒具体支持情况需查阅各函数文档。二、聚合类型枚举aggregation::Kind全景aggregation.hpp 定义了完整的聚合操作枚举enum Kind : int32_t { SUM, // 求和 SUM_OVERFLOW, // 带溢出检测的求和 PRODUCT, // 求积 MIN, MAX, // 最小 / 最大 COUNT_VALID, // 统计有效非空元素数 COUNT_ALL, // 统计全部元素数 ANY, ALL, // 逻辑或 / 逻辑与 SUM_OF_SQUARES, // 平方和 MEAN, // 算术均值 M2, // 与均值差值的平方和 VARIANCE, STD, // 方差 / 标准差 MEDIAN, // 中位数 QUANTILE, // 指定分位数 ARGMAX, ARGMIN, // 最大 / 最小元素的索引 NUNIQUE, // 去重后的元素个数 NTH_ELEMENT, // 取第 n 个元素 ROW_NUMBER, // 当前行的行号相对滚动窗口 EWMA, // 指数加权移动平均 RANK, // 排名 COLLECT_LIST, // 收集为列表 COLLECT_SET, // 收集为去重列表 LEAD, LAG, // 窗口函数行偏移访问 HOST_UDF, // 基于主机端的 UDF 聚合 MERGE_LISTS, MERGE_SETS, // 合并多个列表 / 去重合并 MERGE_M2, // 合并 M2 部分结果 COVARIANCE, CORRELATION, // 协方差 / 相关系数 TDIGEST, MERGE_TDIGEST, // t-digest 分位数估计及其合并 HISTOGRAM, MERGE_HISTOGRAM, // 频数统计及其合并 BITWISE_AGG, // 位运算聚合AND/OR/XOR TOP_K, // 每组前 K 个元素 INVALID // 占位标识非法聚合 };aggregation::Kind::INVALID用于默认构造占位构造aggregation(Kind)时会通过is_valid()校验非法 kind 会触发CUDF_EXPECTS断言失败。aggregation的默认无参构造函数被CUDF_FAIL显式禁止源码注释明确指出其仅为满足编译器要求而存在绝不应被调用。除 Kind 外同一头文件还定义了聚合所需的一组辅助枚举rank_methodFIRST稳定序排名无并列、AVERAGE并列取均值、MIN并列取最小名次、MAX并列取最大名次、DENSE名次连续递增不跳号rank_percentageNONE原始排名、ZERO_NORMALIZEDrank / count、ONE_NORMALIZED(rank-1) / (count-1)bitwise_opAND/OR/XOR用于数值列的位运算聚合correlation_typePEARSON/KENDALL/SPEARMANewm_historyINFINITE/FINITE决定 EWMA 对序列首值的数学处理假设。三、聚合工厂函数aggregation_factoriesmake_*_aggregation全参数解析聚合对象统一由模板工厂函数创建模板参数Base aggregation可替换为上述任一用途派生类从而把同一聚合适配到不同 API。以下按语义分组列出全部工厂及其关键参数均定义于 aggregation.hpp 的addtogroup aggregation_factories组即 aggregation_factories.rst 的渲染内容。3.1 基础统计类工厂函数参数与默认值说明make_sum_aggregation()无求和make_sum_overflow_aggregation()无求和并检测整数/十进制溢出make_product_aggregation()无求积make_min_aggregation()/make_max_aggregation()无最小 / 最大make_count_aggregation(null_handling null_policy::EXCLUDE)null_policy计数默认不统计空值make_any_aggregation()/make_all_aggregation()无逻辑或 / 逻辑与make_sum_of_squares_aggregation()无平方和make_mean_aggregation()无算术均值make_median_aggregation()无中位数make_histogram_aggregation()无统计每个元素出现的频数3.2 方差与标准差make_variance_aggregation(size_type ddof 1); // 方差除数 N - ddof make_std_aggregation(size_type ddof 1); // 标准差除数 N - ddof make_m2_aggregation(); // M2 SUM((x - MEAN)^2)ddofDelta Degrees of Freedom默认值为 1对应样本方差/标准差除数N - 1传入 0 则对应总体方差/标准差。M2是并行计算方差的标准中间量——多个独立集合的 M2 可通过make_merge_m2_aggregation()合并为一次性在所有集合上计算的等价结果。这三个聚合对 chrono时间类型与复合类型会抛出cudf::logic_error。3.3 分位数make_quantile_aggregation(std::vectordouble const quantiles, interpolation interp interpolation::LINEAR);quantiles传入期望的分位点向量如{0.25, 0.5, 0.75}interp指定插值方式默认LINEAR。3.4 索引与元素定位类make_argmax_aggregation()/make_argmin_aggregation()返回最大/最小元素的索引make_nunique_aggregation(null_handling null_policy::EXCLUDE)去重计数默认排除空值make_nth_element_aggregation(size_type n, null_policy null_handling null_policy::INCLUDE)取组/序列第 n 个元素。n的取值范围是[-group_size, group_size)负索引[-group_size, -1]分别对应[0, group_size-1]越界时该组结果为空值make_row_number_aggregation()返回当前行的行号相对滚动窗口make_top_k_aggregation(size_type k, order topk_order order::DESCENDING)每组取前 K 个值默认按降序选取但返回结果不保证排序。3.5 排名RANKmake_rank_aggregation(rank_method method, order column_order order::ASCENDING, null_policy null_handling null_policy::EXCLUDE, null_order null_precedence null_order::AFTER, rank_percentage percentage rank_percentage::NONE);排名聚合只适用于scan 算法输入列是 orderby 列它决定被排名的行的顺序多列排序时应传入包含这些列的 struct 列。返回列类型规则源码注释明确给出除AVERAGE方法外的所有方法以及percentage NONE时返回size_type列AVERAGE方法与percentage ! NONE组合时返回double列。aggregation.hpp 提供了一个完整的赛车数据示例按venue分组、按time排名的silverstone/monza两站数据first: {1, 2, 3, 4, 5, 1, 2, 3, 4, 5} average: {1, 2, 3.5, 3.5, 5, 1, 2.5, 2.5, 4, 5} min: {1, 2, 3, 3, 5, 1, 2, 2, 4, 5} max: {1, 2, 4, 4, 5, 1, 3, 3, 4, 5} dense: {1, 2, 3, 3, 4, 1, 2, 2, 3, 4}同一组数据的百分比排名min方法对比NONE为原始排名ZERO_NORMALIZED为 rank/countONE_NORMALIZED为 (rank-1)/(count-1)后者首名固定为 0.00、末名为 1.00。另有两个使用注意RANK不兼容排他exclusive扫描若 groupby 键已排序则 orderby 列也必须排序才能获得正确结果。3.6 收集与合并类make_collect_list_aggregation(null_policy null_handling null_policy::INCLUDE); make_collect_set_aggregation(null_policy null_handling null_policy::INCLUDE, null_equality nulls_equal null_equality::EQUAL, nan_equality nans_equal nan_equality::ALL_EQUAL);COLLECT_LIST把组内元素收集为列表列COLLECT_SET额外去重。两者在null_handling EXCLUDE时会从每个列表行中剔除空值。COLLECT_SET还可通过nulls_equal/nans_equal控制去重时空值与 NaN 是否视为相等默认均视为相等。针对分布式多节点场景libcudf 提供了一组合并部分结果的聚合make_merge_lists_aggregation()把同一键对应的多个列表合并为一个专门用于合并多份COLLECT_LIST的 partial 结果要求输入列表列非空子列不受此限make_merge_sets_aggregation(nulls_equal, nans_equal)先合并列表再逐列表去重用于把多份COLLECT_LIST/COLLECT_SETpartial 结果合并成最终COLLECT_SET实践中 partial 结果应由COLLECT_LIST生成以免提前去重造成不必要开销make_merge_m2_aggregation()合并COUNT_VALID、MEAN、M2三类 partial 结果输出包含合并后三者数值的 structmake_merge_histogram_aggregation()合并多份HISTOGRAMpartial 结果。3.7 窗口函数与序列类make_lag_aggregation(size_type offset)/make_lead_aggregation(size_type offset)分别访问当前行之前 / 之后offset行的值make_ewma_aggregation(double center_of_mass, ewm_history history)指数加权移动平均。center_of_mass决定历史值对当前值的贡献权重ewm_history决定首值的数学处理——INFINITE假设首值之前存在无限历史FINITE则把首值当作唯一已知数据点。EWMA 还有特殊的空值语义空值会把最近有效值向前传播空值不影响均值同时在计算时仍计为一个有效周期例如序列{1, NULL, 3}计算y_2时y_0被视为相隔两个周期。3.8 双列统计类协方差 / 相关系数make_covariance_aggregation(size_type min_periods 1, size_type ddof 1); make_correlation_aggregation(correlation_type type, size_type min_periods 1);两者均以非空 struct 列的两个子列为输入。min_periods是产出结果所需的最少非空观测数ddof默认 1计算除数N - ddof。相关系数通过correlation_type选择PEARSON/KENDALL/SPEARMAN。3.9 t-digest 与位运算make_tdigest_aggregation(int max_centroids 1000); // 从输入值构建 t-digest make_merge_tdigest_aggregation(int max_centroids 1000); // 合并多个 t-digest make_bitwise_aggregation(bitwise_op op); // AND / OR / XORmax_centroids控制压缩级别与后续查询精度是对 t-digest 大小的上界设为 1000 时每个 t-digest 至多包含 1000 个质心每个 32 字节值越大精度越高。t-digest 输出列为 struct 结构{ list{struct{double mean; double weight;}, ...}, double min, double max }其中min/max来自输入流而非质心用于近似分位数计算两端时的边界。BITWISE_AGG仅支持整型列。3.10 主机 UDFmake_host_udf_aggregation(std::unique_ptrhost_udf_base host_udf);HOST_UDF接受一个派生自host_udf_base的实例完成自定义聚合逻辑是扩展聚合能力的钩子。最后aggregation.hpp 提供类型支持校验函数bool is_valid_aggregation(data_type source, aggregation::Kind kind);它返回给定源数据类型是否支持某种聚合适合在运行时动态分发前做检查。四、归约与扫描aggregation_reductionreduce/segmented_reduce/scanreduction.hpp对应 aggregation_reduction.rst提供三条面向整列 / 分段 / 前缀的计算路径全部接受reduce_aggregation、segmented_reduce_aggregation、scan_aggregation类型的聚合对象。4.1cudf::reduce整列归约std::unique_ptrscalar reduce( column_view const col, reduce_aggregation const agg, data_type output_type, cuda::stream_ref stream cudf::get_default_stream(), rmm::device_async_resource_ref mr cudf::get_current_device_resource_ref()); // 带初始值的重载仅支持 sum / product / min / max / any / all / sum_overflow std::unique_ptrscalar reduce(column_view const col, reduce_aggregation const agg, data_type output_type, std::optionalstd::reference_wrapperscalar const init, ...);核心规则源码注释原文语义除SUM_OVERFLOW外归约不检测溢出当output_type与输入列类型不一致时值可能先提升为int64_t或double计算再转型回output_typeSUM_OVERFLOW是特例对有符号整型或十进制输入做溢出检测返回{sum, overflow_flag}的 struct——溢出时 sum 值未定义以布尔标志为准非算术类型timestamp、string 等仅支持min/max所有空值在操作中被跳过归约失败时输出 scalar 的is_valid() false空输入或全空输入时除少数有明确输出的聚合外一般返回非法 scalarany/all要求输出类型为BOOL8mean/variance/std要求输出类型为浮点。reduction.hpp 给出了完整的聚合-输出类型对照表是编写归约代码的权威参考聚合输出类型支持初始值空输入备注SUM / PRODUCToutput_type是非法 scalar累加进 output_type 变量SUM_OVERFLOWSTRUCT{col.type, BOOL8}是{null, false}{sum, overflow_flag}输入须为有符号整型或十进制SUM_OF_SQUARESoutput_type否非法 scalar—MIN / MAXcol.type是非法 scalar仅支持算术、时间戳、时长、字符串类型ANY / ALLBOOL8是仅 ALL 为 True检测非零元素MEAN / VARIANCE / STDFLOAT32 / FLOAT64否非法 scalaroutput_type 须为浮点MEDIAN / QUANTILEoutput_type否非法 scalaroutput_type 为 FLOAT64 时返回精确值NUNIQUEoutput_type否全空输入返回 1可能处理空行NTH_ELEMENTcol.type否非法 scalar—BITWISE_AGGcol.type否非法 scalar仅支持整型HISTOGRAM / MERGE_HISTOGRAMLIST of col.type否返回空列表—COLLECT_LIST / COLLECT_SETLIST of col.type否返回空列表—TDIGEST / MERGE_TDIGESTSTRUCT否返回空 struct返回 t-digest scalarHOST_UDFoutput_type是非法 scalar自定义 UDF 可忽略 output_type4.2cudf::segmented_reduce分段归约std::unique_ptrcolumn segmented_reduce( column_view const segmented_values, device_spansize_type const offsets, // 长度 num_segments 1 segmented_reduce_aggregation const agg, data_type output_type, null_policy null_handling, ...);offsets定义每个分段的边界第i个分段大小为offsets[i1] - offsets[i]offsets 越界属未定义行为。要点空段对应的结果行为 nullnull_handling INCLUDE时仅当段内全部元素有效结果才有效EXCLUDE时只要段内任一元素有效结果即有效算术类型输入可指定任意算术output_type非算术类型如时间戳必须与输入同类型min/max归约要求输出类型与输入一致any/all要求输出类型为BOOL8非算术输出类型只允许min/max。4.3cudf::scan前缀扫描std::unique_ptrcolumn scan( column_view const input, scan_aggregation const agg, scan_type inclusive, // INCLUSIVE / EXCLUSIVE null_policy null_handling null_policy::EXCLUDE, ...);scan对整列做前缀累积INCLUSIVE表示包含当前元素EXCLUSIVE表示不含。空值默认被跳过但输入第i位为空时输出第i位也为空输入列必须为数值类型否则抛cudf::logic_error。排名聚合RANK正是通过 scan 路径实现的。五、分组聚合aggregation_groupbygroupby类与请求结构groupby.hpp对应 aggregation_groupby.rst提供按键分组的聚合、扫描与移位能力。5.1 请求与结果结构struct aggregation_request { column_view values; // 待聚合的元素 std::vectorstd::unique_ptrgroupby_aggregation aggregations; // 期望的聚合集合 }; struct aggregation_result { std::vectorstd::unique_ptrcolumn results{}; // 每个聚合对应一个结果列 };aggregation_request中values.size()必须等于keys.num_rows()同一请求可携带多个聚合结果列顺序与请求中聚合顺序一致。分组扫描使用类似的scan_request结构。5.2groupby对象构造与aggregateexplicit groupby(table_view const keys, null_policy null_handling null_policy::EXCLUDE, sorted keys_are_sorted sorted::NO, std::vectororder const column_order {}, std::vectornull_order const null_precedence {});null_handling键行含空值时是否纳入分组默认排除keys_are_sorted键已排序时可显著提升性能此时可再传column_order每列升/降序空则默认全升序与null_precedence空值位置空则默认AFTER生命周期警告groupby对象不持有keys的生命周期调用方必须保证keys的 table_view 数据在 groupby 存活期间有效。核心方法是aggregate(std::spanaggregation_request const requests, ...)返回{键表, 结果向量}对keys[i]与所有aggregation_result的第i行对应同一个组。groupby.hpp 给出完整示例输入: keys: {1 2 1 3 1} {1 2 1 4 1} request: values: {3 1 4 9 2}, aggregations: {{SUM}, {MIN}} 结果: keys: {3 1 2} {4 1 2} values: SUM: {9 9 1} MIN: {9 2 1}注意组标签行的顺序是任意的且多次groupby::aggregate调用的组顺序可能不同不能依赖输出顺序。5.3 分组扫描与排序优化groupby::scan(requests)对每个组做累积扫描values[i]与组内values[0..i]聚合结果行顺序同样任意groupby::sort_aggregate(...)与groupby::sort_scan(...)要求请求中指定排序信息基于排序后的组内顺序执行聚合/扫描适合 RANK、LEAD、LAG、NTH_ELEMENT 等依赖行序的聚合groupby::shifts(...)按组做平移配合offsets向量与fill_values填充越界位。5.4 流式分组聚合groupby.hpp 还定义了streaming_groupbyaggregate()接收数据块累积组内状态merge()合并另一个 streaming_groupby 的状态多输入流场景result()输出最终结果column_count()校验列数。它专门服务于无法一次性驻留显存的大规模分组聚合。六、滚动窗口聚合aggregation_rollingrolling_window家族rolling.hpp对应 aggregation_rolling.rst围绕rolling_aggregation提供多种滚动窗口函数。6.1 固定窗口rolling_windowstd::unique_ptrcolumn rolling_window( column_view const input, size_type preceding_window, // 向后前序方向静态窗口大小 size_type following_window, // 向前后续方向静态窗口大小 size_type min_periods, // 窗口内最少有效非空观测数 rolling_aggregation const agg, ...);min_periods语义有效观测数少于min_periods时该元素结果为 null当min_periods 0时返回恒等值——SUM 与 COUNT 返回 0MIN 返回该类型最大值MAX 返回该类型最小值空窗口下的 MEAN 行为未定义。另有带default_outputs列的重载当 LEAD/LAG 的行偏移越出列边界时用该列对应行的默认值代替 null。返回列类型规则源码注释明确COUNT 聚合结果恒为INT32VARIANCE / STD 结果恒为FLOAT64其余算子返回与输入相同的类型。因此做滚动 MEAN 前建议先把低精度整型列转换为FLOAT32/FLOAT64避免精度损失。6.2 窗口边界window_boundsstruct window_bounds { static window_bounds get(size_type value); // 有限边界按天或按行 static window_bounds unbounded(); // 无界边界 };get()构造有限边界unbounded()以size_type最大值表示无界窗口is_unbounded()与value()分别查询边界属性与行/天数值。窗口端点默认包含bounded_closed另提供bounded_open排除端点的强类型包装。6.3 分组滚动窗口grouped_rolling_window(...)分组感知的定长滚动窗口。窗口聚合不能跨组边界输入行需预先按group_keys排序。源码示例场景按user_id分组、对sales_amt列做 3 行窗口当前行 前 2 行 / 后 1 行求和grouped_range_rolling_window(...)基于值范围而非固定行数的分组滚动窗口例如按时间戳范围统计rolling.hpp 以统计超车总数的赛车数据集为示例。其边界同样由window_bounds描述且存在按天与按行两种语义。七、组合实战一条 GPU 聚合流水线将上述 API 组合起来可以构建典型的 libcudf 聚合流水线。以下伪代码演示对sales表按user_id分组求和、再对全表做滚动均值、最后输出总分的完整流程#include cudf/aggregation.hpp #include cudf/groupby.hpp #include cudf/reduction.hpp #include cudf/rolling.hpp // 1. 构造聚合对象工厂函数Base 类型即约束了可用场景 auto sum_agg cudf::make_sum_aggregationcudf::groupby_aggregation(); auto mean_roll cudf::make_mean_aggregationcudf::rolling_aggregation(); // 2. 分组聚合按 user_id 键表对 sales_amt 求和 cudf::groupby gb(keys_view, cudf::null_policy::EXCLUDE); cudf::aggregation_request req{values_view, {std::move(sum_agg)}}; auto [group_keys, results] gb.aggregate({req}); // results[0].results[0] 为各组和 // 3. 滚动均值前 2 行 当前行至少 1 个有效观测 auto rolled cudf::rolling_window(sales_view, /*preceding*/2, /*following*/0, /*min_periods*/1, *mean_roll); // 4. 整列归约求总销售额输出 double auto total cudf::reduce(sales_view, *cudf::make_sum_aggregation(), cudf::data_type(cudf::type_id::FLOAT64));几点工程建议归约/扫描前用cudf::is_valid_aggregation(input.type(), kind)校验类型支持避免运行期抛出std::invalid_argumentgroupby 键已知有序时务必传keys_are_sorted sorted::YES并给出column_order/null_precedence可获得更好的性能涉及 RANK、LEAD/LAG 等依赖行序的聚合优先使用sort_aggregate/sort_scan或预先排序大规模流式数据使用streaming_groupby分布式 partial 结果合并使用MERGE_*系列聚合。八、结语libcudf 的列聚合体系通过aggregation抽象基类 模板工厂函数把归约、分段归约、前缀扫描、groupby、滚动窗口五类计算统一到同一套语义之下。本文所讲的 Kind 枚举、make_*_aggregation工厂参数、reduction 输出类型表、groupby 请求结构和 rolling 的min_periods规则均直接来自 aggregation.hpp、reduction.hpp、groupby.hpp、rolling.hpp 的源码注释。如需查阅各组的完整 Doxygen 文档可依次浏览 aggregation_factories.rst、aggregation_reduction.rst、aggregation_groupby.rst 与 aggregation_rolling.rst它们共同构成 column_aggregation.rst 入口页下的完整聚合 API 文档。赞分享数据分析数据工程机器学习【免费下载链接】cudfcuDF - GPU DataFrame Library项目地址https://gitcode.com/gh_mirrors/cu/cudf点击查看免费下载相关推荐icloudpd 使用教程把 iCloud 照片备份到本地的完整指南icloudpd 使用教程把 iCloud 照片备份到本地的完整指南 icloudpd 是一款跨平台的开源命令行工具可以把你 iCloud 照片库里的照片和数据分析数据工程机器学习免费新手指南把真实地图变成 Minecraft 世界3 个参数配方复刻家园免费新手指南把真实地图变成 Minecraft 世界3 个参数配方复刻家园 Arnis 是一款免费开源的世界生成工具读取 OpenStreetMap 地理数据分析数据工程机器学习cuDF GroupBy 完全指南分组、聚合、apply 与滚动窗口实战cuDF GroupBy 完全指南分组、聚合、apply 与滚动窗口实战 cuDF 是 RAPIDS 生态的 GPU DataFrame 库其 GroupB数据分析数据工程机器学习上一篇抖音无水印下载工具3步轻松保存高清视频的完整指南下一篇阿里Wan2.2开源ComfyUI生态引爆AI视频创作革命创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表