Skip to main content

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 可以按状态形态粗分为几类。组合类工具如 chainzip_longest 把多个输入流合成一个输出流;截断类工具如 islicetakewhiledropwhile 根据计数或条件控制消费边界;分组类工具如 groupby 在相邻 key 变化时切换 group;组合数学类工具如 productpermutationscombinations 生成结构化元组。这个分类有助于判断成本来源:有的工具只保存几个游标,有的工具需要缓存输入池,有的工具会共享底层 iterator。

下面的代码展示 chainislice 在消费时推进上游。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 的时刻。

groupbyitertools 中最容易误判的工具之一。它按相邻元素的 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、mapfilter 和许多 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) 分别返回 25。第三次请求才会在 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 决定是否继续请求 FilterFilter 请求 ParseParse 请求 ChainChain 才访问 pages。异常也沿调用栈返回给 consumer;如果 parse 阶段抛出 ValueError,下游看到的异常发生在消费点。

一次性消费是 iterator pipeline 的基本边界。一个 iterator 保存推进位置,消费后位置改变;再次遍历同一个 iterator 会从当前位置继续,耗尽后不再产出。列表、元组、字符串这类 container 可以反复创建新 iterator;generator、文件对象、mapfilter、大部分 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 指向内层函数 keeplimit 已经离开外层 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 fromitertools.chainchain.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 推进 streamstream 先解析第一行,打印 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 被下游消费时执行,创建 streampurchasesfirst 时都不会调用它。streampurchases 是单次推进对象,消费位置会保存在 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。