首页
/ Pathway CI/CD 集成实践:用 pytest、mypy 与 demo 回放 API 为流式管道建立可靠测试链路

Pathway CI/CD 集成实践:用 pytest、mypy 与 demo 回放 API 为流式管道建立可靠测试链路

2026-09-04 22:49:51作者:姚月梅Lane

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 加入既有项目,都可以自由选择最合适的集成方式;
  • 它天然兼容 mypypytest 等主流工具,可以无缝接入 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_defsstrict_equality 等则保证对已标注部分保持严格。

格式化与 lint[tool.black] 固定 required-version = "24"line-length = 88、排除 docstarget 目录;[tool.isort] 使用 profile = "black" 并声明 known_first_party = ["pathway"];flake8 的配置在 setup.cfg 中:max-line-length = 119docstring-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)

几个关键实践点:

  1. 临时文件 + 高 input_rate:用 tmp_path 落盘小 CSV,input_rate=1000 让回放瞬间完成,测试不依赖网络与真实时钟;
  2. 确定值断言generate_custom_stream 的用例(L11-L36)直接列出 5 行期望值;对随机性强的列(如 noisy_linear_streamy)则只断言确定性列 x
  3. 时间回放与异常路径test_demo_replay_with_timeunit="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.pypyproject.toml 生成 requirements.txt 后再安装,stub 包单独安装。如果你的项目依赖较多,这是保持 mypy job 可复现的实用手法。

构建 + 端到端测试:package_test.yml

.github/workflows/package_test.yml 是一个可被复用的 workflow_call 流水线(由 ubuntu_test.ymlrunner: 'ubuntu-22.04' 调用),体现了"先构建 wheel,再在干净环境中安装并全量测试"的完整部署验证模式:

  1. 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

  2. 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 模块里的用法示例)也成为测试的一部分;

  3. Notify_on_failure:构建或测试失败时通过 webhook 推送 PR 地址与运行地址到 Slack。

对最终用户而言,这条链路证明了:安装发布物 → 跑完整 pytest 套件,可以在任何支持 Python 的 CI runner 上复现,且不需要本地 Rust 工具链(wheel 是预构建的)。

在自己的项目中落地的检查清单

把上述实践迁移到你自己的 Pathway 项目,可以按此顺序建立 CI/CD:

  1. 测试依赖:把 pathway[tests](或按需挑选 pytestpytest-timeoutpytest-xdist)加入你的 dev 依赖;
  2. 单元测试:用 pw.demo.replay_csv / replay_csv_with_time 构造确定性输入,配合 tmp_path 与高 input_rate,对每个 revision 的聚合结果做快照断言;
  3. 时间语义:对涉及窗口、水位线、迟到数据的管道,用带时间列的回放 CSV 覆盖"间隔、乱序、删除"三类样本;
  4. 静态检查:复制本仓库的 black / isort / flake8 / mypy 配置要点(ignore_missing_importscheck_untyped_defs、flake8 的 per-file-ignores);
  5. 构建验证:如果是分发 wheel,参考 package_test.yml 的"maturin 构建 → 工件传递 → 干净 venv 安装 → pytest"四段式;
  6. 防挂起:为流式测试统一配置 pytest-timeout(仓库内取 900 秒),避免回放或连接器异常拖死流水线。

小结

Pathway 的 CI/CD 集成的关键不在于引入特殊工具,而在于两点:标准 Python 工具链(pytest、mypy、black、flake8、isort,加上 Rust 侧的 fmt/clippy/test)全部可直接使用;流式输入的非确定性由 pw.demo 的 session-replay 机制(generate_custom_streamrange_streamnoisy_linear_streamreplay_csvreplay_csv_with_time)转化为确定性输入。本仓库的 pull.ymlpackage_test.yml 展示了从静态检查到 wheel 构建再到全量 pytest 的完整参照实现,test_demo.py 则提供了可复制的回放测试写法——按此组合,你的 Pathway 管道就能像任何 Python 项目一样被自动测试、构建和部署。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
528
588
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
906
1.83 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.53 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.34 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
987
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384