Pathway CI/CD 集成实践:用 pytest、mypy 与 demo 回放 API 为流式管道建立可靠测试链路
Pathway Live Data Framework 是完全 Python 兼容的流处理框架,其管道可以像任何普通 Python 工作流一样被自动测试、构建和部署。本文基于官方部署文档 CI/CD 展开,结合本仓库真实的 GitHub Actions 工作流、demo 模块源码与测试用例,讲清三件事:Pathway 项目如何接入 mypy / pytest 等标准 Python 工具链,如何用内置的 session-replay 机制(pw.demo 等)在标准 CI 中可预测地重现流式场景(包括迟到、乱序的时间语义),以及一套可直接参考的端到端构建-测试流水线长什么样。
为什么 Pathway 管道可以走标准 CI/CD 链路
官方集成文档给出的核心结论是:
- Pathway 是完全 Python 兼容的,可以使用任何你熟悉的 Python 集成工具;
- 无论新建项目还是把 Pathway 加入既有项目,都可以自由选择最合适的集成方式;
- 它天然兼容 mypy、pytest 等主流工具,可以无缝接入 CI/CD 链路(Jenkins、GitHub Actions 等);
- 可以在本地或任何带 Python 的 CI 流水线中,用离线数据片段运行测试;
- 测试可以覆盖数据的时间语义(迟到、乱序到达),方法是比较多个 revision(数据快照版本)上的结果;
- Pathway 提供多种 session-replay 机制,如 demo API,允许在标准 CI/CD 中可预测地重现流式场景。
从工程角度看,这意味着两件事:第一,Pathway 管道是纯 Python 代码,静态检查、类型检查、单测框架全部直接可用;第二,流式系统最大的测试难点——"输入是非确定性的实时流"——被官方通过回放(replay)机制消解掉了,回放输入是确定性的,断言结果才是可重复的。
工具链全景:pytest、mypy、格式化与 Rust 侧检查
本仓库 pyproject.toml 定义了完整的工具配置,可以作为你项目 CI 的直接参照。
依赖与测试依赖。[project.optional-dependencies] 中的 tests extra 声明了官方测试所需的完整工具栈(pyproject.toml):
tests = [
"pytest >= 9.0.3, < 10.0.0",
"pytest-xdist >= 3.8.0, < 4.0.0", # 并行执行
"pytest-rerunfailures >= 16.2, < 17.0", # 失败重试
"pytest-asyncio >= 1.2.0",
"pytest-timeout >= 2.3.0, < 3.0.0",
...
"pathway[all]",
]
其中 pytest-xdist 用于并行加速(CI 中实际以 --dist worksteal -n auto 启动),pytest-timeout 用于防止流式管道中潜在的挂起导致流水线永久卡住。
mypy 配置(pyproject.toml)值得特别注意,它展示了如何对"含 Rust 原生扩展 + 大量第三方库"的项目做类型检查:
[tool.mypy]
python_version = "3.11"
exclude = ["(^|/)target/", "(^|/)examples/", '(^|/)tests(/.*)?/test_.*\.py$']
ignore_missing_imports = true
check_untyped_defs = true
warn_redundant_casts = true
warn_unused_ignores = true
strict_equality = true
ignore_missing_imports = true 让原生扩展(pathway.engine 由 Rust 编译)和缺 stub 的第三方库不阻塞检查,而 check_untyped_defs、strict_equality 等则保证对已标注部分保持严格。
格式化与 lint。[tool.black] 固定 required-version = "24"、line-length = 88、排除 docs 与 target 目录;[tool.isort] 使用 profile = "black" 并声明 known_first_party = ["pathway"];flake8 的配置在 setup.cfg 中:max-line-length = 119、docstring-convention = google,并对 python/pathway/internals/api.py 单独放行 F403, F405(该文件做了通配符再导出)。
核心机制:pw.demo 模块与 session replay
文档中提到的 "session-replay 机制" 在本仓库中最主要的落地就是 python/pathway/demo/init.py。模块 docstring 明确其定位:"This feature empowers you to effectively test and debug your Pathway implementation using realtime data"。它提供 5 个入口函数,全部构建在 pw.io.python.read(自定义 ConnectorSubject 连接器)之上,因此回放出来的表与普通流输入在图执行层面完全同构——这正是回放测试有效的前提:
generate_custom_stream:按函数生成流
pw.demo.generate_custom_stream(
value_generators, # 列名 -> 值生成函数(按行索引 x 取值)
schema=InputSchema, # pw.Schema
nb_rows=None, # 行数;None 表示无限流
autocommit_duration_ms=1000, # 两次提交的最大间隔
input_rate=1.0, # 每秒生成行数
)
从源码实现看(demo/__init__.py),它内部定义了一个匿名 ConnectorSubject 子类,run() 中按 input_rate 休眠后调用 self.next_json(row) 逐行推送,再交给 pw.io.python.read(..., format="json") 读入。input_rate 在测试中可设得很高(如 1000 行/秒),让回放"快进"完成,避免 CI 等待。
range_stream 与 noisy_linear_stream:内置教学流
pw.demo.range_stream(nb_rows=30, offset=0, input_rate=1.0)生成单列value的递增流,常用于验证求和等聚合;pw.demo.noisy_linear_stream(nb_rows=100, input_rate=2.0)生成带确定性随机噪声(random.seed(0))的(x, y)线性流,专供线性回归类教程。两者本质都是对generate_custom_stream的封装。
replay_csv:把静态 CSV 变成流
class InputSchema(pw.Schema):
k: int
v: str
table = pw.demo.replay_csv("input_stream.csv", schema=InputSchema, input_rate=1000)
实现中,autocommit_ms 会按 int(1000.0 / input_rate) 换算,即"行间隔"直接映射为"提交间隔";读取时先以全 str 类型读入再 cast_to_types(**schema.typehints()) 转回目标类型,因此 CSV 内容必须符合标准设定(分隔符 ,、引号 "、无转义)。
replay_csv_with_time:按时间戳语义回放
这是覆盖"时间行为"测试的关键工具。函数签名:
pw.demo.replay_csv_with_time(
path, # CSV 路径
schema=..., # 时间列必须声明为 int 或 float,否则抛 ValueError
time_column="time", # 时间戳列(需为有序正整数)
unit="s", # 时间戳单位:'s' / 'ms' / 'us' / 'ns'
autocommit_ms=100, # 提交间隔
speedup=1, # 相对时间戳加速倍速
)
其原理(见 L311-L330):记录首行的 time_column 值与真实开始时刻,之后每行根据"时间戳差 / speedup"与真实经过时间做差,正差则 time.sleep(tts)。也就是说,它会忠实还原数据之间的时间间隔(含 speedup 加速),从而可以在 CI 中重现"两行数据相隔 2 秒"这类依赖真实时间流逝的场景。
时间语义(迟到、乱序)测试怎么写
官方文档指出:测试可以通过"比较多个 revision 的结果"来覆盖迟到、乱序等时间行为。本仓库的 python/pathway/tests/test_demo.py 展示了这套测试的标准写法:
def test_demo_replay(tmp_path: pathlib.Path):
data = """
k | v
1 | foo
2 | bar
3 | baz
"""
input_path = tmp_path / "input.csv"
write_csv(input_path, data)
class InputSchema(pw.Schema):
k: int
v: str
table = pw.demo.replay_csv(str(input_path), schema=InputSchema, input_rate=1000)
expected = T(data)
assert_table_equality_wo_index(table, expected)
几个关键实践点:
- 临时文件 + 高 input_rate:用
tmp_path落盘小 CSV,input_rate=1000让回放瞬间完成,测试不依赖网络与真实时钟; - 确定值断言:
generate_custom_stream的用例(L11-L36)直接列出 5 行期望值;对随机性强的列(如noisy_linear_stream的y)则只断言确定性列x; - 时间回放与异常路径:
test_demo_replay_with_time用unit="ns"、autocommit_ms=10做带时间列的回放;test_demo_replay_with_time_wrong_schema则验证时间列声明为str时会抛出ValueError("Invalid schema. Time columns must be int or float."),与 源码中的校验 一一对应。
对于更复杂的"迟到/乱序"场景(例如主键更新先于插入到达、删除乱序),可以把乱序样本写成 CSV,用 replay_csv / replay_csv_with_time 逐 revision 回放,再对每个 revision 的聚合结果做快照断言——输入完全确定,输出即可重复比对。
参考实现:本仓库自身的 CI 流水线
本仓库的 CI 配置可直接作为"Pathway 项目接入 CI/CD"的完整范例。
静态检查与单测:pull.yml
.github/workflows/pull.yml 在 PR 与 push 到 main 时触发,Python 版本固定 3.10、Rust 工具链 1.97.1,包含 7 个并行 job:
| Job | 命令 | 对应配置 |
|---|---|---|
| python-black | black . --check(black 24.x) |
[tool.black] |
| python-flake8 | flake8 |
setup.cfg |
| python-mypy | 先执行 .github/ci/get-dependencies.py 同步依赖并安装 type stubs(types-pytz、types-PyYAML、pandas-stubs 等),再 mypy . |
[tool.mypy] |
| python-isort | isort --check . |
[tool.isort] |
| cargo-fmt | cargo fmt -- --check |
— |
| cargo-clippy | cargo clippy --locked --all-targets -- --deny warnings |
— |
| cargo-test | cargo test --locked |
— |
注意 mypy job 中的两步准备(L46-L63):依赖由 .github/ci/get-dependencies.py 从 pyproject.toml 生成 requirements.txt 后再安装,stub 包单独安装。如果你的项目依赖较多,这是保持 mypy job 可复现的实用手法。
构建 + 端到端测试:package_test.yml
.github/workflows/package_test.yml 是一个可被复用的 workflow_call 流水线(由 ubuntu_test.yml 以 runner: 'ubuntu-22.04' 调用),体现了"先构建 wheel,再在干净环境中安装并全量测试"的完整部署验证模式:
-
Build_packages:清理工作区 →
PyO3/maturin-action以--release --strip构建 wheel(Linux 走 manylinux 容器,macOS 构建 universal2 并用delocate-wheel修补依赖)→ 上传target/wheels/工件;版本号从Cargo.toml推导为x.y.(z+1)-dev${BUILD_NUMBER}并写入python/pathway/internals/version.py; -
pytest:矩阵覆盖 Python 3.10 与 3.14,下载工件后用 uv 创建 venv 并
uv pip install "${WHEEL}[tests]",然后运行:python -m pytest -v --confcutdir "${ENV_NAME}" --doctest-modules --pyargs pathway其中
PYTEST_ADDOPTS="--dist worksteal -n auto --timeout=900"提供了并行与 15 分钟级超时保护,--doctest-modules让模块 docstring 中的>>>示例(如 demo 模块里的用法示例)也成为测试的一部分; -
Notify_on_failure:构建或测试失败时通过 webhook 推送 PR 地址与运行地址到 Slack。
对最终用户而言,这条链路证明了:安装发布物 → 跑完整 pytest 套件,可以在任何支持 Python 的 CI runner 上复现,且不需要本地 Rust 工具链(wheel 是预构建的)。
在自己的项目中落地的检查清单
把上述实践迁移到你自己的 Pathway 项目,可以按此顺序建立 CI/CD:
- 测试依赖:把
pathway[tests](或按需挑选pytest、pytest-timeout、pytest-xdist)加入你的 dev 依赖; - 单元测试:用
pw.demo.replay_csv/replay_csv_with_time构造确定性输入,配合tmp_path与高input_rate,对每个 revision 的聚合结果做快照断言; - 时间语义:对涉及窗口、水位线、迟到数据的管道,用带时间列的回放 CSV 覆盖"间隔、乱序、删除"三类样本;
- 静态检查:复制本仓库的 black / isort / flake8 / mypy 配置要点(
ignore_missing_imports、check_untyped_defs、flake8 的 per-file-ignores); - 构建验证:如果是分发 wheel,参考
package_test.yml的"maturin 构建 → 工件传递 → 干净 venv 安装 → pytest"四段式; - 防挂起:为流式测试统一配置
pytest-timeout(仓库内取 900 秒),避免回放或连接器异常拖死流水线。
小结
Pathway 的 CI/CD 集成的关键不在于引入特殊工具,而在于两点:标准 Python 工具链(pytest、mypy、black、flake8、isort,加上 Rust 侧的 fmt/clippy/test)全部可直接使用;流式输入的非确定性由 pw.demo 的 session-replay 机制(generate_custom_stream、range_stream、noisy_linear_stream、replay_csv、replay_csv_with_time)转化为确定性输入。本仓库的 pull.yml 与 package_test.yml 展示了从静态检查到 wheel 构建再到全量 pytest 的完整参照实现,test_demo.py 则提供了可复制的回放测试写法——按此组合,你的 Pathway 管道就能像任何 Python 项目一样被自动测试、构建和部署。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00