将 S3/MinIO + Delta Lake 当作消息队列:Pathway 的 Kafka 替代方案延迟基准测试完整指南
本文是一份以 examples/projects/kafka-alternatives/benchmarks 代码库为主体的实操指南,讲解如何在不引入 Kafka 的前提下,用 Pathway + Delta Lake on S3/MinIO 构建消息队列,并通过官方基准脚本测量端到端延迟。读完本文,你将掌握基准仓库四个脚本各自的职责、环境配置方法、benchmark.py 的完整命令行参数语义,以及延迟百分位(p50/p75/p85/p95/p99)的计算与结果 CSV 的解读方式,从而能够在你自己的 S3/MinIO 环境里复现并扩展这套"以 Delta Lake 为消息队列"的性能测试。
这套基准在验证什么
Kafka 之所以被广泛用作消息中间件,是因为它提供了持久化、可回溯、可多消费者订阅的流式消息通道;但它也带来了 broker、zookeeper 集群等高运维成本。这套基准(benchmark)想要验证的是一个更轻的架构假设:在已有的 S3 兼容对象存储上铺一层 Delta Lake,由 Pathway 负责"持续监控 + 读写" Delta Table,能否获得可用的消息队列延迟。仓库根目录的 benchmark 说明 将该基准定位为 "Benchmark for Delta Lake S3 messaging as a Kafka replacement",其详细思路(架构、实验设计、结果与调优结论)见配套模板文章 Pathway 作为 Python Kafka 替代方案。
整个基准的核心逻辑是:生产者以受控速率把消息写进位于 S3 的 Delta Table,消费者以流式模式从同一个 Delta Table 读回消息,并记录每条消息从"提交时刻"到"被处理时刻"的延迟。仓库中还有一套更贴近真实业务的配套示例 minio-ETL,它演示的是跨时区时间戳的 ETL 处理;而本文聚焦于纯消息吞吐/延迟的基准部分。
基准仓库文件职责一览
| 文件 | 职责 |
|---|---|
| lib.py | 环境配置常量的集中定义,以及两个可复用的 Pathway 封装类 AdhocProducer / AdhocConsumer,它们把 Delta Table 封装成"消息队列"的读端与写端 |
| producer.py | 生产者命令行入口:按指定"每秒消息数"生成消息,通过 AdhocProducer 写入 Delta Lake 消息队列 |
| consumer.py | 消费者命令行入口:通过 AdhocConsumer 从 Delta Table 读消息并计时,输出 p50/p75/p85/p95/p99 百分位延迟报告 |
| benchmark.py | 总调度器:按"消息数/秒"区间依次拉起 producer 与 consumer 子进程,并把结果写入 benchmark-results/ 目录 |
准备工作:配置 S3/MinIO 凭据与常量
基准需要指定 S3 bucket 名称并提供两组访问密钥。这些设置全部集中在 lib.py 顶部的常量区,运行时不需要改动脚本调用逻辑。
TEST_REGION = "eu-central-1"
TEST_MINIO_BUCKET_NAME = "your-bucket"
TEST_MINIO_ENDPOINT = "https://your-endpoint.com"
TEST_MINIO_SETTINGS = pw.io.minio.MinIOSettings(
access_key=os.environ["MINIO_S3_ACCESS_KEY"],
secret_access_key=os.environ["MINIO_S3_SECRET_ACCESS_KEY"],
bucket_name=TEST_MINIO_BUCKET_NAME,
region=TEST_REGION,
endpoint=TEST_MINIO_ENDPOINT,
)
TEST_S3_BUCKET_NAME = "your-bucket"
TEST_S3_SETTINGS = pw.io.s3.AwsS3Settings(
access_key=os.environ["AWS_S3_ACCESS_KEY"],
secret_access_key=os.environ["AWS_S3_SECRET_ACCESS_KEY"],
bucket_name=TEST_S3_BUCKET_NAME,
region=TEST_REGION,
)
需要替换的关键占位值有三个:TEST_BUCKET_NAME(bucket 名)、TEST_ENDPOINT(MinIO 服务的 endpoint URL,使用 AWS S3 时可为空)、TEST_REGION。而访问密钥并非写死在源码里,而是从环境变量读取:
- MinIO 后端读取
MINIO_S3_ACCESS_KEY与MINIO_S3_SECRET_ACCESS_KEY; - AWS S3 后端读取
AWS_S3_ACCESS_KEY与AWS_S3_SECRET_ACCESS_KEY。
两套后端通过 get_s3_backend_settings 统一切换:minio 分支会对 MinIOSettings 调用 create_aws_settings() 将其转成通用的 AwsS3Settings,因此后续所有读写逻辑都只需面对一套 S3 连接配置。启动前请务必先在对象存储中创建好 bucket——与配套 minio-ETL 示例说明 中强调的一致,写入不存在的 bucket 会导致管线失败。
封装层解读:lib.py 如何把 Delta Table 变成消息队列
lib.py 是整个基准的地基。它的两个类分别承担"队列写入"与"队列订阅"角色,非常直观地展示了 Pathway + Delta Lake 的消息队列用法。
AdhocProducer:以受控速率写入 Delta Table
AdhocProducer 的构造过程分三步:
- 定义一个
data: bytes的_MessageTableSchema(消息负载统一为字节串,schema 极简,使测量聚焦于传输而非反序列化); - 通过
pw.io.python.read(...)从自定义的ConnectorSubject子类_MessageTableSubject读取外部消息,autocommit_duration_ms决定外部输入被提交给引擎的间隔; - 通过
pw.io.deltalake.write(table=..., uri=base_path, s3_connection_settings=..., min_commit_frequency=autocommit_duration_ms)把表写向 Delta Lake。
关键点是 write 端的 min_commit_frequency 被显式设置成与 read 端相同的 autocommit_duration_ms。在 deltalake 写连接器 的默认定义中,min_commit_frequency 默认高达 60 000 ms;不显式调小它,消息会被长时间积压在内存批次里,延迟自然无从谈起。基准把两个参数绑定,是为了让"输入提交"与"Delta 提交"的节奏一致,便于观察纯粹的存储链路延迟。
AdhocConsumer:订阅 Delta Table 的新增数据
AdhocConsumer 是队列读端的封装,构造时做三件事:
- 解析
s3://bucket/path形式的base_path,拆出 bucket 名与桶内路径,用于后续轮询探测(若连接设置里已带 bucket 则复用); - 用
pw.io.deltalake.read(uri=..., mode="streaming", s3_connection_settings=..., autocommit_duration_ms=...)以流式模式持续读 Delta Table 的新增版本。相比快照模式,mode="streaming"意味着引擎会持续轮询表的新版本并把增量更新推向下游; - 通过
pw.io.subscribe(table, on_change=...)注册回调,把每条新增/更新(is_addition=True的行)交给用户回调处理。
此外,消费者在 start() 时会先执行 _wait_for_delta_table_creation:用 boto3 客户端以 0.1 秒间隔轮询 list_objects_v2,直到 Delta Table 对应前缀下出现对象为止。这是为了处理生产者还没创建表、消费者就先行启动的竞态——基准调度中消费者先于生产者启动,因此必须先等表出现再进入 pw.run(),否则流式读会因目标表尚不存在而失败。
从源码结构看,lib.py 刻意让消息的读写都只经过 Pathway 的 Delta Lake 连接器与 subscribe 回调,全程没有引入任何 broker 组件——这正是"Delta Lake S3 消息队列"这一方案的核心形态。
生产者如何精确控制流速:producer.py
producer.py 的 CLI 参数为 --lake-path、--rate、--seconds-to-stream、--s3-backend(minio/s3)与 --autocommit-duration-ms。其核心是 create_producer 构造的闭包,它维护了一个简单的"节流环":
- 期望节奏为
duration_per_message = 1.0 / rate_per_second; - 消息内容是当前墙上时钟的时间戳字符串(每 100 条刷新一次),编码为 UTF-8 字节后调用
sender.next(data=message); - 每隔
delay_check_per_iterations = 10000条做一次"计划时间 vs 实际时间"比对:落后于计划则打印Streaming falls behind by ...s;领先则调用time.sleep(...)补足差额并打印Streaming is faster than target by ...s。
这套节流逻辑保证了发往 Delta Table 的速率在统计上是稳定的——只有生成速率确定,消费端测出的延迟才可归因于存储/轮询链路,而不是生产者抖动。消息体携带提交时刻时间戳,正是消费者计算延迟的基准(见下节)。
消费者如何统计延迟百分位:consumer.py
consumer.py 开头定义了三个影响测量的关键常量:
PERCENTILES = [50, 75, 85, 95, 99]
MAX_LATENCY = 100
WARMUP_PERIOD = 30
PERCENTILES:报告输出 50/75/85/95/99 五个百分位;MAX_LATENCY:延迟上限 100 秒,超过的样本会被截断到 100 秒,用于限制直方图桶数;WARMUP_PERIOD:前 30 秒的样本在渲染报告时被跳过。
回调 consume_messages 将每条消息解码回提交时间戳,用 当前时间 - 提交时间 得到单条延迟,记录进 timeline;当处理数量达到 expected_messages 时调用 render_latencies_report 落盘并退出进程。
报告的生成方式是流式的累积直方图:把延迟乘以 100 取整作为桶下标(精度 0.01 秒),样本累积过程中每隔 log_frequency 条输出一行快照,按"当前已处理消息数 × 目标百分位 / 100"的边界判定各百分位对应的延迟值。因此每个速率产生的 CSV 不是一个数字,而是一份"随时间演化的百分位报告",可以画出类似配套文章中 latency-per-rate 的延迟曲线。这种设计还天然过滤了启动阶段的冷启动样本:生产者先于消费者写入的堆积会被 30 秒预热期排除,保证测量的是稳态延迟。
总调度器 benchmark.py:参数与执行流程
benchmark.py 是唯一需要你直接运行的脚本。除原始文档列出的四个参数外,源码还提供了两个可选参数(默认值见下),实际可用参数如下:
| 参数 | 类型 | 必填 | 默认值 | 含义 |
|---|---|---|---|---|
--range-start |
int | 是 | — | 待测"每秒消息数"区间的起始值 |
--range-end |
int | 是 | — | 区间的结束值(含端点) |
--range-step |
int | 是 | — | 区间内的递增步长 |
--seconds-to-stream |
int | 否 | 300 |
每个速率档位的持续流式秒数 |
--s3-backend |
str | 否 | minio |
后端类型,可选 s3 或 minio |
--autocommit-duration-ms |
int | 否 | 1000 |
read/write 两侧的提交间隔(毫秒),同时传给 producer 与 consumer |
每个速率档位的执行流程(见 benchmark.py 主循环)依次是:
- 构造临时 S3 前缀
adhoc/s3-messaging/{rate},完整路径形如s3://{bucket}/adhoc/s3-messaging/{rate}; - 调用 cleanup_files_under_prefix 清理该前缀下上一轮遗留对象,保证测试从空表开始;
- 先启动 consumer 子进程(
expected_messages = rate × seconds_to_stream,log_frequency = rate // 2),再启动 producer 子进程; - 用
assert等待两个子进程均以退出码 0 结束; - 再次清理 S3 前缀,进入下一个速率档位。
需要注意,调度代码是用 subprocess.Popen(["python", "consumer.py", ...]) 拉起子进程的,因此运行时的工作目录必须位于 benchmarks/ 目录内,否则 Python 无法定位 consumer.py / producer.py / lib.py。
运行示例
文档给出的示例是同时测 10 000、20 000、30 000 三档速率,每档持续 10 分钟(600 秒):
python benchmark.py --range-start 10000 --range-end 30000 --range-step 10000 --seconds-to-stream 600
以 MinIO 为后端、每档 5 分钟、提交间隔调为 100 ms 的变体可以写成:
python benchmark.py --range-start 10000 --range-end 30000 --range-step 10000 \
--seconds-to-stream 300 --s3-backend minio --autocommit-duration-ms 100
结果输出:benchmark-results 目录
脚本会在当前目录下创建/写入 benchmark-results/,每档速率生成一个 CSV 文件,命名为 {rate}.csv(例如 benchmark-results/10000.csv)。每个文件的格式为:
time_from_start,p50,p75,p85,p95,p99
<秒>,<延迟秒>,...
首列 time_from_start 是从消费者启动起的运行秒数,后续列分别对应该时刻下 p50/p75/p85/p95/p99 的端到端延迟(单位:秒,保留两位小数)。由于消费者每处理 rate // 2 条就输出一行快照、最后再追加一行最终值,因此 CSV 末行即该档速率的整体稳态百分位结果。若想画延迟-时间曲线,直接用这些行即可;若只要一个最终数字,取末行即可。
参数语义与调优方向:autocommit_duration_ms / min_commit_frequency
阅读配套文章 180.kafka-alternative.md 并结合本仓库源码,可以明确两个核心参数的工程含义:
autocommit_duration_ms:Pathway 输入/输出侧"最多多久强制提交一次"。在 deltalake.read 中默认值为1500ms。调小它,消息在内存中等待的时间变短,延迟随之下降;min_commit_frequency:Delta Lake 写侧两次数据提交之间的最小间隔。在 deltalake.write 中默认值为60_000ms——注意这个默认值很大,若不显式调小,写入方会长时间攒批。基准中AdhocProducer刻意把它与autocommit_duration_ms绑定,而benchmark.py的--autocommit-duration-ms默认取1000。
两者的调优存在明确权衡(这一结论来自配套文章 延迟与调优章节,本文仅转述仓库内已记录的结论):减小批大小/提交间隔能显著降低中低速率下的延迟,但代价是向对象存储发起更多次网络往返与 Delta 版本轮询,因此高速率下吞吐会先于延迟被网络开销限制。配套文章给出的数据为:将相关间隔调到 100 ms 量级后,10 000~30 000 条/秒区间可获得亚秒级延迟(例如 10 000 条/秒时 p50 ≈ 0.26 s、p99 ≈ 0.67 s);而在 70 000 条/秒以上,为换取低延迟所增加的网络调用会反过来拖垮系统,延迟曲线显著恶化。因此:做自己的实验前,先根据目标速率决定批大小策略——低速追求亚秒延迟就调小 --autocommit-duration-ms,高速追求吞吐就适当放大。
把基准跑出可信结果的注意事项
综合 README 说明、benchmark.py 与配套文章 实验方法描述,有几条经验直接影响测量可信度:
- 保留预热期:消费者需先追上生产者抢先写入的数据。
consumer.py内建的WARMUP_PERIOD = 30秒就是为了滤掉这段追赶期;自行修改脚本或缩短档位时长时,不要把这个机制去掉; - 延迟包含对象存储自身开销:测到的是"Pathway + Delta Lake + S3/MinIO"整条链路的延迟,存储端 IO 抖动会直接反映在百分位上,属于预期而非故障;
- 控制网络环境变量:配套实验将 MinIO 与测试机置于同一网络以最小化外网延迟。跨公网复现时,应预期更高的延迟基数;
- 均匀地扩展档位:
--range-step支持你从 10k 一直扫到 250k 条/秒来观察延迟随速率的变化形态;但要意识到单线程、单消费者的设置下,过高速率会先触碰到存储轮询与网络同步的瓶颈; - CPU/内存/网络要充分供给:基准本身也吃资源,资源不足时百分位尾部会被人为放大。
如果你希望在一个更接近生产形态(多时间源、真实 ETL 转换)的场景里观察同样的读写机制,可以进一步查看同目录下的 minio-ETL 示例,它与基准共用同一套 MinIOSettings 连接配置与 Delta Lake 读写模式,可作为从"基准"走向"应用"的桥梁。
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 StartedRust0627
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