Chapter 49: Functional Abstraction
函数式抽象在 Python 标准库里承担一类很具体的任务:把“逐个元素处理数据”和“改造可调用对象”变成可组合的 runtime 对象。读完本章后,读者应能定位一段 itertools / functools 代码背后的 iterator 对象、callable 包装层、消费时机、异常传播和资源生命周期,并能判断一条流水线在内存、可重复消费、调试和封装边界上的后果。
本章的贯穿材料是一段日志事件处理代码。它把多个页面中的行对象摊平,转换成事件,筛出购买事件,再按下游需求逐步取值。表面上这段代码像是普通函数拼接,runtime 里真正流动的是 iterator、generator、partial object、wrapper function 和普通 Python object。
from functools import partial
from itertools import chain, groupby, islice
pages = [
[
{"user": "u1", "kind": "view", "amount": "0"},
{"user": "u1", "kind": "buy", "amount": "12"},
],
[
{"user": "u2", "kind": "buy", "amount": "7"},
{"user": "u1", "kind": "buy", "amount": "5"},
],
]
def parse_event(row, *, source):
return {
"user": row["user"],
"kind": row["kind"],
"amount": int(row["amount"]),
"source": source,
}
parse_from_log = partial(parse_event, source="log")
def purchase_stream(source_pages):
rows = chain.from_iterable(source_pages)
events = (parse_from_log(row) for row in rows)
purchases = (event for event in events if event["kind"] == "buy")
return islice(purchases, 3)
这段代码的主干判断是:上游数据在 next() 请求到达之前通常停留在原输入对象中;每一层函数式工具只保存推进下一步所需的最小状态;真正的成本来自对象分配、callable 包装、动态调用、缓存表、排序物化、共享 iterator 和资源关闭策略。标准库文档把 itertools 描述为用于高效循环的 iterator building blocks,把 functools 描述为对 callable object 的高阶函数和操作,本章把这两个模块放回 Python runtime 路径中解释。Python 3.14 itertools 文档 和 Python 3.14 functools 文档 是本章版本边界的主要依据。
49.1 itertools
itertools 的核心对象是 iterator。iterator 是一个实现迭代协议的对象,调用 iter(obj) 获得 iterator,反复调用 next(iterator) 产出下一个值,耗尽时用 StopIteration 表示正常结束。itertools 中的大部分工具返回新的 iterator 对象,这些对象把“如何推进上游”和“如何产出当前值”封装在自身状态里。
在贯穿材料里,chain.from_iterable(source_pages) 接收一个 iterable,其中每个元素又是一个页面列表。它自身返回一个 chain iterator。下游第一次请求值时,chain iterator 先取得第一个页面,再从该页面取得第一行。当前页面耗尽后,它推进到下一个页面。这里没有提前把所有页面行复制成一个大列表,chain iterator 保存的是当前外层 iterator、当前内层 iterator 和推进位置。
itertools 可以按状态形态粗分为几类。组合类工具如 chain、zip_longest 把多个输入流合成一个输出流;截断类工具如 islice、takewhile、dropwhile 根据计数或条件控制消费边界;分组类工具如 groupby 在相邻 key 变化时切换 group;组合数学类工具如 product、permutations、combinations 生成结构化元组。这个分类有助于判断成本来源:有的工具只保存几个游标,有的工具需要缓存输入池,有的工具会共享底层 iterator。
下面的代码展示 chain 和 islice 在消费时推进上游。trace_rows() 只有在下游 list() 开始拉取时才打印消息。
from itertools import chain, islice
def trace_rows():
for row in ["r1", "r2", "r3", "r4"]:
print("produce", row)
yield row
limited = islice(chain(["head"], trace_rows()), 3)
print(list(limited))
执行顺序体现了 pull 模型:list(limited) 请求第一个值,chain 先产出 "head";第二次请求才进入 trace_rows();islice 满足 3 个元素后停止向上游请求。这个结论决定了流水线的内存形态,也决定了副作用的发生位置。上游函数中打印日志、读取文件、解析 JSON 或访问网络,都可能被推迟到下游真正消费 iterator 的时刻。
groupby 是 itertools 中最容易误判的工具之一。它按相邻元素的 key 变化生成 group,group iterator 和外层 groupby 共享同一个底层 iterator。对日志事件做用户聚合时,必须先保证同一用户的事件在输入中相邻;常见做法是先按 user 排序,然后再 group。
from itertools import groupby
def totals_by_user(events):
key = lambda event: event["user"]
ordered = sorted(events, key=key)
for user, items in groupby(ordered, key):
yield user, sum(event["amount"] for event in items)
print(list(totals_by_user(purchase_stream(pages))))
这里的 sorted() 是明确的物化边界。前面的 purchase_stream(pages) 可以惰性推进,进入 sorted() 后,排序需要先消耗全部输入并构造列表。随后 groupby 再按相邻 key 产生 group。判断 itertools 流水线时,需要先找出这类物化点,因为它们直接改变内存峰值、错误触发时机和上游资源持有时间。
49.2 functools
functools 的主对象是 callable。callable 是可以被调用表达式执行的对象,包括 function、bound method、实现 __call__ 的实例、partial object、decorator 返回的 wrapper,以及 singledispatch 生成的 generic function。functools 的功能集中在两类动作上:改造 callable 的调用入口,或者在 callable 周围添加缓存、分发和元数据保留。
贯穿材料中的 parse_from_log = partial(parse_event, source="log") 创建了一个 partial object。它保存原始函数 parse_event、预绑定的位置参数和关键字参数。调用 parse_from_log(row) 时,partial object 会把新传入的 row 与保存的 source="log" 合并,再调用原始函数。它没有复制函数体,也没有创建新的语法级函数定义;它创建的是一个携带绑定状态的 callable object。
lru_cache 是另一类常见包装。它接收一个函数,返回一个带缓存表的 wrapper。wrapper 在调用时先根据参数构造 cache key,命中时直接返回旧结果,未命中时调用原函数并写入缓存。成本来自 key 构造、字典查找、结果保存和失效策略;收益来自重复调用时减少原函数工作量。这个判断适合放在纯函数、昂贵计算、有限参数空间这类场景中使用。
from functools import lru_cache
@lru_cache(maxsize=128)
def normalize_user(user):
return user.strip().lower()
def normalized_purchase_stream(source_pages):
for event in purchase_stream(source_pages):
event = dict(event)
event["user"] = normalize_user(event["user"])
yield event
这段代码里,normalize_user() 的缓存表绑定在 wrapper 上。每次调用 wrapper 时,runtime 先进入 wrapper 逻辑,再根据参数决定是否进入原始函数。缓存返回的是同一个结果对象引用;当结果对象是可变对象时,调用方修改结果会影响后续命中返回值。因此缓存对象的可变性也是函数式抽象的工程边界。
cached_property 作用在 descriptor 层。它把一个无参方法包装成属性访问:第一次属性读取时调用方法并把结果写入实例属性,后续读取直接从实例字典取值。它适合把昂贵、稳定、与实例状态绑定的计算延迟到第一次访问。它也会改变实例状态,因为第一次读取会写入结果。
wraps 解决 wrapper 的元数据问题。decorator 返回新函数后,外部看到的 callable 已经换成 wrapper;__name__、__qualname__、__doc__、__module__ 和 __wrapped__ 等元数据需要从原函数复制或关联,调试、文档生成和 introspection 才能回到被包装对象。wraps 的意义在于让 runtime 包装层保持可观察身份的连续性。
from functools import wraps
def count_calls(func):
calls = 0
@wraps(func)
def wrapper(*args, **kwargs):
nonlocal calls
calls += 1
return func(*args, **kwargs)
return wrapper
这个 decorator 生成 closure。calls 存在于 cell object 中,wrapper 每次调用都更新同一个 cell。wraps(func) 把原函数元数据复制到 wrapper,并设置 __wrapped__。读这类代码时,先确认外部名字绑定到了 wrapper,再确认 wrapper 保存了哪些 cell、调用了哪个原始 callable、暴露了哪些元数据。
singledispatch 把一个函数转换成基于第一个参数类型分发的 generic function。注册函数被放入分发表,调用时根据第一个参数的 runtime type 和 MRO 查找实现。它适合把“同一操作作用于不同输入类型”的分支从函数内部条件判断迁移到注册表。
from functools import singledispatch
@singledispatch
def event_amount(value):
raise TypeError(f"unsupported event: {type(value).__name__}")
@event_amount.register
def _(event: dict):
return int(event["amount"])
@event_amount.register
def _(amount: int):
return amount
singledispatch 的 runtime 边界很明确:它只按第一个参数分发;annotation 或显式注册类型决定注册入口;真正调用哪个实现取决于 runtime type。它不会把整个函数变成静态重载系统,也不会检查所有参数类型。工程上可以用它压平局部类型分支,过度使用会把控制流分散到多个注册点,增加排查成本。
49.3 lazy evaluation
惰性求值在本章中指一种执行策略:表达式或工具先返回可迭代对象,数据生产推迟到下游调用 next() 时发生。generator expression、generator function、map、filter 和许多 itertools 工具都遵循这种形态。它的直接影响是内存峰值下降、错误触发推迟、副作用顺序跟随消费顺序。
贯穿材料里,events = (parse_from_log(row) for row in rows) 创建 generator object。创建时不会调用 parse_from_log(),也不会读取任何 row。第一次对 events 调用 next() 时,generator frame 才开始执行:先向 rows 请求一行,再调用 partial object,再把返回的 dict yield 给下游。
下面的例子把错误触发时间写得更直观。
values = (10 // number for number in [5, 2, 0])
print(next(values))
print(next(values))
print(next(values))
values 创建时没有执行除法。前两次 next(values) 分别返回 2 和 5。第三次请求才会在 10 // 0 处抛出 ZeroDivisionError。因此调试惰性流水线时,异常位置需要沿消费路径往上追踪:先看哪个下游触发了 next(),再看当前 generator frame 执行到哪个表达式,最后看上游对象在该时刻返回了什么。
惰性求值降低内存峰值的前提是中间步骤保持逐项传递。islice(purchases, 3) 只请求最多 3 个购买事件;如果前面有 100 万行日志,上游也只推进到找到 3 个购买事件所需的位置。这个收益在输入巨大、只需前缀、处理可以逐项完成时很稳定。
惰性求值也会延长上游资源的生命周期。文件对象、数据库游标、网络响应和锁如果被 generator 持有,资源释放会跟随 generator 消费完成、显式关闭或垃圾回收时机。工程写法应把资源所有权放进可见的 context manager,或者让调用方承担关闭责任。
def read_lines(path):
with open(path, encoding="utf-8") as file:
for line in file:
yield line.rstrip("\n")
这段 generator 在进入第一次 next() 后打开文件,在 generator 正常耗尽或被关闭时退出 with。调用方只消费一部分后长期保留 generator,文件也可能保持打开状态。设计惰性接口时,需要说明调用方是否会完全消费、是否会调用 close()、是否由更外层 context manager 管理资源。
49.4 iterator pipeline
iterator pipeline 是一串 iterator 对象组成的数据通道。每一层暴露相同协议:接收上游 iterable,返回下游 iterable;真正推进由最末端消费者发起。这个模型让代码可以逐步组合,也让状态分散在多个 iterator 内部。
贯穿材料的 pipeline 可以写成 pages → chain → generator parse → generator filter → islice → consumer。每一层只知道相邻两端:chain 知道如何从页面中取行;parse generator 知道如何把行变成事件;filter generator 知道如何筛出购买事件;islice 知道计数上限;consumer 知道需要多少结果。理解这条链时,先画出相邻对象关系,再沿 next() 方向追踪一次请求。
这张图的重点是方向:数据值从左到右返回,控制请求从右到左发出。consumer 调用 next(Slice),Slice 决定是否继续请求 Filter,Filter 请求 Parse,Parse 请求 Chain,Chain 才访问 pages。异常也沿调用栈返回给 consumer;如果 parse 阶段抛出 ValueError,下游看到的异常发生在消费点。
一次性消费是 iterator pipeline 的基本边界。一个 iterator 保存推进位置,消费后位置改变;再次遍历同一个 iterator 会从当前位置继续,耗尽后不再产出。列表、元组、字符串这类 container 可以反复创建新 iterator;generator、文件对象、map、filter、大部分 itertools 返回值通常是单次推进对象。
stream = purchase_stream(pages)
print(list(stream))
print(list(stream))
第一行 list(stream) 消耗 stream。第二行再次转换同一个 stream,只能从剩余位置继续;如果第一行已经耗尽,第二行返回空列表。排查流水线缺数据时,先检查对象是 reusable iterable 还是 iterator,再检查是否被日志、调试打印、list()、sum()、any()、all() 或 sorted() 提前消费。
短路操作会改变上游消费量。any() 找到真值后停止,all() 遇到假值后停止,islice() 达到数量后停止,takewhile() 条件失败后停止。短路可以减少工作量,也会让上游剩余数据停留在未消费状态。后续继续使用同一个上游 iterator 时,会从短路停止的位置继续。
49.5 higher-order function
高阶函数在 Python runtime 中的含义是:函数或 callable object 可以像普通对象一样被传递、保存、返回和调用。它依赖对象引用语义、call protocol、closure、descriptor binding 和普通名字绑定。它的工程价值是把变化的行为抽出来,交给调用方或注册表注入。
贯穿材料中的 partial(parse_event, source="log") 是高阶用法:partial 接收函数对象,返回新的 callable object。sorted(events, key=key) 也是高阶用法:sorted 接收 key function,并在排序过程中多次调用它。groupby(ordered, key) 继续复用同一个 key function,以相同标准识别相邻分组。
高阶函数的成本来自多层间接调用。直接调用 parse_event(row, source="log") 只有一次 Python 函数调用;通过 partial 调用会先进入 partial object,再进入原函数;通过 decorator 调用会先进入 wrapper,再进入原函数;通过 bound method 调用还会涉及 descriptor 生成的绑定对象或解释器层面的 method binding 优化。通常这类成本小于 I/O 和复杂计算,大量小对象、高频调用、内层循环中会变得可观察。
closure 是高阶函数最常见的状态承载方式。外层函数返回内层函数时,内层函数引用的外层局部变量会存入 cell object。这个 cell 和函数对象一起存活,多个 wrapper 也可以共享同一个 cell。
def amount_filter(limit):
def keep(event):
return event["amount"] >= limit
return keep
large_purchase = amount_filter(10)
selected = (event for event in purchase_stream(pages) if large_purchase(event))
large_purchase 指向内层函数 keep。limit 已经离开外层 frame 的普通局部变量区域,但它被 closure cell 保存下来。调用 large_purchase(event) 时,runtime 从 cell 中读取 limit。读高阶函数代码时,检查顺序是:外部名字绑定到哪个 callable;callable 保存了哪些参数、cell 或实例状态;调用时进入了几层 wrapper;最终业务逻辑在哪个函数体中执行。
高阶函数也会影响调试可读性。行为被作为对象传递后,错误栈上可能出现 wrapper、<lambda>、_ 这类低信息名字。工程代码应给关键 callable 明确命名,在 decorator 上使用 wraps 保留元数据,在注册式分发中把注册点集中到可搜索位置。
49.6 callable adaptation
callable adaptation 指把不同来源的行为对象整理成统一的调用入口。Python 的调用表达式 obj(...) 会先确认对象可调用,再把参数传入对象的 call path。function、method、class、实现 __call__ 的实例、partial、decorator wrapper 和 singledispatch generic function 都可以进入这条路径。
这种统一入口建立在对象模型上。普通函数对象实现调用能力;函数放在类上作为 descriptor 访问时,会生成 bound method,把实例作为第一个参数绑定进去;实现 __call__ 的实例把调用转发到实例方法;partial 保存原 callable 和预绑定参数;decorator wrapper 是普通函数对象;singledispatch 返回带 register 能力的 callable。不同对象的外观统一,内部状态来源不同。
class AmountParser:
def __init__(self, default=0):
self.default = default
def __call__(self, row):
value = row.get("amount")
if value is None:
return self.default
return int(value)
parse_amount = AmountParser(default=0)
amounts = (parse_amount(row) for page in pages for row in page)
parse_amount 是一个实例,但调用表达式会进入 AmountParser.__call__。实例状态 default 成为 callable 行为的一部分。和 closure 相比,这种写法把状态放进实例属性;和 partial 相比,它可以同时保存多项状态并暴露更多方法。选择哪种 adapter,取决于状态规模、命名需求、可测试性和 introspection 需求。
callable adaptation 的边界在于参数契约。一个 pipeline 中所有 callable 都要接受上游传来的对象形态,并返回下游期望的对象形态。map(func, iterable) 假设 func 接收一个元素;starmap(func, iterable) 假设每个元素可以展开成参数序列;sorted(key=func) 假设 key function 接收一个元素并返回可比较 key;singledispatch 假设第一个参数的 runtime type 足以决定实现。
排查 callable adaptation 的错误时,按固定顺序看:第一,当前变量绑定到哪个 callable object;第二,调用表达式传入哪些 positional 和 keyword 参数;第三,adapter 是否预先绑定了参数或实例;第四,wrapper 是否改变返回值、异常或元数据;第五,最终业务函数是否收到预期对象。这个顺序能把 TypeError、错误分发、缓存污染和 wrapper 吞异常等问题拆开。
49.7 generator composition
generator composition 是把多个 generator 或 iterable 连接成一个惰性数据流。常见写法包括 generator expression、嵌套 generator function、yield from、itertools.chain 和 chain.from_iterable。它们的共同目标是让组合层少保存数据,按需把下游请求转交给上游。
贯穿材料中,chain.from_iterable(source_pages) 和 for row in page 的嵌套循环等价于把多个页面连接成一条行流。yield from 可以把这种转发写成 generator function:
def iter_rows(source_pages):
for page in source_pages:
yield from page
def iter_purchases(source_pages):
for row in iter_rows(source_pages):
event = parse_from_log(row)
if event["kind"] == "buy":
yield event
yield from page 把当前 generator 的产出委托给 page 的 iterator。下游请求值时,iter_rows 先从当前 page 产出所有行,再切换到下一个 page。和 chain.from_iterable 相比,这种写法更容易插入日志、异常处理和资源管理;和手写 for item in page: yield item 相比,它表达了清晰的委托关系。
generator expression 适合局部转换和过滤。它创建 generator object,并捕获表达式中需要的外部名字。多个 generator expression 串联时,每层 generator 都有自己的 frame 和暂停点。暂停点意味着局部变量、异常状态和上游 iterator 引用会随 generator 存活。
yield from 在 generator protocol 中还承担更完整的委托语义:send()、throw()、close() 可以沿委托链传递给子 generator,子 generator 的返回值也可以被外层接收。普通数据处理代码常用到的是值转发和关闭传播;协程风格代码会接触更完整的协议边界。
composition 还要处理异常边界。上游 generator 抛出的异常会穿过组合层交给消费者;组合层可以捕获、补充上下文、转换异常或执行清理。工程上应让异常携带足够定位信息,例如当前文件、当前 page、当前 row 编号,从而让最末端的转换失败也能对应到具体输入位置。
49.8 stream processing
stream processing 把数据视为持续到来的序列。Python 标准库中的 iterator pipeline 是同步、pull-based 的 stream 模型:下游每次 next() 请求形成需求,上游根据需求生产一个值。这个模型天然支持背压,因为下游停止请求时,上游同步停止推进。
贯穿材料中的 islice(purchases, 3) 就是一个背压点。consumer 只要 3 个购买事件,islice 达到数量后停止请求,filter、parse、chain 和 pages 都随之停止推进。对于文件、socket、分页 API、压缩流等输入,这种停止请求的能力可以把 CPU、内存和 I/O 使用限制在下游实际需求范围内。
stream processing 的第一条边界是一次性消费。stream 经常绑定外部资源或内部游标,天然适合“从头到尾处理一次”。需要重复分析时,应主动选择物化策略,例如把结果写入 list、临时文件、数据库表或缓存层。物化点应清楚命名,因为它把惰性流变成可重复读取的数据结构。
stream processing 的第二条边界是 cleanup。同步 iterator 没有自动知道“业务已经不再需要我”的通用信号;generator 支持 close(),context manager 支持确定退出,文件对象支持关闭。把外部资源放入 stream 时,应设计清楚关闭路径。
from contextlib import closing
def first_purchase(source_pages):
with closing(purchase_stream(source_pages)) as stream:
return next(stream, None)
这里 closing() 会在退出 with 时调用 stream 的 close()。对于 generator,这会触发 generator 内部的 finally 或 context manager 退出路径。对于普通 itertools iterator,是否有 close() 取决于对象本身;因此资源管理最好放在拥有资源的 generator 内部,或者放在外层明确的 context manager 中。
stream processing 的第三条边界是异常时的状态恢复。某一行解析失败后,上游 iterator 位置已经推进到该行;重试同一个 iterator 往往无法回到失败前。可靠系统会把 checkpoint、输入 offset、幂等写入和错误记录作为 stream 设计的一部分。标准库 iterator 提供组合原语,业务系统仍然需要定义失败策略。
综合本章,函数式抽象的判断顺序可以固定为六步:先确认输入是 reusable iterable 还是单次 iterator;再确认每一层返回的对象类型和保存状态;然后沿 next() 追踪消费方向;接着标出 sorted()、list()、缓存和分组等物化点;随后检查 callable adapter 的绑定参数、closure、descriptor 和 wrapper;最后检查异常、短路和资源关闭路径。这个顺序能把简洁的函数式代码还原成可观察、可调试的 runtime 结构。
最小自检任务
阅读下面代码,判断两次打印分别输出什么,并说明 parse() 在什么时刻执行、stream 是否可以重复消费、partial 在调用路径里保存了什么状态。
from functools import partial
from itertools import islice
rows = [
{"kind": "view", "amount": "0"},
{"kind": "buy", "amount": "3"},
{"kind": "buy", "amount": "4"},
]
def parse(row, *, source):
print("parse", row["amount"])
return {"kind": row["kind"], "amount": int(row["amount"]), "source": source}
parse_log = partial(parse, source="log")
stream = (parse_log(row) for row in rows)
purchases = (event for event in stream if event["kind"] == "buy")
first = islice(purchases, 1)
print(list(first))
print(list(purchases))
答案要点
第一次 list(first) 会从 first 向上游发起请求。islice 需要 1 个购买事件,于是 filter generator 推进 stream。stream 先解析第一行,打印 parse 0,得到 view 事件后被过滤;随后解析第二行,打印 parse 3,得到 buy 事件后返回给 islice。第一次打印的列表包含第二行转换后的事件:[{"kind": "buy", "amount": 3, "source": "log"}]。
第二次 list(purchases) 继续使用同一个 purchases generator。前两行已经被消费,generator 从第三行继续推进,打印 parse 4,返回第三行转换后的 buy 事件。第二次打印的列表是 [{"kind": "buy", "amount": 4, "source": "log"}]。
parse() 在 generator 被下游消费时执行,创建 stream、purchases 和 first 时都不会调用它。stream 和 purchases 是单次推进对象,消费位置会保存在 generator 状态中。partial 保存了原始 callable parse 和预绑定关键字参数 source="log";每次 parse_log(row) 被调用时,它把当前 row 与保存的 source 合并后调用 parse()。
本章知识点总结
- iterator 对象:
itertools工具通常返回 iterator,并把当前上游、当前位置和必要辅助状态保存在对象内部。 - 惰性推进:generator expression 和多数 iterator 工具把数据生产推迟到下游调用
next()的时刻。 - 消费方向:iterator pipeline 的控制请求从消费者向上游传递,数据值沿相反方向返回。
- 物化边界:
sorted()、list()、缓存表和部分组合工具会把惰性流转换成已保存数据。 - groupby 边界:
groupby按相邻 key 分组,group iterator 与外层对象共享底层 iterator。 - partial 状态:
partial保存原 callable 和预绑定参数,并在调用时合并新参数后转发。 - 缓存包装:
lru_cache在 wrapper 上维护缓存表,收益和风险都取决于参数 key 与返回对象可变性。 - 元数据保留:
wraps让 decorator wrapper 保留原函数可观察身份,支撑调试和 introspection。 - 类型分发:
singledispatch基于第一个参数的 runtime type 选择注册实现。 - 高阶成本:closure、partial、bound method 和 wrapper 会增加间接调用层,内层高频路径需要评估成本。
- callable 统一:function、method、实例、partial 和 wrapper 都能通过调用表达式进入各自 call path。
- 生成器组合:
yield from、generator expression 和chain都能把多段 iterable 连接成惰性数据流。 - 背压模型:同步 iterator pipeline 是 pull-based 模型,下游停止请求时上游同步停止推进。
- 资源关闭:stream 设计需要明确 generator、context manager、文件和外部连接的关闭路径。
- 排查顺序:先看 iterable 类型,再看 iterator 状态,再沿
next()追踪消费、物化、callable adapter、异常和 cleanup。