- 流处理
- 消息队列
- 后端
【免费下载链接】faust
Python Stream Processing
本指南聚焦 Faust(Python Stream Processing 框架)内置的监控统计模块faust.web.apps.stats。该模块以 Blueprint(蓝图)形式对外提供两个 HTTP 端点:/返回所有注册 Sensor(传感器/监控器)的实时统计快照(JSON),/assignment/返回当前节点的分区分配信息(活跃分区与备用分区)。读完本文,你将掌握如何开启并访问这两个端点、如何解读 JSON 输出中的监控指标与分配字段、以及如何将这些端点与 Faust 的传感器体系和分区分配器(Partition Assignor)在源码层面对应起来,为生产环境的可观测性排查提供直接抓手。
模块定位:Faust 的“调试用”内置统计端点
在 Faust 的 Web 体系中,faust.web.apps.stats是一个默认只在调试模式(debug)下挂载的内置 Web App。它并不参与生产流量,而是为开发与排障阶段提供进程内的运行状态快照。
从 faust/web/base.py 可以看到 Faust 把内置蓝图分为三组:
DEFAULT_BLUEPRINTS: _BPList = [ ('/router', 'faust.web.apps.router:blueprint'), ('/table', 'faust.web.apps.tables.blueprint'), ] PRODUCTION_BLUEPRINTS: _BPList = [ ('', 'faust.web.apps.production_index:blueprint'), ] DEBUG_BLUEPRINTS: _BPList = [ ('/graph', 'faust.web.apps.graph:blueprint'), ('', 'faust.web.apps.stats:blueprint'), ]其中DEFAULT_BLUEPRINTS(/router、/table)始终挂载;当应用处于 debug 模式时追加DEBUG_BLUEPRINTS(/graph、/stats),否则追加PRODUCTION_BLUEPRINTS。挂载逻辑在Web.__init__中:
blueprints = list(self.default_blueprints) if self.app.conf.debug: blueprints.extend(self.debug_blueprints) else: blueprints.extend(self.production_blueprints)也就是说:faust.web.apps.stats只在你以--debug方式启动faust worker时才会被注册。这一事实来自源码 faust/web/base.py,生产环境默认不会暴露这两个统计端点。
模块骨架:一个 Blueprint + 两个 View
模块整体结构非常精简,完整源码见 faust/web/apps/stats.py:
"""HTTP endpoint showing statistics from the Faust monitor.""" from collections import defaultdict from typing import List, MutableMapping, Set from faust import web from faust.types.tuples import TP __all__ = ['Assignment', 'Stats', 'blueprint'] TPMap = MutableMapping[str, List[int]] blueprint = web.Blueprint('monitor') @blueprint.route('/', name='index') class Stats(web.View): """Monitor statistics.""" async def get(self, request: web.Request) -> web.Response: """Return JSON response with sensor information.""" return self.json( {f'Sensor{i}': s.asdict() for i, s in enumerate(self.app.sensors)}) @blueprint.route('/assignment/', name='assignment') class Assignment(web.View): """Cluster assignment information.""" @classmethod def _topic_grouped(cls, assignment: Set[TP]) -> TPMap: tps: MutableMapping[str, List[int]] = defaultdict(list) for tp in sorted(assignment): tps[tp.topic].append(tp.partition) return dict(tps) async def get(self, request: web.Request) -> web.Response: """Return current assignment as a JSON response.""" assignor = self.app.assignor return self.json({ 'actives': self._topic_grouped(assignor.assigned_actives()), 'standbys': self._topic_grouped(assignor.assigned_standbys()), })它导出了三个符号:blueprint(Blueprint 实例)、Stats(根路径/的 View)、Assignment(/assignment/路径的 View)。两个 View 都继承web.View并实现async def get,因此它们处理的是 HTTP GET 请求。
Blueprint 机制回顾
blueprint = web.Blueprint('monitor')中的'monitor'是该蓝图的名字,同时会被用作视图名的命名空间前缀。参照 faust/web/blueprints.py 的实现:@blueprint.route()会生成一个FutureRoute暂存起来,直到蓝图被注册到 App 时才真正创建视图;视图名由_view_name拼接为f'{name}:{handler_name}',因此这里生成的路由名称分别是monitor:index和monitor:assignment。
注册时通过BlueprintManager._apply_blueprint调用bp.register(web.app, url_prefix=prefix)并随后调用bp.init_webserver(web)(见 faust/web/base.py)。由于DEBUG_BLUEPRINTS中该蓝图的前缀为空字符串'',最终这两个端点的 URL 就是http://<host>:<port>/与http://<host>:<port>/assignment/。
端点一:/—— 传感器统计快照
Stats视图遍历self.app.sensors中注册的所有传感器(SensorDelegate持有的是Set[SensorT],见 faust/sensors/base.py),并为每个传感器生成Sensor{i}键(i为从 0 开始的序号),其值来自每个传感器的asdict():
return self.json( {f'Sensor{i}': s.asdict() for i, s in enumerate(self.app.sensors)})谁来提供统计数据?
统计数据的实际来源是各个Sensor子类的asdict()方法。Faust 自带的Monitor监控器(faust/sensors/monitor.py)实现了最完整的状态导出,其asdict()返回的字段包括:
messages_active、messages_received_total、messages_sent、messages_sent_by_topicmessages_s、messages_received_by_topicevents_active、events_total、events_s、events_runtime_avgevents_by_task、events_by_streamcommit_latency、send_latency、send_errorsassignment_latency、assignments_completed、assignments_failedtopic_buffer_full、tables等
(完整实现见 faust/sensors/monitor.py 及后续行。)
因此,在 debug 模式下访问/,你会得到类似如下的 JSON 结构:
{ "Sensor0": { "messages_active": 0, "messages_received_total": 12345, "messages_sent": 9876, "messages_sent_by_topic": {"topic-a": 5000, "topic-b": 4876}, "messages_s": 12.34, "messages_received_by_topic": {"topic-a": 6000, "topic-b": 6345}, "events_active": 0, "events_total": 54321, "events_s": 56.78, "events_runtime_avg": 0.0001, "events_by_task": {...}, "events_by_stream": {...}, "commit_latency": [...], "send_latency": [...], "send_errors": 0, "assignment_latency": [...], "assignments_completed": 3, "assignments_failed": 0, "topic_buffer_full": {...}, "tables": {...} } }底层事件采集机制
这些数字并不是凭空产生的,而是 Sensor 接口定义的一系列回调在运行时被驱动调用后累积出来的。Sensor 基类在 faust/sensors/base.py 中定义了完整的钩子集合,包括:
- 消息生命周期:
on_message_in(消费者收到消息)、on_stream_event_in/on_stream_event_out(事件进入/离开流)、on_message_out(所有流处理完,可提交偏移) - 表操作:
on_table_get、on_table_set、on_table_del - 提交/发送:
on_commit_initiated、on_commit_completed、on_send_initiated、on_send_completed、on_send_error - 分区分配与再平衡:
on_assignment_start、on_assignment_completed、on_assignment_error、on_rebalance_start、on_rebalance_return、on_rebalance_end - Web 请求:
on_web_request_start、on_web_request_end
SensorDelegate(同样位于 faust/sensors/base.py)会把每一次回调扇出给所有已注册的传感器,并把每个传感器各自返回的中间状态按传感器维度的字典回传。这就是/端点能同时展示多个Sensor{i}快照的原因——你注册了几个传感器,响应里就有几个条目,每个条目的键名取决于该传感器asdict()的字段。
端点二:/assignment/—— 分区分配快照
Assignment视图用于展示当前节点在集群中的分区分配结果,它不依赖传感器,而是直接读取分区分配器(Partition Assignor):
assignor = self.app.assignor return self.json({ 'actives': self._topic_grouped(assignor.assigned_actives()), 'standbys': self._topic_grouped(assignor.assigned_standbys()), })响应结构
响应包含两个键:
actives:当前节点活跃持有的分区(负责实际消费和处理)standbys:当前节点作为备用(standby)持有的分区(用于故障切换/恢复,配合 standby 表)
分区按 topic 分组,分组逻辑由_topic_grouped完成:
@classmethod def _topic_grouped(cls, assignment: Set[TP]) -> TPMap: tps: MutableMapping[str, List[int]] = defaultdict(list) for tp in sorted(assignment): tps[tp.topic].append(tp.partition) return dict(tps)它把TP((topic, partition)元组)集合按 topic 排序后归并成{"topic": [partition0, partition1, ...]}的映射。示例响应:
{ "actives": { "events": [0, 1, 2], "word-counts": [0] }, "standbys": { "word-counts": [1, 2] } }数据来源:assigned_actives / assigned_standbys
assigned_actives()与assigned_standbys()在 faust/assignor/partition_assignor.py 中实现:
def assigned_standbys(self) -> Set[TP]: return { TP(topic, partition) for topic, partitions in self._assignment.standbys.items() for partition in partitions } def assigned_actives(self) -> Set[TP]: return { TP(topic, partition) for topic, partitions in self._assignment.actives.items() for partition in partitions }可以看到,两者分别把分配器内部维护的self._assignment.actives/self._assignment.standbys(TopicToPartitions映射)展平为Set[TP]集合,_topic_grouped再把它们重新按 topic 分组。这一“展平再分组”的往返看似冗余,实则是为了在 JSON 中以稳定的形式呈现:/assignment/端点始终输出{topic: [partitions]},方便脚本直接解析。
app.assignor由 App 在启动时实例化(通过app.conf.assignor配置项指定具体实现类,默认是基于集中式协调的LeaderAssignor,相关参考见 faust/types/settings/settings.py 附近对 assignor 配置的说明)。
如何开启并访问这两个端点
由于 stats 蓝图属于DEBUG_BLUEPRINTS,开启方式就是让 Faust App 运行在 debug 模式:
使用 CLI 启动 worker 并开启 debug 与 web:
faust -A my_app worker --debug --with-web--debug使app.conf.debug = True,从而挂载DEBUG_BLUEPRINTS(包含 stats);--with-web启用内置 Web 服务器(对应web_enabled设置,见 faust/types/settings/settings.py)。
Web 服务器默认绑定
--web-host/--web-port(默认端口 6066,范围 1024–65535,见 faust/types/settings/settings.py)。访问端点:
curl http://localhost:6066/ curl http://localhost:6066/assignment/第一个返回所有传感器的统计快照,第二个返回分区分配快照。
手动注册(1.7 版本提供的替代方式)
如果你不想依赖 debug 模式,也可以像 Faust 1.7 changelog 所记录的那样,在代码里手动把 stats 蓝图挂到任意 URL 前缀下:
app.web.blueprints.add('/stats/', 'faust.web.apps.stats:blueprint')(该用法记录在 docs/history/changelog-1.7.rst 中。)执行后端点变为:
curl http://localhost:6066/stats/ curl http://localhost:6066/stats/assignment/blueprints.add(prefix, blueprint)接受字符串形式的符号路径(faust.web.apps.stats:blueprint),BlueprintManager会在apply()阶段通过symbol_by_name解析并注册(见 faust/web/base.py)。注意add()必须在 Web 服务器启动之前调用,否则会抛出RuntimeError: Cannot add blueprints after server started。
与相关内置模块的关系
faust.web.apps.stats只是 Faust 内置 Web App 家族的一员。同目录下的其他内置应用(见 faust/web/apps/init.py)包括:
faust.web.apps.router:路由/服务发现信息;faust.web.apps.tables:表(Table)状态浏览;faust.web.apps.graph:Agent/流拓扑图;faust.web.apps.production_index:生产模式首页。
其中graph与stats一样同属DEBUG_BLUEPRINTS(见 faust/web/base.py)。如果你想把这些端点暴露到生产环境,需要自行注册对应蓝图并评估安全策略,因为统计信息可能包含内部拓扑与性能细节。
使用建议与注意事项
- 这是调试工具,不是生产监控:stats 端点默认仅在 debug 模式挂载,若需在生产环境使用,请通过
app.web.blueprints.add(...)显式注册,并确认 Web 端口不会对外部网络开放。 - 响应结构随传感器实现变化:
/端点的具体字段完全取决于各传感器asdict()的实现。若要稳定解析,建议同时控制Monitor类或自定义Sensor的asdict()输出;传感器可通过App(..., Monitor=...)配置替换(相关配置见 faust/types/settings/settings.py 中Monitor的设置说明)。 /assignment/只反映当前节点:它输出的是本进程assignor视角下的分配,要获得全集群视角需在每个节点上分别查询,或结合各节点的结果汇总。- 分区号可能无序但值完整:
_topic_grouped按sorted(assignment)输出,topic 有序、每个 topic 内的 partition 号按升序排列,适合直接做断言与可视化。
通过这两个端点,你可以在调试阶段快速回答两类关键问题:“这个节点现在处理了多少消息、事件,性能如何?”(/)以及“这个节点被分配了哪些分区,活跃/备用各是哪些?”(/assignment/),再配合 faust/sensors/base.py 的传感器回调和 faust/assignor/partition_assignor.py 的分配器实现,即可把观测结果一路追到源码层。
- 流处理
- 消息队列
- 后端
【免费下载链接】faust
Python Stream Processing
相关推荐
Faust 内置 Router Web 应用:基于表分区路由的 HTTP 端点源码详解
Faust 内置 Router Web 应用:基于表分区路由的 HTTP 端点源码详解 本文基于 faust.web.apps.router 模块 https:
流处理消息队列后端Docker rm 别名解析:tldr 别名页机制与 docker container rm 实战指南
Docker rm 别名解析:tldr 别名页机制与 docker container rm 实战指南 docker rm 是 Docker CLI 中 doc
流处理消息队列后端ArchiveBox 实时进度监控 API 深入解析:progressmonitor 模块与 progress.json 端点的完整实现
ArchiveBox 实时进度监控 API 深入解析:progressmonitor 模块与 progress.json 端点的完整实现 本篇指南围绕 Arch
后端数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考