资讯动态

Pathway 表操作(Table Operations)完全指南:列选择、过滤、聚合、连接与数据变形

发布时间:2026/9/8 20:17:26 来源:尧图企业网站定制
Pathway 表操作Table Operations完全指南列选择、过滤、聚合、连接与数据变形【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本指南以 Pathway Live Data Framework 官方文档中 Table Operations Overview 一章为主体系统讲解在流式数据管道中对pw.Table实施的全部基础转换从列的创建/重命名/删除到算术、比较、布尔运算符的逐列应用再到groupby/reduce聚合、多表join、union/concat、单元格更新与列展开。文章同时结合仓库内 Table 实现源码 与 reducers 源码 的注释与实际行为帮助读者在读完本文后能够独立编写、验证并排障自己的 Pathway 数据转换逻辑。Pathway 是一个面向流处理、实时分析、LLM 管道与 RAG 的 Python ETL 框架。其核心抽象pw.Table是一张随输入持续更新的活表对表施加的每个转换都会被框架自动转换为增量incremental计算无需开发者手工管理状态与重算逻辑。理解本文中的表操作即掌握了在 Pathway 中表达任意数据变换的基本语法。快速开始如何在本地验证下面的示例本文所有示例都可以在本地交互式环境中原样运行。Pathway 提供了两组调试用辅助 API在 pathway 源码 中大量 doctest 均基于它们编写pw.debug.table_from_markdown(...)用 Markdown/ASCII 文本直接声明一张内存表便于快速构造测试数据pw.debug.compute_and_print(table, include_idFalse)运行计算并把结果打印到控制台。例如import pathway as pw t pw.debug.table_from_markdown( | colA | colB 1 | valA | 1 2 | valB | 4 ) result t.select(new_colt.colB * 10) pw.debug.compute_and_print(result, include_idFalse)include_idFalse表示省略行 IDid列只打印数据列若需观察每行内部 ID 可去掉该参数。行 ID 是 Pathway 表的核心概念——每条记录对应一个 ID连接、更新、聚合都以 ID 为纽带详见后文。建列与重命名select与赋值运算符在 Pathway 中创建新列使用select与赋值运算符。select接收若干kwargs键为输出列名值为列表达式或字面量返回一张只包含这些新列的表t.select(new_colt.colA t.colB) # 由两列运算生成新列 t.select(new_coldefault value) # 用常量创建新列从源码实现看见 python/pathway/internals/table.pyselect会对每个表达式求值得到一列再组合为与输入表处于同一 universe行集合的新表输出列名由kwargs的键决定。值得注意如果赋值的键是id则等价于对表重新索引源码注释明确写到Assigning to id reindexes the table。select之后表中原有列不会保留因此选择列与重命名列常由它一并完成重命名一列t.select(new_colt.old_col)更直接的方式是使用专门的renamet.rename(new_colt.old_col)它只改名而保留其余列。Pathway 中列引用有三种等价写法按场景自由选择写法含义示例点号t.colA引用t表的colA列t.select(t.colA)下标t[colB]以字符串引用列适合列名含特殊字符t.select(t[colB])pw.this.colC引用当前上下文中的表的列t.select(pw.this.colC)列的选择、剔除与表的引用语义选择列与展开选择select配合星号可以把整张表的所有列原样输出t.select(*pw.this) # 保留全部列pw.this是当前表的占位符其定义位于 python/pathway/internals/thisclass.py在.select()、.filter()等方法体内指代被操作的表本身因此上面的写法等价于把当前表所有列逐一展开。若只想保留大部分列可结合下一步的without使用。剔除列withoutt.without(t.colA, t.colB) # 移除这两列保留其余列without接受要删除的列引用或列名字符串与select的方向正好相反适合列数较多、只需去掉少量列的场景。引用当前表与 join/window 中的左右表引用当前表用pw.thist1.select(new_colpw.this.colA pw.this.colB)在连接join或窗口window运算中同时涉及两张表时用pw.left与pw.right分别引用左表与右表t1.join(t2, pw.left.colA pw.right.colB).reduce(*pw.left, pw.right.colC)pw.left/pw.right是 join 上下文中取代pw.this的专用引用。上例先按t1.colA t2.colB匹配两表行再在reduce中把左表所有列展开、同时保留右表colC列。参考索引Reference Indexingix_refPathway 允许引用其他表的行在列中存放指向另一表行的指针运行时按需取回该行的字段t_selected_ids.select(selectedt.ix_ref(column).name)ix_ref把column存放目标表行 ID 的列转换为指向目标表的引用随后可通过.name等语法读取被引用行的属性。这是免 join 的高效关联手段完整机制见专门的 Indexing Grouped Tables 手册 与 Indexes 手册。重索引with_id_from把表的行 ID 换成另一张表t_new_ids中new_id_source列给出的 IDt.with_id_from(t_new_ids.new_id_source)运算符在列上做逐行计算Pathway 支持在列表达式上直接使用 Python 运算符它们会被框架翻译为逐行的增量计算。运算符统一在select等变换中构造新列。算术运算符含义运算符示例加法t.select(new_colt.colA t.colB)减法-t.select(new_colt.colA - t.colB)乘法*t.select(new_colt.colA * t.colB)除法/t.select(new_colt.colA / t.colB)整除向下取整//t.select(new_colt.colA // t.colB)取模%t.select(new_colt.colA % t.colB)幂运算**t.select(new_colt.colA ** t.colB)例如对销售流水做含税金额计算t.select(totalt.price * (1 t.tax_rate))即为合法表达式。比较运算符比较运算产生布尔列常配合filter使用含义运算符示例等于t.select(new_colt.colA t.colB)不等于!t.select(new_colt.colA ! t.colB)大于t.select(new_colt.colA t.colB)小于t.select(new_colt.colA t.colB)大于等于t.select(new_colt.colA t.colB)小于等于t.select(new_colt.colA t.colB)布尔运算符布尔列可进一步组合。注意 Pathway 借用位运算符表示逻辑运算因为 Python 的关键字and/or无法被重载到列表达式上含义运算符示例与Andt.select(new_colt.colA t.colB)或Or\|t.select(new_colt.colA \| t.colB)非Not~t.select(new_col~t.colA)异或XOR^t.select(new_colt.colA ^ t.colB)建议用括号把布尔组合括起来以保证可读性例如t.select(flag(t.colA 10) (t.colB 5))。过滤行filterfilter按一个布尔表达式保留满足条件的行t.filter(~pw.this.column) # 保留 column 为 False/None 之外的行等价于非空过滤 t.filter(pw.this.column value) # 保留 column 严格大于 value 的行从实现细节看python/pathway/internals/table.pyfilter会先对过滤表达式做类型检查其类型必须是bool否则直接抛出TypeError提示Filter argument of Table.filter() has to be bool。此外结果表的 schema 与输入表一致行 ID 集合是self.id的子集若过滤表达式就是对某一列的非 None 判断如t.filter(pw.this.column)框架会从类型层面将该列从Optional[T]自动收缩为T源码中以unoptionalize更新列类型后续使用该列时无需再手工处理空值。如需把一个表按条件拆成两半可使用与filter相邻的splitpositive, negative t.split(pw.this.flag)返回两个行集合互斥的互补表常在需要按规则分流下游处理的场景使用。处理缺失值pw.coalesce在流式数据中空值None十分常见。pw.coalesce返回给定参数中第一个非 None 的值用于为列提供回退值或合并多个候选列t.select(new_colpw.coalesce(t.colA, t.colB)) # colA 为空则取 colB t.select(new_colpw.coalesce(t.colA, 10)) # colA 为空则使用常量 10注意第二个示例的括号需完整闭合pw.coalesce(t.colA, 10)常量为兜底值。这在清洗上游脏数据、给缺失字段补默认值等场景非常实用。聚合groupby与reduce把表按某些列的值分组、再对组内其他列做汇总是数据分析最常用的操作。Pathway 用两个操作符组合完成t.groupby(pw.this.column).reduce(sumpw.reducers.sum(pw.this.value))groupby按一个或多个列引用分组返回一个中间类型GroupedTable参数必须是ColumnReference若误传字符串会得到带提示的ValueErrorDid you mean table.arg ?。它还支持三个可选的扩展参数id指定结果行 ID 取自哪一列、sort_by作为某些 reducers 的排序键以及instance将数据划分成相互独立的实例分别聚合。reduce对每个分组执行聚合并产出一张普通表。源码实现表明t.reduce(...)本质上就是t.groupby().reduce(...)的简写——即不写分组列时把整张表归并成一行这对求全局汇总值很方便。Pathway 的 groupby/reduce 是持续增量的输入流不断到达时各组聚合结果随之更新框架负责维护组状态这正是流式实时聚合与批处理GROUP BY的本质区别。完整的语法细节、多列分组与交互式示例见专门的 Groupby Reduce 手册。内置 Reducers 全览Pathway 在 python/pathway/internals/reducers.py 中实现了丰富的内置 reducer全部通过pw.reducers.xxx访问。下面的语义描述直接对应各函数 docstring 中可验证的实际行为Reducer语义示例any返回组内任意一个聚合值对同一组同时作用于多个列时取值来自同一行保持一致t.groupby(t.colA).reduce(col_anypw.reducers.any(t.colB))argmax返回使arg取得最大值的那个行 ID默认或指定列的值t.groupby(t.colA).reduce(col_argmaxpw.reducers.argmax(t.colB))argmin返回使arg取得最小值的行 ID 或指定列值t.groupby(t.colA).reduce(col_argminpw.reducers.argmin(t.colB))avg组内平均值实现为sum(expr) / count(expr)t.groupby(t.colA).reduce(col_avgpw.reducers.avg(t.colB))earliest处理时间processing time最早的聚合值t.groupby(t.colA).reduce(col_minpw.reducers.earliest(t.colB))latest处理时间最晚的聚合值t.groupby(t.colA).reduce(col_maxpw.reducers.latest(t.colB))max组内最大值t.groupby(t.colA).reduce(col_maxpw.reducers.max(t.colB))min组内最小值t.groupby(t.colA).reduce(col_minpw.reducers.min(t.colB))ndarray收集全部聚合值为一个 numpy 数组skip_nonesTrue可剔除 Nonet.groupby(t.colA).reduce(col_arraypw.reducers.ndarray(t.colB))sorted_tuple收集全部聚合值为一个排序后的元组t.groupby(t.colA).reduce(col_tuplepw.reducers.sorted_tuple(t.colB))sum求和支持 int、float 与数组int 按 64 位累加t.groupby(t.colA).reduce(col_sumpw.reducers.sum(t.colB))tuple收集全部聚合值为元组多列之间的元素次序一致t.groupby(t.colA).reduce(col_tuplepw.reducers.tuple(t.colB))unique组内值全部相同才返回该值否则运行时抛异常t.groupby(t.colA).reduce(col_uniquepw.reducers.unique(t.colB))几个值得展开的实现要点sum的strict参数对 float 求和时默认strictFalse每个批次在单个浮点上做增量更新内存与速度俱佳但在值频繁更新/删除时可能引入数值不稳定设strictTrue则每个批次都基于组内全部值重新求和更稳定但更慢、更耗内存源码 docstring 有明确说明。这对金融等对精度敏感的流式计算是重要开关。argmin/argmax的第二个参数默认返回极值所在行的内部 ID但传入第二个列引用后将返回该极值行中指定列的值。源码 doctest 给出经典用法reduce(minpw.reducers.argmin(table.age, table.name))直接得到年龄最小者的姓名省去一次按 ID 回表查询。时间类 reducer 的区分min/max按数值大小取极值而earliest/latest按到达系统的处理时间取首末值二者适用于不同业务语义请勿混用。除上述之外同一源码文件中还提供了count、count_distinct与基于 HyperLogLog 的近似的count_distinct_approximateprecision需在 4~18 之间仅适用于 append-only 表可按需查阅 reducers.py。Pathway 还支持自定义有状态 reducer可封装业务专用的增量聚合逻辑参见 Custom Reducers 手册。连接两张表joinjoin把两张表中满足匹配条件的行关联起来从而合并两边列是流式关联分析的基础。标准形态为t1.join(t2, pw.left.column pw.right.column).select(...)在join(...)的第二个参数中写匹配条件相等谓词是最常用形式也可写更复杂的表达式例如加偏置的关联、比较运算。完成匹配后通常紧跟select或reduce决定输出列——此时即可使用pw.left与pw.right引用两侧表的任意列例如orders.join(customers, orders.customer_id customers.id).select( order_idorders.id, customer_namecustomers.name, )更复杂的情形多对多关联、按组内行做关联、输出列的精简与 ID 归属在专门的 Join 手册 中有完整讲解。若需要保留未匹配上的行可查阅框架提供的 left/right/outer 等连接变体。纵向合并Union 与 Concatenation当多张表结构相同、需要把行摞在一起时用 union 或 concat操作语法语义Uniont1 t2返回两表行集的并集Union 并回写t1 t2把t2并入t1就地修改t1引用所指表重索引拼接pw.Table.concat_reindex(t1, t2)相同 schema 的表内容拼接并为每一行分配全新的合成 ID注意两个操作语义上的关键差异union保留参与行各自原有的行 ID而concat_reindex的实现源码注释将其类比为 PySpark 的 union要求所有表schema 完全一致拼接后所有行都被重新索引、生成新的合成 ID从而保证最终表的 ID 两两不同、互不冲突。源码中它正是通过先为每张表with_id_from(...)打上互不重叠的新 ID再调用底层的Table.concat完成的。拼接流式、分批到达的同类数据例如多个 Kafka topic 或文件分片解析结果时concat_reindex是首选——因为它不要求各批次的 ID 预先全局唯一。按行更新单元格update_cells与如果要用一张表的内容去覆盖另一张表对应单元格Pathway 提供update_cells及其运算符别名t.update_cells(t_new) t t_new它的语义见源码 docstring 与实现非常明确结果表保留t左表的全部列与行 ID用t_new中同 ID 行的值覆盖t的对应单元格冲突时以t_new为准前提约束t_new的列集合必须是t列集合的子集且t_new的行 ID 集合必须是t行 ID 集合的子集后者通常需要借助pw.universes.promise_is_subset_of之类机制声明或运行时校验。从实现看若左右两表 ID 集合完全相同框架会发出警告并建议改用with_columns按行更新而非按单元格覆盖。update_cells在增量场景的典型用法是用实时迟到数据修正值去订正主表中已输出的历史记录例如温度传感器按时间窗口先报近似值、后补精确值。把一个数组/可迭代列展开成多行flatten若某列存放的是列表、元组、JSON 数组等可迭代对象flatten可把每个元素摊开为一行实现一行变多行t.flatten(t.col_to_flatten)常见的触发场景包括一个订单携带多个商品明细的嵌套数组、一次抓取返回的 JSON 列表等。flatten 之后通常紧跟select/reduce对元素列做加工或再按元素重新分组聚合。列级函数应用pw.apply、pw.make_tuple与 UDF对每个单元格应用任意 Python 函数、或把多列折叠成一列可用下面的操作需求运算符示例对每列逐格应用一个函数在select中调用pw.applyt.select(new_colpw.apply(func, pw.this.col))把多列折叠成单个元组列pw.make_tuplet.select(new_colpw.make_tuple(t.a, t.b, t.c))当func是普通 Python 函数时pw.apply会在运行时逐格调用并把结果纳入 Pathway 的类型系统要求对返回类型做标注。注意它逐行执行、天然比矢量化的列运算符慢因此能用运算符表达的逻辑优先用运算符pw.apply主要面向无法用内置表达式完成的任意逻辑正则处理、日期格式化、调用第三方库函数等。如果需要把复杂的逐行逻辑包装成可复用、可注册的一等公民变换Pathway 还支持完整定义自己的用户自定义函数UDF包括异步与多返回值形态参见 User-defined Functions 手册 与异步转换 Asynchronous Transformations 章节。一个组合示例把上面的操作串起来下面把本文各操作组合成一条完整、可直接运行的小管道同样使用调试 API 验证import pathway as pw raw pw.debug.table_from_markdown( order_id | customer | category | quantity | unit_price 1 | Alice | book | 2 | 12.5 2 | Bob | book | 1 | 12.5 3 | Alice | toy | 5 | 8.0 4 | Bob | book | 3 | 12.5 5 | Carol | None | 2 | 3.0 ) # 1) 缺失分类兜底2) 计算行小计3) 只保留 book 订单 enriched raw.select( customerraw.customer, categorypw.coalesce(raw.category, unknown), amountraw.quantity * raw.unit_price, ).filter(pw.this.category book) # 4) 按客户聚合小计与订单笔数 summary enriched.groupby(pw.this.customer).reduce( customerpw.this.customer, totalpw.reducers.sum(pw.this.amount), orderspw.reducers.count(), ) # 5) 按金额取最大一单的客户名 top summary.reduce(bestpw.reducers.argmax(summary.total, summary.customer)) pw.debug.compute_and_print(summary, include_idFalse) pw.debug.compute_and_print(top, include_idFalse)这段代码演示了本指南覆盖的常量列与运算建列、coalesce缺失值处理、布尔列filter、groupby().reduce()多 reducer 汇总以及argmax第二参数取极值行字段的用法——它们是 Pathway 流式 ETL 中最常用的原子操作。延伸阅读Groupby Reduce 手册分组聚合的完整教程Join 手册join 的详细语法与内部语义Indexing Grouped Tablesix_ref参考索引的机制Custom Reducers编写自定义有状态 reducerUser-defined Functions注册与使用 UDFIterate循环迭代变换当单个转换不够、需要循环直到收敛时的扩展操作。此外在 Pathway 的select/filter/join/groupby之上还有完整的 SQL 接口与更偏数据科学的窗口操作感兴趣的读者可从 SQL 手册 入口继续阅读开发者文档目录 docs/2.developers/4.user-guide 下的其他章节。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

免费获取报价