下面的函数统一处理普通响应和流式事件响应:
def compact_generate_response(
response: Response | Generator,
) -> FlaskResponse:
"""统一处理普通响应和流式事件响应。"""
if isinstance(response, Response):
return json(response)
def generate() -> Generator:
yield from response
return FlaskResponse(
stream_with_context(generate()),
status=200,
mimetype="text/event-stream",
)函数有两条分支:Response 通过 json() 一次返回;Generator 产生的数据经过 yield from 转发,再由 Flask 以 SSE 响应发送。这里先关注控制流。
读懂这段代码需要拆开五个概念:
Generator为什么能暂停和恢复yield from如何转发另一个迭代器的数据Flask如何把Iterable作为响应体SSE如何在文本流中划分事件stream_with_context()为何与 Flask Request Context 有关
这些概念处于不同层次:
Generator ──> 按需产生数据
│
▼
Flask Response ──> 把 Iterable 作为 HTTP 响应体
│
├── 普通文本流:text/plain
│
└── SSE 事件流:text/event-stream + SSE 消息格式
stream_with_context ──> Generator 执行期间保留 Flask 请求上下文从 Iterable 到 Generator
Python 的 for 循环并不要求对象必须是列表。只要对象能够提供一个迭代器,for 就能逐个取值。
Iterable 与 Iterator
Iterable(可迭代对象)表示“可以开始一次迭代”的对象。列表、元组、字符串和字典都是 Iterable:
numbers = [10, 20, 30]
for number in numbers:
print(number)for 循环会先调用 iter(numbers) 得到 Iterator(迭代器),再反复调用迭代器的 next():
numbers = [10, 20, 30]
iterator = iter(numbers)
print(next(iterator)) # 10
print(next(iterator)) # 20
print(next(iterator)) # 30
print(next(iterator)) # 抛出 StopIteration两者的职责如下:
| 概念 | 核心能力 | 例子 |
|---|---|---|
| Iterable | iter(obj) 能返回 Iterator | list、tuple、str |
| Iterator | next(obj) 能返回下一个值,耗尽时抛出 StopIteration | iter(list) 的结果、Generator |
列表本身不是 Iterator。对列表反复调用 iter(),可以得到相互独立的迭代器。Generator 既是 Iterable,也是 Iterator;iter(generator) 返回它自己,只会从当前位置继续,耗尽后不能重新开始。
from collections.abc import Iterable, Iterator
numbers = [1, 2, 3]
iterator = iter(numbers)
print(isinstance(numbers, Iterable)) # True
print(isinstance(numbers, Iterator)) # False
print(isinstance(iterator, Iterator)) # TrueGenerator 是什么
Generator(生成器)是创建 Iterator 的一种方式。包含 yield 的普通 def 会定义一个生成器函数;调用这个函数得到生成器对象,此时函数体不会执行。
def trace():
print("A")
yield 1
print("B")
yield 2
print("C")
generator = trace()
print(generator)执行到这里时,A、B、C 都没有打印。后续调用 next() 才会推进函数:
print(next(generator)) # 先打印 A,再得到 1
print(next(generator)) # 先打印 B,再得到 2
print(next(generator)) # 打印 C,然后抛出 StopIteration时间线是:
generator = trace()
│
└── 创建生成器对象,函数体尚未执行
next(generator)
├── print("A")
├── yield 1
└── 暂停,保留局部变量和执行位置
next(generator)
├── 从上次暂停处继续
├── print("B")
├── yield 2
└── 再次暂停
next(generator)
├── print("C")
├── 函数结束
└── 抛出 StopIterationfor 循环执行的仍是这套协议,只是它捕获了 StopIteration,所以日常遍历时看不到这个异常:
for value in trace():
print(f"收到 {value}")Python 官方文档把这种行为描述为:每次 yield 暂停执行并保留局部状态,生成器恢复时从暂停位置继续。详细语义见 Yield expressions。
return 与 yield 改变的是时间行为
下面的普通函数先构造列表,再把结果交给调用方:
def build_numbers(limit: int) -> list[int]:
result = []
for number in range(limit):
result.append(number)
return result生成器函数则在调用方索取下一个值时才继续计算:
from collections.abc import Iterator
def generate_numbers(limit: int) -> Iterator[int]:
for number in range(limit):
yield number二者的差异不是最终能否得到相同数字,而是计算、内存占用和消费发生的时间:
return
生成全部数据 ──> 保存在容器中 ──> 返回容器 ──> 开始消费
yield
生成一个 ──> 消费一个 ──> 生成下一个 ──> 消费下一个这种按需推进的行为称为惰性计算(lazy evaluation)。它适合大文件逐行处理、数据库分批读取、无限序列、日志流和模型增量输出。惰性并不等于低内存:如果调用方执行 list(generate_numbers(...)),所有结果仍会被收集到内存中。
还有一个需要注意的是生成器中的异常会延迟到迭代发生时。生成器对象创建成功,不代表后面的数据库访问、网络访问或数据转换成功。
yield from 如何转发上游数据
转发上游 Generator 时,可以写成:
def generate():
yield from response假设上游生成器逐个产生三个字符串:
from collections.abc import Iterator
def upstream() -> Iterator[str]:
yield "A"
yield "B"
yield "C"只考虑普通遍历时,下面两种写法产生相同结果:
def forward_with_loop() -> Iterator[str]:
for item in upstream():
yield item
def forward_with_yield_from() -> Iterator[str]:
yield from upstream()数据流是逐项转发,不会先收集成列表:
upstream yield "A" ──> forward yield "A" ──> 调用方收到 "A"
upstream yield "B" ──> forward yield "B" ──> 调用方收到 "B"
upstream yield "C" ──> forward yield "C" ──> 调用方收到 "C"可以把 yield from iterable 理解成为 “把这个可迭代对象产生的值继续向外产生”。但它不只等同于 for item in iterable: yield item:外层生成器会把 send()、throw()、close() 等操作委托给子迭代器,并接收子生成器结束时的返回值。因此,官方文档把它称为“向子迭代器委托”。
def child():
yield "working"
return "finished"
def parent():
result = yield from child()
print(result)
print(list(parent()))
# 打印 finished
# 得到 ['working']后文的 Flask 响应只需要逐项转发,不依赖 send() 或子生成器返回值。
Flask 如何把 Generator 作为响应体
普通 Flask 视图是返回已生成的内容:
@app.get("/hello")
def hello():
return "Hello"如果响应体改为 Iterable,Flask、Werkzeug 与 WSGI 服务器组成的响应链就会在发送响应时去迭代它。
WSGI(Web Server Gateway Interface)是 Web 服务器与 Python Web 应用或框架之间的标准接口。服务器调用应用,应用返回一个 Iterable,服务器再迭代它并把产生的内容交给客户端。本文提到“WSGI 响应链”时,指的就是 Flask 返回响应之后、服务器消费响应体的这一层。完整规范见 PEP 3333。
下面的伪代码体现了服务器消费响应体的过程:
for chunk in response_iterable:
write_to_client(chunk)for 内部会调用 next()。每次生成器执行到 yield,一段数据就交还给响应链;下一次迭代再从暂停处继续。因此,业务代码里没有显式写 next(generator),消费方仍然会推进生成器。
import time
from flask import Flask, Response
app = Flask(__name__)
@app.get("/text-stream")
def text_stream():
def generate():
for letter in ("A", "B", "C"):
yield f"{letter}\n"
time.sleep(1)
return Response(generate(), mimetype="text/plain")这里的 mimetype="text/plain" 是有意为之。这个例子只说明 Flask 如何发送普通文本流,还没有进入 SSE。A\n、B\n、C\n 也不是完整的 SSE 消息;如果把媒体类型改为 text/event-stream,响应体还需要遵守下一节介绍的 SSE 分帧格式。
这条路径可以按时间展开:
1. 客户端请求 /text-stream
2. Flask 调用 text_stream()
3. generate() 创建生成器对象,函数体尚未执行
4. text_stream() 返回 Response
5. WSGI 响应链开始迭代生成器
6. 第一次 next() 推进到 yield "A\n"
7. 响应链取得并发送这一段数据
8. 后续迭代依次取得 "B\n"、"C\n"
9. 生成器结束,响应体迭代结束Flask 官方的 Streaming Contents 使用生成器和直接响应实现这种模式。这里的“发送一段”描述的是应用向 WSGI 响应链交出一段内容,并不保证客户端立刻看见它。调试中间件、反向代理、压缩层和客户端都可能缓冲数据。
SSE 是有格式的单向事件流
流式 HTTP 响应表示响应体可以逐段产生,SSE(Server-Sent Events)则在这条响应流上规定了一种文本事件格式。它允许服务器通过一个持续存在的 HTTP 响应向浏览器单向发送事件;浏览器不能通过同一条 SSE 连接向服务器反向发送消息。
SSE 响应使用:
Content-Type: text/event-stream事件流使用 UTF-8 文本。一个事件由若干字段行组成,空行表示当前事件结束:
event: message
data: {"text":"你好"}
event: message
data: {"text":"世界"}
对应的 Python 字符串必须保留结尾的两个换行符:
def event_stream():
yield 'event: message\ndata: {"text":"你好"}\n\n'event 是事件类型,浏览器可以按名称监听;省略 event 时,EventSource 会把它作为默认的 message 事件。data 是事件数据,可以是普通文本,也可以放一段 JSON 字符串。连续的多行 data: 会由客户端用换行符连接。以冒号开头的行是注释,常用于心跳:
: keep-alive
SSE 的字段、分帧和重连行为见 MDN 的 Using server-sent events。mimetype="text/event-stream" 只是在响应头中声明协议;如果生成器产出的字符串没有遵守 SSE 分帧格式,浏览器的 EventSource 仍无法按事件解析。
stream_with_context() 保留请求上下文
stream_with_context() 处理的是 Flask Request Context。Flask 收到请求时建立 Request Context,使 request、session 等代理指向当前请求的数据;同时存在的不同请求都各自有自己的上下文。
普通视图可以在请求处理期间读取 request:
from flask import request
@app.get("/user")
def user():
return request.args.get("name", "anonymous")request 不是所有请求共享的全局变量,而是上下文本地代理(context-local proxy)。请求处理结束后,Request Context 会被退出。
普通视图的生命周期是:
请求到达
↓
Flask 建立 Request Context
↓
执行视图函数,request 可用
↓
视图返回响应
↓
Request Context 被弹出未使用 stream_with_context()
流式响应的 Generator 在视图返回后才被消费。下面的代码虽然把 generate() 定义在视图内部,request 的读取操作却发生在视图返回以后:
from flask import Response, request
@app.get("/broken")
def broken():
def generate():
yield request.args.get("name", "anonymous")
return Response(generate(), mimetype="text/plain")这段代码的时间线是:
请求到达
↓
建立 Request Context
↓
调用 broken()
↓
generate() 只创建生成器,尚未读取 request
↓
broken() 返回 Response
↓
视图返回,原 Request Context 不再保持
↓
WSGI 响应链调用 next(generator)
↓
执行 request.args.get(...)
↓
RuntimeError: Working outside of request context使用 stream_with_context()
Flask 提供 stream_with_context() 包装响应迭代器。在迭代器被消费期间,它会保持当前请求上下文可用:
from flask import Response, request, stream_with_context
@app.get("/working")
def working():
def generate():
yield request.args.get("name", "anonymous")
return Response(
stream_with_context(generate()),
mimetype="text/plain",
)它也能作为装饰器使用:
@app.get("/working")
def working():
@stream_with_context
def generate():
yield request.args.get("name", "anonymous")
return Response(generate(), mimetype="text/plain")使用后的时间线是:
视图创建生成器
↓
stream_with_context 包装生成器
↓
视图返回 Response
↓
响应链迭代包装后的生成器
↓
生成器执行期间 Request Context 可用
↓
request.args 读取成功
↓
生成器耗尽,上下文退出stream_with_context() 不负责让响应“流起来”。没有它,Response(generate()) 也是可迭代响应;生成器不访问 request、session、g、current_app 等上下文代理时,无须这个包装。它解决的是生成器执行时间晚于视图函数所造成的上下文生命周期问题。Flask 对这一行为的说明见 Streaming with Context 和 The Request Context。
这个包装不会把上下文传到新建的工作线程。线程需要使用请求值时,可以在视图处于请求上下文期间把值提取成普通数据,再传给线程。线程、Queue 与流式消费者之间的关系见 Python Thread 线程基础。
一个可运行的 Flask SSE 例子
下面的程序同时提供普通 JSON 响应和 SSE 响应。流式路由在生成器中读取查询参数,因此使用 stream_with_context:
import json
import time
from collections.abc import Iterator
from flask import Flask, Response, jsonify, request, stream_with_context
app = Flask(__name__)
@app.get("/blocking")
def blocking():
"""一次返回 JSON。"""
name = request.args.get("name", "World")
return jsonify({"message": f"Hello, {name}"})
@app.get("/events")
def events():
"""返回 SSE 事件流。"""
@stream_with_context
def generate() -> Iterator[str]:
name = request.args.get("name", "World")
messages = ["开始生成", f"Hello, {name}", "生成完成"]
for index, message in enumerate(messages, start=1):
data = json.dumps(
{"index": index, "message": message},
ensure_ascii=False,
)
yield f"event: message\ndata: {data}\n\n"
time.sleep(1)
yield "event: done\ndata: {}\n\n"
return Response(
generate(),
status=200,
mimetype="text/event-stream",
headers={"Cache-Control": "no-cache"},
)
if __name__ == "__main__":
app.run(debug=True)保存为 app.py 后,安装 Flask 并启动:
python -m pip install Flask
python app.py普通接口会一次得到 JSON:
curl 'http://127.0.0.1:5000/blocking?name=Youyou'{ "message": "Hello, Youyou" }curl -N 关闭 curl 的输出缓冲,便于观察事件逐个到达:
curl -N 'http://127.0.0.1:5000/events?name=Youyou'终端会依次出现:
event: message
data: {"index": 1, "message": "开始生成"}
event: message
data: {"index": 2, "message": "Hello, Youyou"}
event: message
data: {"index": 3, "message": "生成完成"}
event: done
data: {}
第一条消息可以立即产生,后面的消息间隔约一秒。服务端的内部状态依次变化:
客户端建立连接
↓
generate() 读取 name=Youyou
↓
yield 第 1 个事件,生成器暂停
↓
等待 1 秒
↓
yield 第 2 个事件,生成器暂停
↓
等待 1 秒
↓
yield 第 3 个事件,生成器暂停
↓
等待 1 秒
↓
yield done 事件
↓
生成器结束,HTTP 响应结束浏览器可以用 EventSource 接收同源 SSE 地址:
const source = new EventSource('/events?name=Youyou');
source.addEventListener('message', (event) => {
const data = JSON.parse(event.data);
console.log(data.index, data.message);
});
source.addEventListener('done', () => {
source.close();
});
source.onerror = (error) => {
console.error('SSE connection error', error);
};组合应用:统一封装普通响应与 SSE
理解各层机制后,可以重新阅读开头的响应封装。为了不依赖某个项目自定义的 Response 类型,下面用 Mapping 表示普通结果,并用 Flask 自带的 jsonify() 创建 JSON 响应:
from collections.abc import Generator, Mapping
from typing import Any
from flask import Response as FlaskResponse
from flask import jsonify, stream_with_context
def compact_generate_response(
response: Mapping[str, Any] | Generator[str, None, None],
) -> FlaskResponse:
if isinstance(response, Mapping):
return jsonify(dict(response))
def generate() -> Generator[str, None, None]:
yield from response
return FlaskResponse(
stream_with_context(generate()),
status=200,
mimetype="text/event-stream",
)类型分支
def serialize_if_blocking(response):
if isinstance(response, Mapping):
return jsonify(dict(response))这段代码约定业务层会返回两类对象:表示普通结果的 Mapping,或者逐项产生 SSE 消息的 Generator。Mapping 转成 dict 后由 jsonify() 创建一次性 JSON 响应;Generator 进入流式分支。
类型标注描述了这个边界,但运行时仍要由调用方遵守。传入整数或 None 会进入流式分支,直到服务器尝试迭代时才报错。若业务还有其他响应类型,应当在这里增加对应分支或提前校验。
上游事件转发
def generate() -> Generator[str, None, None]:
yield from response内层生成器不改变事件,只转发上游 response 产生的内容。它也为日志、编码、心跳、异常事件转换或资源清理保留了边界。
如果只有原样转发,而且 response 满足 stream_with_context() 接受的迭代器协议,也可以写成:
def build_stream_response(response):
return FlaskResponse(
stream_with_context(response),
status=200,
mimetype="text/event-stream",
)因此,额外定义 generate() 不是保持请求上下文的必要条件。它在这个组合函数中负责生成器委托,也保留了扩展点。
Flask 响应封装
def build_stream_response():
return FlaskResponse(
stream_with_context(generate()),
status=200,
mimetype="text/event-stream",
)这部分分别完成四件事:
| 组件 | 职责 |
|---|---|
generate() | 提供可迭代的响应体 |
stream_with_context(...) | 迭代期间保持当前 Request Context |
FlaskResponse(...) | 创建状态码为 200 的 HTTP 响应 |
text/event-stream | 声明响应体使用 SSE 格式 |
调用链如下:
上游业务 Generator
│ yield event
▼
内层 generate()
│ yield from response
▼
stream_with_context()
│ 维持请求上下文
▼
Flask Response / WSGI 响应链
│ 迭代并交出每个 chunk
▼
SSE 客户端流式响应的边界
响应头必须在正文开始前确定
HTTP 会先发送状态行和响应头,再发送响应体。第一个正文 chunk 发出以后,不能再把状态码从 200 改成 500,也不能追加依赖响应头传递的 Cookie。
因此,能在开始流式输出前完成的参数校验、权限判断和资源检查,应当在返回 Response 前完成。Flask 官方文档说明:如果生成器会读取 session,视图中也应先访问它,以便 Flask 设置正确的 Vary: Cookie;不要在生成器中修改 session,因为携带修改结果的 Set-Cookie 响应头可能已经发出。
流中异常不再是普通 JSON 错误
生成器可能在发送若干事件后失败:
HTTP 200 和响应头已发送
↓
客户端收到 event 1
↓
客户端收到 event 2
↓
生成器抛出异常此时服务端无法改写已发出的 HTTP 响应;它可以终止连接,或者捕获异常并把错误编码成 SSE 事件:
event: error
data: {"message":"generation failed"}
错误事件属于应用协议,并不会把已发送的 HTTP 状态码改成 500。客户端需要同时处理协议内的 error 事件和连接本身的异常。
中间层可能缓冲
应用每次 yield 一段内容,不等于浏览器每次都会立刻显示一段。开发调试器、WSGI 中间件、反向代理、压缩、CDN 和客户端缓冲都可能合并小块内容。排查时可以先用 curl -N 直接访问应用,再逐层加入代理,以判断缓冲发生在哪里。
客户端断开后要释放资源
用户关闭页面、网络中断或代理超时后,响应迭代可能被关闭。持有数据库游标、文件、订阅或后台任务的生成器需要在 finally 中释放资源:
def generate():
subscription = subscribe()
try:
yield from subscription
finally:
subscription.close()若后台生产者和 HTTP 消费者通过 Queue 连接,还要设计停止信号、超时和生产者取消机制。否则客户端离开后,后台任务仍可能运行或向无人消费的 Queue 写入。
同步 Generator 不等于 async generator
本文使用的是同步生成器:
def generate():
yield "data"异步生成器使用 async def 和 yield:
async def generate():
yield "data"异步生成器实现的是异步迭代协议,通过 async for 消费,不能因为名字相同就直接交给只接受同步 Iterable 的 WSGI 响应链。具体接入方式取决于 Web 框架和部署协议,不能用同步示例直接替换。
概念之间的最终关系
这段响应封装不是四个孤立 API 的组合,而是一条由执行时间串起来的数据通路:
Generator 保存执行状态并按需产生数据
↓
yield 把一个值交给消费方并暂停
↓
yield from 把上游迭代器产生的值继续向外委托
↓
WSGI 响应链迭代 Generator,把每段内容交给客户端
↓
SSE 规定每段文本怎样组成事件
↓
stream_with_context 让延迟执行的生成器仍能访问当前请求上下文Python - 高级特性整理了迭代器和生成器在 Python 语法体系中的位置;Python Thread 线程基础连接 Worker Thread、Queue、Generator 与 HTTP/SSE 响应,说明后台生产和前台消费怎样组成数据流。