- 流处理
- 消息队列
- 后端
【免费下载链接】faust
Python Stream Processing
导读
Faust 的 LiveCheck 模块提供了一种独特的测试理念:让测试用例直接运行在生产或预发布环境中,以被动观察者(passive observer)的身份追踪请求在分布式系统中的流转,并在每个环节校验契约与不变量。本篇文章以 faust/livecheck/case.py 中的Case基类为核心,完整讲解测试用例的配置属性、状态机、触发机制、信号同步、停滞检测与故障恢复,并结合 runners.py、models.py、signals.py 等源码,让你能够从零开始编写一个可投入生产环境运行的 LiveCheck 测试用例。
一、LiveCheck 与 Case 的定位
LiveCheck 是 Faust 内置的"面向生产环境的端到端测试"子系统。与单元测试、集成测试不同,LiveCheck 的核心思想是:
- 追踪请求:测试请求在 Kafka 消息中携带专属头(header),随事件在整个微服务链路中流转;
- 定义契约:开发者在每个处理环节用断言(assert)声明"这里必须满足的约束",比如"订单金额不能超过阈值"、"积分必须增加";
- 概率采样执行:测试并不针对每个请求执行,而是按概率抽样(默认 50%),因此可以做到 100%、30%、50% 乃至 0.1% 的采样率;
- 被动监控:测试用例是被动观察者,能够探测子系统宕机、服务停滞等异常,并发出告警。
从源码结构看,Case是这一切的"单元"——每一个测试用例就是一个Case子类,且每个Case本质上是一个 mode 的Service(见 case.py 中class Case(Service)),这意味着它可以被挂载到 LiveCheck 应用中随应用启停,并拥有定时任务(Service.timer、Service.task)能力。正如 docs/userguide/livecheck.rst 所述:"每个 LiveCheck 测试用例都是一个流处理器,因此可以使用表(table)来存储数据"。
二、Case 类的核心属性与构造参数
Case的构造全部通过关键字参数完成。下面依据 case.py 的实现,列出完整属性与默认值。
2.1 基本属性
| 属性 | 默认值 | 含义 |
|---|---|---|
name | 由子类名生成 | 测试用例名称;若未显式指定,会从子类限定名(qualname)生成 |
active | True | 是否启用该测试;设为False时执行会被跳过(跳过原因记为case inactive) |
status | State.INIT | 测试当前状态,见下文状态机 |
Runner | TestRunner | 执行测试的运行器类,可替换以自定义执行逻辑 |
label | 只读属性 | 形如ClassName: name的人类可读标签,用于日志 |
2.2 执行策略属性
| 属性 | 默认值 | 含义 |
|---|---|---|
probability | 0.5 | 真实流量经过时执行测试的概率。取值1.0表示对每个请求都执行,0.001表示千分之一采样 |
frequency | None | 使用假数据(make_fake_request())执行测试的间隔秒数。若生产流量足够频繁,可设为None,完全依赖真实流量触发 |
warn_stalled_after | 1800.0(30 分钟) | 若在指定秒数内没有任何测试被接收/执行,触发停滞(STALL)警告 |
max_consecutive_failures | 30 | 连续失败达到该次数时,判定整个测试套件失败(抛出SuiteFailed) |
test_expires | timedelta(hours=3) | 测试执行的过期时间;从发起时刻算起超过该时长的执行会被跳过(expired) |
realtime_logs | False | 为True时执行日志实时输出;为False时先缓冲,待测试结束后统一输出 |
2.3 统计与历史属性
| 属性 | 默认值 | 含义 |
|---|---|---|
max_history | 100 | 历史记录最大条数,用于latency_history、frequency_history、runtime_history三个双端队列 |
latency_history/frequency_history/runtime_history | 空deque() | 分别记录延迟偏差、实际触发间隔、单次运行耗时 |
runtime_avg/latency_avg/frequency_avg | None | 上述历史的中位数统计值,由采样任务每 10 秒更新 |
从 case.py 可以看到,Case注册了一个@Service.timer(10.0)的_sampler任务,每 10 秒调用_sample(),用statistics.median计算三个历史队列的中位数,并输出形如Stats: (median) frequency=... latency=... runtime=...的日志。使用中位数而非平均值,可以避免个别长尾耗时污染统计结果。
2.4 HTTP 探测属性
当测试用例需要主动向外部 HTTP 服务发起请求(例如make_fake_request()生成假请求)时,以下参数控制超时与重试策略:
| 属性 | 默认值 | 含义 |
|---|---|---|
url_timeout_total | 300.0(5 分钟) | 请求总超时(秒) |
url_timeout_connect | None | 连接建立超时(秒),None表示不单独限制 |
url_error_retries | 10 | 遇到客户端错误(ClientError)时的最大重试次数 |
url_error_delay_min | 0.5 | 首次重试前的等待秒数 |
url_error_delay_backoff | 1.5 | 每次重试的退避倍率(指数退避) |
url_error_delay_max | 5.0 | 重试等待的最大秒数上限 |
state_transition_delay | 60.0 | 失败状态持续超过该秒数后才允许状态迁移(如从 FAIL 恢复为 PASS) |
2.5 其他内部状态
| 属性 | 默认值 | 含义 |
|---|---|---|
last_test_received | None | 最近一次收到测试的时间戳(单调时钟),停滞检测依赖它 |
last_fail | None | 套件最近一次失败的时间戳 |
consecutive_failures | 0 | 当前连续失败次数 |
total_failures | 0 | 累计失败总次数 |
total_by_state | Counter() | 按状态统计的计数器 |
signals/total_signals | {}/0 | 该用例注册的信号字典及其数量 |
构造时,signals中的每个信号都会通过sig.clone(case=self)绑定到当前用例,并写入self.__dict__,使信号属性可以被直接访问(见 case.py)。
三、状态机:测试用例的生命周期
测试用例的status由 faust/livecheck/models.py 中的State枚举定义:
| 状态 | 含义 |
|---|---|
INIT | 初始状态,属于 OK 状态 |
PASS | 测试通过,属于 OK 状态 |
SKIP | 测试被跳过(如用例被禁用、执行已过期),属于 OK 状态 |
FAIL | 断言失败(不变量被破坏) |
ERROR | 执行过程抛出异常 |
TIMEOUT | 等待信号或执行超时 |
STALL | 长时间没有测试活动,套件停滞 |
其中OK_STATES = frozenset({State.INIT, State.PASS, State.SKIP}),State.is_ok()用于判断当前状态是否正常。
状态迁移发生在以下关键节点(详见 case.py):
- 失败:
on_test_failed/on_test_error/on_test_timeout统一调用_set_test_error_state(state),累加consecutive_failures与total_failures,更新total_by_state;当连续失败达到max_consecutive_failures时,抛出SuiteFailed并调用on_suite_fail; - 通过:
on_test_pass记录运行耗时到runtime_history,如果存在历史失败记录且当前时间晚于失败时间,则调用_maybe_recover_from_failed_state()尝试恢复; - 恢复:
_maybe_recover_from_failed_state在状态不是PASS、且失败持续时间超过state_transition_delay(默认 60 秒)时,将状态置回PASS。这一设计避免了瞬时抖动导致的频繁告警; - 套件失败:
on_suite_fail(exc, new_state)中,仅当当前状态处于 OK 或失败已持续超过state_transition_delay时才真正发布TestReport(防止告警刷屏),否则仅更新状态。
四、编写第一个测试用例:run() 与信号
4.1 自定义 Case 的骨架
Case是一个抽象基类,run方法必须由子类实现,否则抛出NotImplementedError(见 case.py)。一个典型的自定义用例形如:
from faust.livecheck import Case, Signal class OrderFlowCase(Case): # 声明一个信号:等待外部系统(如支付服务)确认 order_confirmed: Signal async def run(self, order_id=None, **kwargs): """核心不变量校验逻辑(可附带任意关键字参数)。""" assert order_id is not None, 'missing order id' # 等待信号在 30 秒内被触发 await self.order_confirmed.wait(timeout=30.0) # ... 其余断言说明(依据 app.py 的case()装饰器实现):
case()装饰器会扫描子类的类型注解,凡是BaseSignal子类的属性都会被自动提取并实例化为信号对象,同时赋予从 1 开始递增的index;- 装饰器还负责把用例注册进
LiveCheck.cases字典(键为用例名),并通过venusian.attach挂到livecheck.case扫描类别下,实现自动发现; - 若未指定
name,默认使用类的限定名(qualname),并处理__main__.前缀。
4.2 信号:测试与系统的同步原语
信号(Signal)是 LiveCheck 测试同步的关键机制,定义在 faust/livecheck/signals.py:
signal.send(value, key=None):通知测试某个外部事件已经完成。若当前处于测试上下文中,key默认为当前测试的id;若不在测试上下文中则必须显式传key且force=True;signal.wait(key=None, timeout=None):等待信号完成,返回事件值。等待过程会调用runner.on_signal_wait记录等待,并在收到后通过on_signal_received计算信号延迟(写入runner.signal_latency[signal.name]);- 底层通过 LiveCheck 应用内部的
_resolved_signals字典与_can_resolve事件实现等待唤醒,超时则抛出TestTimeout; - 收到的事件会经过
_verify_event校验key、signal_name、case_name三者一致,确保信号确实属于当前测试。
信号的实际传输依赖 Kafka:Signal.send将SignalEvent发布到livecheck-bus主题,由 LiveCheck 应用的_populate_signalsagent 消费并回调case.resolve_signal(见 app.py)。
4.3 请求追踪的实现细节
LiveCheck 之所以能"追踪请求流过整个微服务架构",靠的是两个补丁机制(见 app.py 与 app.py):
LiveCheckSensor.on_stream_event_in:当流开始处理事件时,从事件 headers 中解析TestExecution(TestExecution.from_headers),压入current_test_stack;on_stream_event_out在事件处理结束时弹出;LiveCheck.on_produce_attach_test_headers:通过连接 Faust 的on_produce_message信号,在当前测试上下文存在时,把测试的LiveCheck-Test-Id、LiveCheck-Test-Name、LiveCheck-Test-Timestamp、LiveCheck-Test-Expires四个头附加到外发的 Kafka 消息上。
这样一来,测试上下文会沿着消息传递链一路传播,下游系统也能感知"当前正在处理一次测试"。TestExecution的as_headers()/from_headers()负责头与对象之间的序列化(见 models.py)。
五、触发机制:概率采样与假请求
5.1 基于概率的触发:maybe_trigger
case.py 提供了maybe_trigger异步上下文管理器:每次进入时以uniform(0, 1) < self.probability决定是否真正触发执行。若触发,则创建一个TestExecution并把它发送到pending_tests主题,同时将执行对象压入current_test_stack。这就实现了"100% / 50% / 0.1% 采样率"。
async with self.maybe_trigger(id='my-test-1', order_id=123): ... # 只有概率命中的执行才会进入此块5.2 立即触发:trigger
trigger(id, *args, **kwargs)(见 case.py)无条件创建TestExecution:包含唯一id(默认faust.utils.uuid生成)、用例名、时间戳、测试参数、过期时间,然后通过self.app.pending_tests.send(key=id, value=execution)发布到livecheck主题。
5.3 定期假请求:_send_frequency 与 make_fake_request
当frequency被设置时,@Service.task修饰的_send_frequency会按间隔运行定时器,且仅在当前 LiveCheck 应用是 leader 时调用make_fake_request()(见 case.py)。make_fake_request是一个空占位方法(...),子类应覆盖它来向被测系统发送合成请求,例如:
async def make_fake_request(self): await self.post_url( 'http://localhost:6066/order/init/sell/', json={'symbol': 'AAPL'}, )配合get_url/post_url/url_request(见 case.py),用例可以发起带超时与指数退避重试的 HTTP 请求:ClientError时按url_error_delay_min起步、每次乘以url_error_delay_backoff、封顶url_error_delay_max重试,超过url_error_retries次仍未成功则抛出ServiceDown并转入on_suite_fail——这正是"被动监控子系统宕机"能力的来源。
六、测试执行与报告:TestRunner 与 TestExecution
6.1 TestExecution:一次测试执行的载体
TestExecution(models.py)记录一次测试执行的完整元数据:id、case_name、timestamp、test_args、test_kwargs、expires。它还提供:
ident/shortident:日志用长/短标识符(case_name:id);human_date/was_issued_today:人类可读的时间描述;is_expired:判断执行是否已过期(当前时间超过expires)。
6.2 TestRunner:执行的调度中枢
faust/livecheck/runners.py 的TestRunner.execute()定义了完整的执行流程:
- 压入
current_test_stack; - 若
case.active为假 → 跳过(case inactive);若test.is_expired→ 跳过(expired); - 通过
_prepare_args/_prepare_kwargs对参数调用maybe_model进行模型解析(支持把原始数据解析成 Faust Record); on_start→ 记录开始时间并调用case.on_test_start(更新频率/延迟历史);- 执行
case.run(*args, **kwargs),按异常类型分流:AssertionError→on_failed→ 状态FAIL,抛出TestFailed;TestSkipped→on_skipped→ 状态SKIP;TestTimeout→on_timeout→ 状态TIMEOUT;- 其他异常 →
on_error→ 状态ERROR,抛出TestRaised; - 正常返回 →
on_pass→ 状态PASS,输出Test OK in ~x.x seconds;
- 最终
_finalize_report组装TestReport(含case_name、state、runtime、signal_latency、error、traceback)并调用case.post_report发布到livecheck-report主题。
realtime_logs机制也在这里体现:为False时执行中的日志先缓冲在runner.logs,on_pass时统一冲刷输出;为True时实时输出(见 runners.py)。
6.3 报告的发布
case.post_report最终调用self.app.post_report(report)(app.py),以report.test.id为 key 发布到livecheck-report主题。报告可以作为告警源,也可以接入监控面板。
七、停滞检测与套件故障
7.1 停滞检测:_check_frequency
case.py 中@Service.task修饰的_check_frequency是停滞哨兵:启动后先sleep(warn_stalled_after),随后按相同间隔轮询。若距last_test_received(或启动时刻)超过warn_stalled_after秒仍无任何测试活动,则抛出SuiteStalled(消息形如Test stalled! Last received ... ago (warn_stalled_after=...)),并调用on_suite_fail(exc, State.STALL)。注意该异常被捕获后并不传播,哨兵任务持续运行;同时在未停滞时调用_maybe_recover_from_failed_state()恢复状态。
7.2 套件故障异常体系
faust/livecheck/exceptions.py 定义了完整的异常层级:
| 异常 | 父类 | 触发场景 |
|---|---|---|
LiveCheckError | Exception | 通用基类 |
SuiteFailed | LiveCheckError | 整个套件失败 |
ServiceDown | SuiteFailed | 依赖的 HTTP 服务无响应(重试耗尽) |
SuiteStalled | SuiteFailed | 长时间无测试执行 |
TestSkipped | LiveCheckError | 测试被跳过 |
TestFailed | LiveCheckError | 断言失败 |
TestRaised | LiveCheckError | 执行抛出未预期异常 |
TestTimeout | LiveCheckError | 等待信号或执行超时 |
八、把 Case 接入 LiveCheck 应用并运行
8.1 应用集成
LiveCheck 应用本身是faust.App的子类(app.py),并为目标应用提供LiveCheck.for_app(app)工厂方法,自动注入LiveCheckMiddleware(Web 中间件,用于从 HTTP 请求头恢复测试上下文)、LiveCheckSensor,并将app.livecheck指向自身。
case()装饰器是注册用例的标准方式(app.py),其签名覆盖了前文Case的全部可配置参数,且默认warn_stalled_after=timedelta(minutes=30):
from faust.livecheck import LiveCheck livecheck = LiveCheck('my-livecheck', broker='kafka://localhost:9092') @livecheck.case(name='orders.flow', probability=0.5, frequency=30.0) class OrderFlowCase(Case): ...应用内置三个 Kafka 主题(默认名称):livecheck(待执行测试)、livecheck-bus(信号通信)、livecheck-report(测试报告),并发度分别由test_concurrency=100、bus_concurrency=30控制(app.py)。
8.2 运行示例
参照 docs/userguide/livecheck.rst 的教程,使用仓库自带的 examples/livecheck.py 可以快速体验完整流程:
终端一:启动被测应用(订单系统 worker):
$ python examples/livecheck.py worker -l info终端二:为该应用启动 LiveCheck 实例:
$ python examples/livecheck.py livecheck -l info产生真实流量(示例中测试执行概率为 50%,因此至少需要触发两次才能看到 LiveCheck 终端出现执行活动):
$ python examples/livecheck.py post_order --side=sell或直接访问浏览器:
http://localhost:6066/order/init/sell/
8.3 测试用例的单元测试
仓库在 t/unit/livecheck/test_case.py 中为Case提供了覆盖全面的单元测试,包括构造参数、_sampler采样、maybe_trigger概率触发、trigger、execute、on_test_start频率统计、_check_frequency停滞检测、on_suite_fail套件失败等场景,是理解Case各方法行为边界的绝佳参考资料。
九、结语:Case 设计要点回顾
编写生产环境测试用例时,请牢记Case的设计精髓:
- 继承并实现
run():run是唯一必须覆盖的方法,里面用断言声明不变量; - 用
Signal做跨服务同步:等待外部事件完成,而不是盲目 sleep; - 用
probability控制采样率:流量大时调低、流量小时调高,避免对生产造成压力; - 用
frequency+make_fake_request制造假流量:当真实流量不足以满足warn_stalled_after时主动喂数据; - 理解状态机与告警节流:
state_transition_delay和max_consecutive_failures共同决定了"什么时候告警、什么时候恢复"; - 通过报告主题集成监控:
livecheck-report主题上的TestReport可直接接入告警管道。
从 case.py 的源码可以看出,Case将"测试编写"与"生产运行"两层关注点彻底解耦:你只需要像写单元测试一样写断言,剩下的概率采样、请求追踪、停滞检测、故障恢复与报告发布,全部由框架在背后完成。
- 流处理
- 消息队列
- 后端
【免费下载链接】faust
Python Stream Processing
相关推荐
Continue VS Code 扩展端到端(E2E)测试指南:从环境搭建到测试用例编写
Continue VS Code 扩展端到端(E2E)测试指南:从环境搭建到测试用例编写 导读 本文基于 Continue 开源仓库中 VS Code 扩展的
人工智能AI Agent代码智能体开发工具工具调用RAGTraefik Proxy Helm Chart监控与可观测性:Prometheus+Grafana完美集成指南
Traefik Proxy Helm Chart监控与可观测性:Prometheus+Grafana完美集成指南 Traefik Proxy是一款现代化的云原生
Eclipse Che端到端测试终极指南:如何编写高质量E2E测试用例
Eclipse Che端到端测试终极指南:如何编写高质量E2E测试用例 Eclipse Che作为基于Kubernetes的企业级云开发环境,其稳定性和可靠性至
开发工具云原生
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考