首页
/ 将 S3/MinIO + Delta Lake 当作消息队列:Pathway 的 Kafka 替代方案延迟基准测试完整指南

将 S3/MinIO + Delta Lake 当作消息队列:Pathway 的 Kafka 替代方案延迟基准测试完整指南

2026-09-07 19:42:44作者:吴年前Myrtle

本文是一份以 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_KEYMINIO_S3_SECRET_ACCESS_KEY
  • AWS S3 后端读取 AWS_S3_ACCESS_KEYAWS_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 的构造过程分三步:

  1. 定义一个 data: bytes_MessageTableSchema(消息负载统一为字节串,schema 极简,使测量聚焦于传输而非反序列化);
  2. 通过 pw.io.python.read(...) 从自定义的 ConnectorSubject 子类 _MessageTableSubject 读取外部消息,autocommit_duration_ms 决定外部输入被提交给引擎的间隔;
  3. 通过 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 是队列读端的封装,构造时做三件事:

  1. 解析 s3://bucket/path 形式的 base_path,拆出 bucket 名与桶内路径,用于后续轮询探测(若连接设置里已带 bucket 则复用);
  2. pw.io.deltalake.read(uri=..., mode="streaming", s3_connection_settings=..., autocommit_duration_ms=...)流式模式持续读 Delta Table 的新增版本。相比快照模式,mode="streaming" 意味着引擎会持续轮询表的新版本并把增量更新推向下游;
  3. 通过 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-backendminio/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 后端类型,可选 s3minio
--autocommit-duration-ms int 1000 read/write 两侧的提交间隔(毫秒),同时传给 producer 与 consumer

每个速率档位的执行流程(见 benchmark.py 主循环)依次是:

  1. 构造临时 S3 前缀 adhoc/s3-messaging/{rate},完整路径形如 s3://{bucket}/adhoc/s3-messaging/{rate}
  2. 调用 cleanup_files_under_prefix 清理该前缀下上一轮遗留对象,保证测试从空表开始;
  3. 先启动 consumer 子进程(expected_messages = rate × seconds_to_streamlog_frequency = rate // 2),再启动 producer 子进程;
  4. assert 等待两个子进程均以退出码 0 结束;
  5. 再次清理 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 中默认值为 1500 ms。调小它,消息在内存中等待的时间变短,延迟随之下降;
  • min_commit_frequency:Delta Lake 写侧两次数据提交之间的最小间隔。在 deltalake.write 中默认值为 60_000 ms——注意这个默认值很大,若不显式调小,写入方会长时间攒批。基准中 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 与配套文章 实验方法描述,有几条经验直接影响测量可信度:

  1. 保留预热期:消费者需先追上生产者抢先写入的数据。consumer.py 内建的 WARMUP_PERIOD = 30 秒就是为了滤掉这段追赶期;自行修改脚本或缩短档位时长时,不要把这个机制去掉;
  2. 延迟包含对象存储自身开销:测到的是"Pathway + Delta Lake + S3/MinIO"整条链路的延迟,存储端 IO 抖动会直接反映在百分位上,属于预期而非故障;
  3. 控制网络环境变量:配套实验将 MinIO 与测试机置于同一网络以最小化外网延迟。跨公网复现时,应预期更高的延迟基数;
  4. 均匀地扩展档位--range-step 支持你从 10k 一直扫到 250k 条/秒来观察延迟随速率的变化形态;但要意识到单线程、单消费者的设置下,过高速率会先触碰到存储轮询与网络同步的瓶颈;
  5. CPU/内存/网络要充分供给:基准本身也吃资源,资源不足时百分位尾部会被人为放大。

如果你希望在一个更接近生产形态(多时间源、真实 ETL 转换)的场景里观察同样的读写机制,可以进一步查看同目录下的 minio-ETL 示例,它与基准共用同一套 MinIOSettings 连接配置与 Delta Lake 读写模式,可作为从"基准"走向"应用"的桥梁。

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