ARTICLE DETAIL

资讯详情

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

PyArrow Compute Functions 完全指南:逐元素运算、分组聚合、连接、表达式过滤与自定义 UDF

PyArrow Compute Functions 完全指南:逐元素运算、分组聚合、连接、表达式过滤与自定义 UDF PyArrow Compute Functions 完全指南逐元素运算、分组聚合、连接、表达式过滤与自定义 UDF【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow本文以 Apache Arrow 官方 Python 文档 compute.rst 为主体系统讲解pyarrow.compute模块的完整能力版图从直接调用计算函数pc.sum、pc.equal等、通过Table.group_by做分组聚合到Table/Dataset的 join 操作、基于表达式Expression的行过滤、实验性的用户自定义函数UDF注册以及数组/标量上的标准 Python 运算符重载。读完本文你将能够针对数组、标量、表乃至数据集选择合适的 compute 入口并能通过函数注册机制把 NumPy 等生态的算法接入 Arrow 计算体系。1. 计算模块的总体结构函数注册表驱动的统一 APIArrow 支持对可能具有不同类型的输入执行逻辑计算操作logical compute operations。标准计算操作由pyarrow.compute模块提供可以直接调用 import pyarrow as pa import pyarrow.compute as pc a pa.array([1, 1, 2, 3]) pc.sum(a) pyarrow.Int64Scalar: 7从源码结构看pyarrow.compute中暴露的每个顶层函数并不是逐一手写的。compute.py 中的_make_global_functions()会遍历 C 侧的全局FunctionRegistryfunction_registry()为每个已注册的 compute 函数动态生成一个带签名、带文档字符串的 Python 包装器hash_aggregate类函数因为不可直接调用见第 3 节而被显式跳过。每个包装器的文档由_decorate_compute_function从 C 侧的函数元数据func._doc中的 summary、description、参数名、选项类拼装而成选项参数还能以关键字参数或options整体两种形式传入参见 compute.py 中的_handle_options。这套动态生成的机制对应着 C 侧的统一 Compute API函数存储在arrow::compute::FunctionRegistry中按名称查找计算输入统一表示为DatumScalar、Array、ChunkedArray 等形状的标签联合。C 侧的函数注册与内核kernel实现位于 cpp/src/arrow/compute/其中 registry.cc 管理注册表kernels/目录存放各函数的具体实现api_scalar.cc、api_aggregate.cc 等提供具体 API。完整函数清单见 C 计算文档中的 Available functions 一节Python 文档中的.. arrow-computefuncs::指令即从该注册表自动生成函数列表。一个实用提示大多数 compute 函数同时支持数组含 chunked与标量输入但部分函数强制要求特定输入形状。例如sort_indices要求其第一个也是唯一的输入必须是数组。2. 标准计算函数逐元素运算、多值返回与表级操作以下示例展示了数组间与标量间的基本调用方式 a pa.array([1, 1, 2, 3]) b pa.array([4, 1, 2, 8]) pc.equal(a, b) pyarrow.lib.BooleanArray object at ... [ false, true, true, false ] x, y pa.scalar(7.8), pa.scalar(9.3) pc.multiply(x, y) pyarrow.DoubleScalar: 72.542.1 多值返回StructScalar如果一个 compute 函数返回多个值结果会以StructScalar形式给出。可以通过调用其values()方法提取各个字段 a pa.array([1, 1, 2, 3]) pc.min_max(a) pyarrow.StructScalar: [(min, 1), (max, 3)] a, b pc.min_max(a).values() a pyarrow.Int64Scalar: 1 b pyarrow.Int64Scalar: 32.2 超越逐元素对表进行排序compute 函数不仅能做元素级运算还能作用于整张表。例如按某一列排序并取回行索引 t pa.table({x:[1,2,3],y:[3,2,1]}) i pc.sort_indices(t, sort_keys[(y, ascending)]) i pyarrow.lib.UInt64Array object at ... [ 2, 1, 0 ]完整的 PyArrow 计算函数参考清单见官方 API 文档中的pyarrow.compute参考compute.rst 末尾指向的api.compute参考页跨语言的 C 函数全集可在 C Compute Functions 文档 中查阅。3. 分组聚合group_by 与 hash_* 聚合函数分组聚合函数grouped aggregation不能像普通函数那样直接调用而必须通过pyarrow.Table.group_by能力使用。group_by会返回一个分组声明grouping declaration在其上应用哈希聚合函数 t pa.table([ ... pa.array([a, a, b, b, c]), ... pa.array([1, 2, 3, 4, 5]), ... ], names[keys, values]) t.group_by(keys).aggregate([(values, sum)]) pyarrow.Table keys: string values_sum: int64 ---- keys: [[a,b,c]] values_sum: [[3,7,5]]上例中传给aggregate的sum聚合底层就是hash_sumcompute 函数。group_by的 Python 实现见 table.pxi 中的Table.group_by方法。3.1 一次执行多个聚合 t pa.table([ ... pa.array([a, a, b, b, c]), ... pa.array([1, 2, 3, 4, 5]), ... ], names[keys, values]) t.group_by(keys).aggregate([ ... (values, sum), ... (keys, count) ... ]) pyarrow.Table keys: string values_sum: int64 keys_count: int64 ---- keys: [[a,b,c]] values_sum: [[3,7,5]] keys_count: [[2,2,1]]3.2 为聚合函数提供选项每个聚合函数都可以提供选项。例如用CountOptions改变 null 值的计数方式 table_with_nulls pa.table([ ... pa.array([a, a, a]), ... pa.array([1, None, None]) ... ], names[keys, values]) table_with_nulls.group_by([keys]).aggregate([ ... (values, count, pc.CountOptions(modeall)) ... ]) pyarrow.Table keys: string values_count: int64 ---- keys: [[a]] values_count: [[3]] table_with_nulls.group_by([keys]).aggregate([ ... (values, count, pc.CountOptions(modeonly_valid)) ... ]) pyarrow.Table keys: string values_count: int64 ---- keys: [[a]] values_count: [[1]]CountOptions的mode对应三种计数语义默认只统计非 null 值、只统计 null 值、或统计全部值。3.3 受支持的分组聚合函数清单所有支持的分组聚合函数都可以在aggregate中带或不带hash_前缀使用。其底层实现在 C 侧注册例如 hash_aggregate.cc 中注册了hash_count、hash_count_all、hash_first、hash_last、hash_min_max、hash_min、hash_max、hash_any、hash_all、hash_count_distinct、hash_distinct、hash_one、hash_list等函数数值类函数另见 hash_aggregate_numeric.cc。结合 C 计算文档 的函数表主要的分组聚合函数及其行为如下函数名输入类型输出类型选项类说明hash_all/hash_anyBooleanBooleanScalarAggregateOptions若skip_nullsfalse则按 Kleene 逻辑处理 nullhash_approximate_medianNumericFloat64ScalarAggregateOptions近似中位数hash_countAnyInt64CountOptionsCountMode 控制是否统计 nullhash_count_all无参Int64—行计数hash_count_distinctAnyInt64CountOptions去重计数hash_distinctAnyList of input typeCountOptions收集组内去重值hash_first/hash_lastNumeric, Binary输入类型ScalarAggregateOptions结果依赖输入数据顺序hash_first_lastNumeric, BinaryStructScalarAggregateOptions首尾值组合返回hash_min/hash_max非嵌套、非 binary/string输入类型ScalarAggregateOptionshash_min_max非嵌套类型StructScalarAggregateOptions{min: 输入类型, max: 输入类型}hash_meanNumericDecimal/Float64ScalarAggregateOptionsdecimal 输入保持精度与 scalehash_product/hash_sumNumericNumericScalarAggregateOptions输出为 Int64/UInt64/Float64 或 Decimal128/256取决于输入hash_skew/hash_stddev/hash_variance/hash_kurtosisNumericFloat64SkewOptions/VarianceOptionsdecimal 参数先转为 Float64hash_listAnyList of input type—将组内值收集为 list 数组hash_oneAny输入类型—每组返回一个任意值偏向非 null 值hash_pivot_widerBinary/String/Integer AnyStructPivotWiderOptions宽表透视hash_tdigestNumericFixedSizeList[Float64]TDigestOptions近似分位数固定内存占用C 文档给出的一个典型例子值得注意对包含 null 键与 null 值的输入做分组求和时null 会被当作一个独立的键值参与分组例如keynull的组单独成组其sum(x)只由非 null 值累加得出。这一点在编写涉及缺失分组键的 ETL 逻辑时需要留意。4. Table 与 Dataset 的 Join 操作pyarrow.Table与pyarrow.dataset.Dataset都通过各自的join方法支持连接操作Python 侧实现见 table.pxi 中的Table.join。方法接受要连接进来的右侧表/数据集以及一个或多个连接键。默认执行left outer join也可以请求以下任意连接类型left semiright semileft antiright antiinnerleft outerright outerfull outer4.1 单键连接只提供表和连接键即可完成基本连接 table1 pa.table({id: [1, 2, 3], ... year: [2020, 2022, 2019]}) table2 pa.table({id: [3, 4], ... n_legs: [5, 100], ... animal: [Brittle stars, Centipede]}) joined_table table1.join(table2, keysid)结果是一张由table1与table2在id键上执行 left outer join 得到的新表 joined_table pyarrow.Table id: int64 year: int64 n_legs: int64 animal: string ---- id: [[3,1,2]] year: [[2019,2020,2022]] n_legs: [[5,null,null]] animal: [[Brittle stars,null,null]]4.2 指定连接类型full outer join通过join_type参数可以请求其他连接类型例如全外连接。注意 join 结果中各列可能来自不同分块示例中用combine_chunks()合并后再sort_by使输出可读 table1.join(table2, keysid, join_typefull outer).combine_chunks().sort_by(id) pyarrow.Table id: int64 year: int64 n_legs: int64 animal: string ---- id: [[1,2,3,4]] year: [[2020,2022,2019,null]] n_legs: [[null,null,5,100]] animal: [[null,null,Brittle stars,Centipede]]4.3 复合键连接可以提供一个以上的连接键使连接发生在两个键上。例如为table2增加year列后按(id, year)连接 table2_withyear table2.append_column(year, pa.array([2019, 2022])) table1.join(table2_withyear, keys[id, year]) pyarrow.Table id: int64 year: int64 n_legs: int64 animal: string ---- id: [[3,1,2]] year: [[2019,2020,2022]] n_legs: [[5,null,null]] animal: [[Brittle stars,null,null]]4.4 Dataset 级别的连接Dataset.join具备同样的能力可以直接把两个数据集连接起来 import pyarrow.dataset as ds ds1 ds.dataset(table1) ds2 ds.dataset(table2) joined_ds ds1.join(ds2, keysid) joined_ds.head(5) pyarrow.Table id: int64 year: int64 n_legs: int64 animal: string ---- id: [[3,1,2]] year: [[2019,2020,2022]] n_legs: [[5,null,null]] animal: [[Brittle stars,null,null]]5. 表达式过滤pc.field、布尔组合与惰性执行Table与Dataset都可以用布尔类型的Expression进行过滤。表达式从pyarrow.compute.field构建起对一个或多个字段施加比较与转换运算即可组合出所需的过滤表达式。大多数 compute 函数都可以用于对field做转换。pc.field与pc.scalar的 Python 实现分别见 compute.py。注意二者与pa.field/pa.scalar的本质区别pyarrow.scalar()创建 Arrow 内存模型中的Scalar对象而pyarrow.compute.scalar()创建的是代表标量值的Expression对象用于计算表达式、谓词和数据集过滤。5.1 用位运算构造偶数过滤器下面构建一个找出列nums中所有偶数的过滤器 even_filter (pc.bit_wise_and(pc.field(nums), pc.scalar(1)) pc.scalar(0))其原理1的二进制是00000001只有末位为1的数即奇数与1做bit_wise_and才会得到非零结果因此num 1 0恰好刻画偶数。构建好过滤器后将其传给Table.filter即可只保留匹配行 table pa.table({nums: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], ... chars: [a, b, c, d, e, f, g, h, i, l]}) table.filter(even_filter) pyarrow.Table nums: int64 chars: string ---- nums: [[2,4,6,8,10]] chars: [[b,d,f,h,l]]5.2 组合过滤器and / or / not多个过滤器可以用、|、~分别表达 and、or、not。例如~even_filter就是筛选所有奇数 table.filter(~even_filter) pyarrow.Table nums: int64 chars: string ---- nums: [[1,3,5,7,9]] chars: [[a,c,e,g,i]]也可以把even_filter与pc.field(nums) 5组合筛选出大于 5 的偶数 table.filter(even_filter (pc.field(nums) 5)) pyarrow.Table nums: int64 chars: string ---- nums: [[6,8,10]] chars: [[f,h,l]]5.3 Dataset 的惰性过滤Dataset同样可以用Dataset.filter方法过滤。该方法返回一个新的Dataset实例过滤器只有在真正访问数据时才被应用惰性求值因此多个filter调用可以链式组合而不产生中间物化 dataset ds.dataset(table) filtered dataset.filter(pc.field(nums) 5).filter(pc.field(nums) 2) filtered.to_table() pyarrow.Table nums: int64 chars: string ---- nums: [[3,4]] chars: [[c,d]]6. 用户自定义函数UDF实验性 API注意官方文档明确标注该 API 为实验性experimental。PyArrow 允许定义并注册自定义 compute 函数。注册后这些函数可以从 Python、C 以及任何包装 Arrow C 的实现如 R 的arrow包按注册名调用。UDF 支持范围限于标量函数scalar function即对数组或标量执行逐元素操作的函数其输出通常不依赖于参数中值的顺序。这类函数大致对应 SQL 表达式中使用的函数或 NumPy 的 universal functions。6.1 注册一个 UDF注册 UDF 需要定义函数名、函数文档、输入类型和输出类型使用pyarrow.compute.register_scalar_function import numpy as np function_name numpy_gcd function_docs { ... summary: Calculates the greatest common divisor, ... description: ... Given x and y find the greatest number that divides\n ... evenly into both x and y. ... } input_types { ... x : pa.int64(), ... y : pa.int64() ... } output_type pa.int64() def to_np(val): ... if isinstance(val, pa.Scalar): ... return val.as_py() ... else: ... return np.array(val) def gcd_numpy(ctx, x, y): ... np_x to_np(x) ... np_y to_np(y) ... return pa.array(np.gcd(np_x, np_y)) pc.register_scalar_function(gcd_numpy, ... function_name, ... function_docs, ... input_types, ... output_type)UDF 实现函数的第一个参数始终是context上例中命名为ctx它是pyarrow.compute.UdfContext的实例。该上下文暴露若干有用属性特别是UdfContext.memory_pool用于在 UDF 内部做内存分配应使用 Arrow 内存池分配而非普通 Python 对象默认分配。6.2 直接调用 UDFcall_function pc.call_function(numpy_gcd, [pa.scalar(27), pa.scalar(63)]) pyarrow.Int64Scalar: 9 pc.call_function(numpy_gcd, [pa.scalar(27), pa.array([81, 12, 5])]) pyarrow.lib.Int64Array object at ... [ 27, 3, 1 ]可见 UDF 同样遵循标量/数组广播语义(scalar, scalar)产出标量(scalar, array)产出数组。6.3 在 Dataset 中使用 UDF更一般地UDF 可以用在任何可以按名称引用 compute 函数的地方。例如通过Expression._call在数据集的列上调用它。考虑数据在表中需要计算某一列与标量 30 的 GCD复用上面注册的numpy_gcd data_table pa.table({category: [A, B, C, D], value: [90, 630, 1827, 2709]}) dataset ds.dataset(data_table) func_args [pc.scalar(30), ds.field(value)] dataset.to_table( ... columns{ ... gcd_value: ds.field()._call(numpy_gcd, func_args), ... value: ds.field(value), ... category: ds.field(category) ... }) pyarrow.Table gcd_value: int64 value: int64 category: string ---- gcd_value: [[30,30,3,3]] value: [[90,630,1827,2709]] category: [[A,B,C,D]]注意ds.field()._call(...)返回的是一个pyarrow.compute.Expression。传给该函数调用的参数都是表达式而非标量值注意pyarrow.scalar与pyarrow.compute.scalar的区别后者产生表达式。该表达式在投影算子projection operator执行时才被求值。投影表达式的限制在上例中我们用表达式为表添加了新列gcd_value。为表添加这种动态计算的新列称为投影projection对投影表达式中可以使用的函数有明确限制投影函数必须为每个输入行恰好输出一个值该输出值应完全由该行本身计算得出不得依赖其他行。因此上面一直使用的numpy_gcd是合法的投影函数而累积求和cumulative sum不合法因为它对某一行的结果依赖此前各行丢弃 null 行drop nulls也不合法因为它对某些行不产生输出值。7. 标准 Python 运算符数组与标量的运算符重载PyArrow 为数组和标量支持标准 Python 运算符的逐元素运算。目前支持范围限于部分标准 compute 函数算术运算、-、/、%、**、位运算、|、^、、及其他。这些运算符尽可能使用底层内核的带检查checked版本并带有相应的约束——例如两个字符串数组不能相加。使用示例 import pyarrow as pa arr pa.array([-1, 2, -3]) val pa.scalar(42.7) arr val pyarrow.lib.DoubleArray object at ... [ 41.7, 44.7, 39.7 ] val ** arr pyarrow.lib.DoubleArray object at ... [ 0.023419203747072598, 1823.2900000000002, 0.000012844475506953143 ] arr 2 pyarrow.lib.Int64Array object at ... [ -4, 8, -12 ]7.1 底层机制隐式类型提升与 common numeric type运算符重载最终落在 compute 内核上其类型行为值得了解。从 C Compute 文档 看当内核与参数类型不完全匹配时函数可能先对参数做隐式转换比较与算术内核要求同类型参数支持通过把参数提升到可容纳任一侧任意值的数值类型来对不同数值类型求值这一类型即common numeric type输入类型公共数值类型说明int32, int32int32int16, int32int32最大宽度 32提升 LHS 至 int32uint16, int32int32一侧有符号覆盖无符号uint32, int32int64加宽以容纳 uint32 的范围uint16, uint32uint32全部无符号保持无符号int16, uint32int64uint64, int16int64int64 无法容纳所有 uint64 值float32, int32float32RHS 提升为 float32float32, float64float64float32, int64float32尽管 int64 更宽仍提升为 float32特别地比较uint64列与int16列时如果某个uint64值无法用公共类型int64表示如2 ** 63可能抛出错误。此外算术函数默认版本不检测溢出结果通常回绕多数函数另有带_checked后缀的溢出检查变体检测到溢出时返回Invalid错误状态——Python 侧运算符正是尽可能使用这些 checked 内核。8. 小结按场景选择计算入口对数组/标量做逐元素运算或整体归约直接用pyarrow.compute顶层函数如pc.sum、pc.equal、pc.min_max多值返回用StructScalar.values()解包按键分组汇总使用Table.group_by(...).aggregate([...])其中sum/count等聚合名对应hash_*compute 函数并可通过CountOptions等选项类微调 null 语义表/数据集连接Table.join/Dataset.join默认 left outer可指定 8 种连接类型与复合键行级过滤用pc.field/pc.scalar构建Expression配合 | ~组合交给Table.filter立即执行或Dataset.filter惰性执行接入自有算法通过register_scalar_function注册实验性 UDF用pc.call_function直接调用或在 Dataset 投影中以表达式形式按名称调用交互式脚本对数值数组/标量可直接使用、-、**、等 Python 运算符底层映射到 checked 内核并遵循 common numeric type 提升规则。以上行为均以当前仓库中 Python Compute 文档、C Compute 文档 与 pyarrow/compute.py 等源码为准UDF 一节为实验性 API生产使用前建议关注上游变更。【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表