Faust LiveCheck TestRunner 深度解析:生产环境端到端测试执行器的设计与运行原理 流处理消息队列后端【免费下载链接】faustPython Stream Processing项目地址https://gitcode.com/gh_mirrors/fa/faust点击查看免费下载导读faust.livecheck.runners是 FaustPython Stream Processing中 LiveCheck 端到端测试框架的执行引擎其核心类TestRunner负责执行一次测试并跟踪其状态。本文围绕 faust.livecheck.runners 参考文档该 RST 通过automodule指令自动收集源码 docstring 成文展开结合仓库源码逐层拆解TestRunner的状态机、异常分发、日志缓冲、信号延迟统计与报告生成机制并给出基于 examples/livecheck.py 的实战用法。读完本文你将能理解 LiveCheck 测试从触发、执行、等待信号到产出报告的完整生命周期并能在自己的 Faust 应用中编写和调试 LiveCheck 用例。一、TestRunner 是什么LiveCheck 的执行单元在 LiveCheck 架构中三层对象各司其职对象职责源码位置LiveCheckApp管理用例注册、Kafka 主题pending_tests/bus/reports、并发调度faust/livecheck/app.pyCase定义测试用例声明信号、配置概率/频率、实现run()与make_fake_request()faust/livecheck/case.pyTestRunner执行并跟踪单次测试执行TestExecutionfaust/livecheck/runners.pyTestRunner是Case与一次具体执行之间的临时实体每次有测试执行事件到达pending_tests主题LiveCheck._execute_testsagent 会把TestExecution交给对应Case.execute()而Case.execute()会为这次执行新建一个TestRunner实例faust/livecheck/case.pyasync def execute(self, test: TestExecution) - None: Execute test using :class:TestRunner. t_start monotonic() runner self.Runner(self, test, startedt_start) with current_execution_stack.push(runner): # resolve_models await runner.execute()关键点Case是长期存活的服务被注册为运行时依赖而TestRunner是每次执行临时创建、用完即弃的对象其生命周期恰好覆盖一次TestExecution从开始到产出TestReport的完整过程。二、TestRunner 的核心状态与数据结构2.1 实例属性一览TestRunner在初始化时faust/livecheck/runners.py持有case所属Case实例test本次执行的TestExecution含 id、case_name、timestamp、test_args/test_kwargs、expiresstarted启动时刻monotonic()单调时钟由Case.execute()传入ended/runtime结束时刻与运行时长秒logs日志缓冲列表List[Tuple[str, Tuple]]消息模板 参数signal_latencyDict[str, float]记录每个信号从发出到被解析的延迟report/error最终TestReport与异常对象。同时TestRunner通过CompositeLogger把日志挂接到case.log.logger上并自定义格式化器_format_logfaust/livecheck/runners.py让每条日志都带上本次执行的短标识def _format_log(self, severity: int, msg: str, *args, **kwargs) - str: return f[{self.test.shortident}] {msg}shortident形如case名:id缩写id 最多显示 15 个字符超出加[...]后缀见 faust/livecheck/models.py。2.2 状态机State 枚举执行结果状态定义在 faust/livecheck/models.pyclass State(Enum): INIT INIT PASS PASS FAIL FAIL ERROR ERROR TIMEOUT TIMEOUT STALL STALL SKIP SKIP def is_ok(self) - bool: return self in OK_STATES OK_STATES frozenset({State.INIT, State.PASS, State.SKIP})其中STALL不是单次执行产生的状态而是Case级套件停滞长时间无测试到达时由_check_frequency定时任务上报的状态faust/livecheck/case.py。TestRunner本身负责流转到PASS/FAIL/ERROR/TIMEOUT/SKIP。三、execute()一次测试执行的主流程execute()是TestRunner的核心入口faust/livecheck/runners.py流程如下推入当前测试上下文with current_test_stack.push(test)把本次执行压入上下文栈使应用内任意位置的livecheck.current_test都能取到当前测试faust/livecheck/locals.py前置过滤若case.active为假 → 跳过skip(case inactive)若test.is_expired当前时间 ≥expires→ 跳过skip(expired)参数准备调用_prepare_args/_prepare_kwargs处理测试参数执行await self.case.run(*args, **kwargs)然后按异常类型分发见下节。3.1 参数准备maybe_model 归一化_prepare_args与_prepare_kwargs统一走_prepare_valfaust/livecheck/runners.pydef _prepare_val(self, arg: Any) - Any: return maybe_model(arg)maybe_model会把普通 dict 转换为对应的 FaustRecord模型若已注册保证测试参数在跨 Kafka 序列化/反序列化后仍是类型完整的模型对象。3.2 异常分发矩阵核心execute()的try/except块实现了异常 → 状态 → 回调 → 报告的完整映射这是理解 LiveCheck 行为的关键case.run() 抛出触发回调结果状态重新抛出asyncio.CancelledError无静默通过—不重抛TestSkippedon_skippedSKIPTestSkippedTestTimeouton_timeoutTIMEOUTTestTimeoutAssertionError断言失败on_failedFAILTestFailed包装LiveCheckError其他on_errorERROR原异常其他任意Exceptionon_errorERRORTestRaised包装无异常on_passPASS—代码对应 faust/livecheck/runners.py而异常类体系定义在 faust/livecheck/exceptions.py基类LiveCheckError派生出SuiteFailed、TestSkipped、TestFailed、TestRaised、TestTimeout等。值得注意的工程细节AssertionError被包装为TestFailed重新抛出便于上层 agentLiveCheck._execute_tests统一按LiveCheckError捕获并忽略faust/livecheck/app.py普通Exception被包装为TestRaised——即测试代码本身报错了与断言未满足被区分对待单元测试 t/unit/livecheck/test_runners.py 用参数化用例逐条验证了上述映射含回调只调用一次、on_pass不被调用等断言。3.3 skip()带回溯的跳过skip()faust/livecheck/runners.py主动构造TestSkipped异常、立即raise并捕获目的是让异常携带完整回溯后再调用on_skippedasync def skip(self, reason: str) - NoReturn: exc TestSkipped(fTest {self.test.ident} skipped: {reason}) try: raise exc except TestSkipped as exc: # save with traceback await self.on_skipped(exc) raise四、状态回调on_* 系列方法与报告定稿4.1 结果回调的共性流程on_failed/on_error/on_timeout/on_passfaust/livecheck/runners.py都遵循同一模式self.end()记录ended monotonic()并计算runtimefaust/livecheck/runners.py设置self.state与self.error输出对应日志log.exception或log_info回调Case上的on_test_failed/on_test_error/on_test_timeout/on_test_pass触发Case级统计失败计数、连续失败、历史采样等见 faust/livecheck/case.py调用_finalize_report()产出报告。on_pass还会把运行时长humanize_seconds(..., microsecondsTrue)格式化写入日志并在输出前刷新缓冲日志_flush_logs示例日志形如[test_order:ab3f1c2d…] Test OK in ~0.05 seconds √4.2 日志缓冲机制realtime_logslog_info()faust/livecheck/runners.py实现了可选的缓冲日志当case.realtime_logs为False时日志先追加进self.logs列表只有测试通过时才由_flush_logsfaust/livecheck/runners.py一次性批量输出——避免频繁的日志 IO 干扰高并发测试若开启realtime_logsTrue则每条日志实时打印。注意on_start与on_signal_wait也使用log_info因此在非实时模式下这些中间日志会在测试结束时统一回放。4.3 信号等待与延迟统计on_signal_wait与on_signal_received是TestRunner与信号系统faust/livecheck/signals.py的接口点当测试中的Signal.wait(timeout...)开始等待时Signal.wait会先调用runner.on_signal_wait(signal, timeout)记录等待哪个信号、第几个/共几个faust/livecheck/signals.pyTestRunner打印∆ index/total NAME (timeouts)...形式的日志信号被解析resolve后runner.on_signal_received(signal, time_start, time_end)计算延迟并写入self.signal_latency[signal.name]faust/livecheck/runners.py。async def on_signal_received(self, signal, time_start, time_end) - None: latency time_end - time_start self.signal_latency[signal.name] latency这些延迟数据最终会进入TestReport.signal_latency用于观测分布式链路中每一步的响应时间。4.4 报告定稿TestReport_finalize_report()faust/livecheck/runners.py在有错误时抓取traceback.format_tb生成回溯文本然后构造TestReport并交给case.post_report(report)self.report TestReport( case_nameself.case.name, stateself.state, testself.test, runtimeself.runtime, signal_latencyself.signal_latency, errorstr(error) if error else None, tracebacktb, ) await self.case.post_report(self.report)TestReport的字段定义在 faust/livecheck/models.pycase_name、state、test可为空、runtime、signal_latency、error、traceback。报告最终经Case.post_report → LiveCheck.post_report发布到reportslivecheck-report主题faust/livecheck/case.py、faust/livecheck/app.py供外部监控/告警系统消费。单元测试 t/unit/livecheck/test_runners.py 分别验证了有错误error/traceback 填充与无错误二者为 None两种报告路径。五、与 Case / LiveCheck App 的联动5.1 测试执行的完整调用链一次 LiveCheck 测试从触发到报告完整链路为Case.trigger() / maybe_trigger() # 创建 TestExecution发到 pending_tests 主题 ↓ LiveCheck._execute_tests (agent) # 消费 pending_tests查找注册的 Case ↓ Case.execute(test) # 新建 TestRunner推入 current_execution_stack ↓ TestRunner.execute() # 前置过滤 → 参数准备 → case.run() ↓ case.run() 中 Signal.wait() / Signal.send() # 通过 livecheck-bus 主题跨进程同步 ↓ TestRunner.on_* 回调 → _finalize_report() # 状态定稿 ↓ Case.post_report → LiveCheck.post_report # 发布到 livecheck-report 主题5.2 上下文栈如何追踪测试TestRunner的执行被两层本地栈包裹faust/livecheck/locals.pycurrent_test_stack保存当前TestExecution让流处理器agent通过livecheck.current_test获取测试上下文current_execution_stack保存当前TestRunner让Signal.wait能通过case.current_execution取到 runner 并调用其on_signal_wait/on_signal_receivedfaust/livecheck/signals.py。配合LiveCheckSensorfaust/livecheck/app.py在流事件进入/离开时自动 push/pop 测试上下文LiveCheck 就能在 Kafka 消息穿越多个 agent 时持续追踪这条消息属于哪次测试——这正是 LiveCheck 端到端追踪HTTP/Kafka 头LiveCheck-Test-Id等的实现基础见 faust/livecheck/models.py 与 faust/livecheck/app.py。5.3 并发与配置LiveCheck默认用test_concurrency100个并发 agent 执行测试、bus_concurrency30个并发 agent 处理信号faust/livecheck/app.pyTestRunner因而必须轻量且无共享状态——每次执行独立实例化靠本地栈隔离上下文这是它在高并发下安全运行的保证。六、实战在 stock 订单示例中观察 TestRunner仓库自带的 examples/livecheck.py 完整演示了 TestRunner 的运转。启动方式见 docs/userguide/livecheck.rst$ python examples/livecheck.py worker -l info # 终端 1订单系统 worker $ python examples/livecheck.py livecheck -l info # 终端 2LiveCheck 实例随后访问http://localhost:6066/order/init/sell/或运行python examples/livecheck.py post_order --sidesell下单。由于用例默认probability0.5约一半请求会触发测试执行。用例定义examples/livecheck.pylivecheck.case(warn_stalled_after5.0, frequency0.5, probability0.5) class test_order(Case): order_sent_to_db: Signal[Order] order_sent_to_kafka: Signal[None] order_cache_in_redis: Signal[None] order_executed: Signal[str] async def run(self, side: str) - None: # 1) wait for order to be sent to database. order await self.order_sent_to_db.wait(timeout30.0) # contract: order id matches test execution id assert order.id self.current_execution.test.id assert order.side side # 2) wait for order to be sent to Kafka await self.order_sent_to_kafka.wait(timeout30.0) # 3) wait for redis index to be updated. await self.order_cache_in_redis.wait(timeout30.0) assert await livecheck.cache.client.sismember( forder.{order.user_id}.orders, order.id) # 4) wait for execution agent to execute the order. await self.order_executed.wait(timeout30.0) async def make_fake_request(self) - None: await self.get_url(http://localhost:6066/order/init/sell/?fake1)可以看到run()中的每个Signal.wait(timeout30.0)都会触发TestRunner.on_signal_wait日志若任一断言失败如某系统篡改了订单 sideTestRunner会走on_failed→State.FAIL→ 生成带回溯的TestReport若 30 秒内信号未到达Signal.wait内部抛TestTimeoutfaust/livecheck/signals.pyrunner 转入TIMEOUT状态并同样产出报告。而make_fake_request配合frequency0.5与warn_stalled_after5.0会让 LiveCheck 在无真实流量时主动注入 fake 请求避免套件因无测试而进入STALL告警。七、小结faust.livecheck.runners.TestRunner是 LiveCheck 框架中承上启下的执行引擎向上承接Case的调度与信号同步向下沉淀每次执行的可观测数据状态、时长、信号延迟、错误回溯最终以TestReport形式发布到 Kafka 供监控消费。理解其状态机与异常分发矩阵是排查 LiveCheck 测试为什么被跳过、为什么超时、为什么失败的第一把钥匙。相关源码可继续深入阅读执行器实现faust/livecheck/runners.py用例基类与套件状态机faust/livecheck/case.py模型与状态定义faust/livecheck/models.py异常体系faust/livecheck/exceptions.py信号机制faust/livecheck/signals.py上下文栈faust/livecheck/locals.py应用装配与 agent 调度faust/livecheck/app.py单元测试异常分发/报告路径验证t/unit/livecheck/test_runners.py完整实战示例examples/livecheck.py使用指南docs/userguide/livecheck.rst赞分享流处理消息队列后端【免费下载链接】faustPython Stream Processing项目地址https://gitcode.com/gh_mirrors/fa/faust点击查看免费下载相关推荐SQL Server R Services SSMS 自定义报表基于 RDL 的配置巡检、资源监控与 R 脚本执行统计指南SQL Server R Services SSMS 自定义报表基于 RDL 的配置巡检、资源监控与 R 脚本执行统计指南 本篇技术指南聚焦于 sql ser流处理消息队列后端Faust LiveCheck 实战指南面向生产环境的端到端测试与异常监控Faust LiveCheck 实战指南面向生产环境的端到端测试与异常监控 导读 LiveCheck 是 FaustPython Stream Proces流处理消息队列后端Optimism op-e2e 端到端测试套件分类架构、运行方式与设计原则深度解析Optimism op e2e 端到端测试套件分类架构、运行方式与设计原则深度解析 op e2e 是 Optimism 单仓库中的 Go 集成测试模块负责对区块链Web3后端上一篇企业级低代码平台终极指南如何用JeecgBoot快速构建AI应用下一篇Roo Code完整指南5分钟打造你的AI编程团队创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考