Pathway 实现无 Kafka 的亚秒级实时流处理:基于 S3/MinIO 与 Delta Tables 构建 Python 数据管线
本文围绕官方 ETL 模板文章的核心方案展开:以 Pathway(Rust 内核的 Python 流处理框架)+ Delta Lake + S3 兼容对象存储(如 MinIO) 取代 Apache Kafka,作为消息队列与流处理层的轻量替代。你将掌握该架构的三大组件如何协同、如何在本地复现 minio-ETL 示例管线(从 MinIO 写入、读取、时区统一转换到回写 Delta Tables),以及如何通过 autocommit_duration_ms 与 min_commit_frequency 两个关键参数在"延迟"与"吞吐"之间做取舍。仓库内的可运行代码位于 examples/projects/kafka-alternatives,Delta Lake 连接器的 Python 实现位于 python/pathway/io/deltalake/init.py。
为什么考虑用 S3 替代 Kafka
Apache Kafka 是构建实时流式数据管线的事实标准分布式事件存储,但它给"轻量使用场景"带来了不成比例的成本与运维负担。
成本
- 自建集群:部署与维护自管理 Kafka 集群每月成本可超过 2 万美元(据 Confluent 估算),涵盖硬件、维护与人力开销。
- 托管服务:即便使用 Confluent Cloud 之类的托管服务,仅处理每秒 5 万条小消息每月就需约 1800 美元,且不含风险管理、运营等额外支出。
复杂性
- Kafka 需要配置与维护 brokers、ZooKeeper、producers、consumers 等多个组件。
- 保障可用性、容错与可扩展性需要持续投入专业运维。
- 团队需要付出学习成本去掌握 Kafka 的内部机制。
如果你已经在使用 S3 兼容对象存储,完全可以换一种思路:把对象存储当作"消息队列底座",让一个流处理引擎负责提供 Delta Tables 本身所欠缺的"持续监听、增量读取与即时处理"能力,从而免去整个 Kafka 集群。
无 Kafka 架构的三大核心组件
该架构把基础设施压缩为三个可清晰分工的组件:
- S3 兼容对象存储(如 MinIO):消息以对象形式存放在可扩展、低成本的对象存储上。S3(Simple Storage Service)将数据以"bucket 内带唯一 key 的对象"方式管理,本质是一个 key-value 存储:按 key 存取、无目录层级/顺序概念——这与传统文件系统的差异正是对象存储能支撑海量流式数据的原因。
- Delta Tables:基于 Delta Lake 实现的 ACID 存储层,让对象存储"表现得像数据库表":记录数据变更、保证一致性、支持 schema 演进与高性能分析,并提供追加式(append-only)写入与读取机制。它可以独立于 Kafka 充当"数据流落地点"。
- Pathway Live Data Framework:充当处理引擎,持续监控、读取 Delta Tables 上的新增数据并完成转换、回写,是整套架构中真正提供"流式能力"的一环。
一句话概括分工:Delta Tables 负责"数据如何被组织",Pathway 负责"数据如何流动"。
为什么这个组合能带来流式语义
单看 Delta Tables,它只是一个存储层,不会主动推送数据。Pathway 通过 Delta Lake 连接器弥补这一点。从 read 的源码文档 可以看到:
- 连接器以
"streaming"(默认)模式运行,引擎会持续轮询 Delta Lake 中新增的版本,把每次数据变更反映到内部状态(对应mode参数;也可传"static"一次性读取现有数据)。 - 读取相对于数据版本是原子的:若一个 Delta 事务内发生了插入、删除或更新,这些变更会作为一个整体、在同一个 minibatch 内一次性进入计算图。
- 写入侧(write 源码)默认以
"stream_of_changes"模式输出变更流,并附带time(计算批次)与diff(1表示新增、-1表示删除)两个整型列——这正是你在结果 CSV 中会看到的列名由来。
从源码结构看,这套"读增量版本 + 追加式写入 + minibatch 提交"的机制,使得对象存储上也能获得接近消息队列的连续处理体验,且读写均基于同一套 S3 凭证体系与连接器配置。
这套 Kafka 替代方案适用的场景
官方基准显示(见下文),Pathway + Delta Tables + MinIO 在吞吐上"一个数量级低于 Kafka 的理论上限",但对大量"近实时即可"的业务绰绰有余:
- 物流与 IoT:从物流、工业自动化、农业设备的海量传感器持续采集数据。
- 金融服务:交易处理、欺诈检测、行情数据流的近实时/实时处理。
- Web 与移动端分析:用户行为监控、实时竞价、广告曝光追踪。
关键判断标准是:当你已拥有 S3 存储、对延迟要求为亚秒级而非毫秒级、并且希望显著减少组件与成本时,该方案具备明显吸引力。文章同时提醒:Pathway 本身也常被用于超低延迟场景(F1 赛车实时分析、Intel/NATO 等关键任务),但"无 Kafka"架构强调的是以简化基础设施为前提满足足够好的实时性。
实战:在 MinIO 上构建完整流式数据管线
仓库中给出了两套可运行代码:管线项目 examples/projects/kafka-alternatives/minio-ETL 与基准测试项目 examples/projects/kafka-alternatives/benchmarks。本节先复现前者——一个跨时区消息的时间戳统一管线。
项目结构与文件职责
官方模板将该管线描述为如下结构(.env 与 README.md 为文档约定文件):
.
├── base.py # MinIO 访问配置(bucket、endpoint、凭据)
├── etl.py # 核心 ETL:读取 Delta Tables → 转换 → 回写
├── producer.py # 以纽约/巴黎两个时区生成带时间戳的消息并写入 Delta Tables
├── read-results.py # 读取 ETL 输出流,保存为 CSV
└── .env # 存放 MinIO 凭据(python-dotenv 加载)
当前仓库中实际包含 4 个 Python 模块(base.py、etl.py、producer.py、read-results.py),凭据通过 python-dotenv 的 load_dotenv() 从环境变量注入,下文会逐一说明。
第 1 步:准备 MinIO 实例与 Bucket
- 启动一个 MinIO 实例(自托管或托管均可),准备好凭据
MINIO_S3_ACCESS_KEY与MINIO_S3_SECRET_ACCESS_KEY。 - 预先创建目标 bucket:Delta Lake 写入者不会自动创建不存在的 bucket,向不存在的 bucket 写入会导致管线失败。
在 .env 中填入凭据:
MINIO_S3_ACCESS_KEY = *******
MINIO_S3_SECRET_ACCESS_KEY = *******
第 2 步:在 base.py 配置访问参数
base.py 负责统一定义所有模块共享的连接配置。注意其 URL 均以 s3:// 开头——连接器源码 明确说明:路径以 s3:///s3a:// 开头走 S3 存储,其余路径走本地文件系统:
import os
from dotenv import load_dotenv
import pathway as pw
load_dotenv()
bucket = "your-bucket"
base_path = f"s3://{bucket}/"
str_repr = "%Y-%m-%d %H:%M:%S.%f %z"
s3_connection_settings = pw.io.minio.MinIOSettings(
bucket_name=bucket,
access_key=os.environ["MINIO_S3_ACCESS_KEY"],
secret_access_key=os.environ["MINIO_S3_SECRET_ACCESS_KEY"],
endpoint="your-url-endpoint.com",
)
要点说明:
- 需把
bucket换成你的真实 bucket 名,把endpoint换成 MinIO 的实际访问端点。 MinIOSettings是 S3 兼容对象存储的连接配置类型;若使用标准 AWS S3,可改用pw.io.s3.AwsS3Settings(同样支持bucket_name/region/access_key/secret_access_key)。如使用真实 AWS 且已登录认证环境,连接器文档说明也可省略凭据,引擎会自动使用当前已授权用户身份(见 write docstring)。str_repr是消息中日期字符串的 strftime 格式,必须带%z时区信息,后续 ETL 会按此解析。
第 3 步:producer.py —— 模拟多时区消息生产者
producer.py 模拟"跨时区数据源":每秒随机生成一条消息,归属纽约时区(America/New_York)或巴黎时区(Europe/Paris),写入对应的 Delta Table(timezone1 / timezone2,充当类似 Kafka topic 的角色)。
写入侧用的是 deltalake 库的 write_deltalake,以 mode="append" 向 Delta Table 追加数据。它的 storage_options 需要补齐 S3 兼容存储的连接信息,包括 AWS_REGION、bucket、endpoint,以及 MinIO 特有的 AWS_S3_ALLOW_UNSAFE_RENAME(MinIO 不支持原子 rename)与 AWS_STORAGE_ALLOW_HTTP(本地 HTTP 端点需要):
def send_message(timezone: ZoneInfo, topic: str, i: int):
timestamp = datetime.datetime.now(timezone)
lake_path = base_path + topic
message_json = {"date": timestamp.strftime(str_repr), "message": str(i)}
df = pd.DataFrame([message_json])
write_deltalake(
lake_path,
df,
mode="append",
storage_options={
"AWS_ACCESS_KEY_ID": os.environ["MINIO_S3_ACCESS_KEY"],
"AWS_SECRET_ACCESS_KEY": os.environ["MINIO_S3_SECRET_ACCESS_KEY"],
"AWS_REGION": "eu-central-1",
"AWS_S3_ALLOW_UNSAFE_RENAME": "True",
"AWS_BUCKET_NAME": "your-bucket",
"endpoint": "your-url-endpoint.com",
"AWS_STORAGE_ALLOW_HTTP": "true",
},
)
注意:文档中提示需在
producer.py的storage_options中同步配置AWS_REGION等参数,示例代码已包含;实际运行时请替换 bucket 与 endpoint 为你的值。
第 4 步:etl.py —— 读取、转换与回写 Delta Tables
etl.py 是本管线的核心,完整展示了"读 Delta Tables → 流式转换 → 写 Delta Tables"三段式。
读取两个时区的 Delta Table:
timestamps_timezone_1 = pw.io.deltalake.read(
base_path + "timezone1",
schema=InputStreamSchema,
s3_connection_settings=s3_connection_settings,
autocommit_duration_ms=100,
)
timestamps_timezone_2 = pw.io.deltalake.read(
base_path + "timezone2",
schema=InputStreamSchema,
s3_connection_settings=s3_connection_settings,
autocommit_duration_ms=100,
)
其中 InputStreamSchema 定义读入字段:
class InputStreamSchema(pw.Schema):
date: str
message: str
统一不同时区的时间戳。官方文章给出一段示意代码 pw.transform.unify_timezones(...),而在当前仓库的实际实现中,时区统一是通过显式的解析与拼接完成的两步转换:先用 dt.strptime 按带时区格式解析字符串为 datetime,再用 dt.timestamp(unit="ms") 转成毫秒级 Unix 时间戳,最后用 concat_reindex 将两条流合并为一条统一的表(concat_reindex 定义见 table.py):
def convert_to_timestamp(table: pw.Table) -> pw.Table:
table = table.select(
date=pw.this.date.dt.strptime(fmt=str_repr, contains_timezone=True),
message=pw.this.message,
)
return table.select(
timestamp=pw.this.date.dt.timestamp(unit="ms"),
message=pw.this.message,
)
timestamps_timezone_1 = convert_to_timestamp(timestamps_timezone_1)
timestamps_timezone_2 = convert_to_timestamp(timestamps_timezone_2)
timestamps_unified = timestamps_timezone_1.concat_reindex(timestamps_timezone_2)
将统一后的流回写到 Delta Table(timezone_unified):
pw.io.deltalake.write(
timestamps_unified,
base_path + "timezone_unified",
s3_connection_settings=s3_connection_settings,
min_commit_frequency=100,
)
# 启动计算
pw.run(monitoring_level=pw.MonitoringLevel.NONE)
这里出现了两个贯穿全篇的关键参数,它们的语义可直接对照源码:
| 参数 | 所在方法 | 仓库默认值 | 本示例取值 | 作用 |
|---|---|---|---|---|
autocommit_duration_ms |
pw.io.deltalake.read |
1500 |
100 |
连接器每隔该毫秒数将收到的更新提交并推入 Pathway 计算图(read 源码) |
min_commit_frequency |
pw.io.deltalake.write |
60_000 |
100 |
两次向存储提交数据之间的最小时间间隔(毫秒);每次提交都会新建文件并写一条事务日志,故官方建议限制提交频率以降低开销(write 源码) |
第 5 步:运行管线四部曲
- 运行 Producer:启动
producer.py生成数据流并写入两个 Delta Table。 - 启动 ETL 管线:另开一个终端运行
etl.py,Pathway 将持续消费 Delta Table 上的新数据并完成转换与回写。 - 停止管线:
producer.py结束后等待约 10 秒,按Ctrl+C停止 Pathway。 - 读取并保存结果:运行
read-results.py(源码 通过pw.io.csv.write(table, "./results.csv")将timezone_unified的流写为 CSV),然后用cat results.csv查看。
输出类似:
timestamp,message,time,diff
1724403133066.217,"0",1724403268648,1
1724403138388.917,"1",1724403269748,1
1724403144896.706,"3",1724403270848,1
1724403147095.393,"4",1724403271948,1
1724403149295.165,"5",1724403272948,1
1724403151499.115,"6",1724403274048,1
1724403153736.456,"7",1724403275148,1
1724403140576.744,"2",1724403276248,1
1724403158229.244,"9",1724403276248,1
可见 message 均已带上统一换算为毫秒 Unix 时间戳的 timestamp 列,diff=1 表示行新增——这正是 stream_of_changes 输出语义的直观体现。
至此,一个"以 Delta Tables 为消息载体、Pathway 驱动流动"的轻量流处理管线就构建完成了。
基准测试与延迟/吞吐调优
基准方法说明
官方基准(benchmark.py 源码)将消息写入 MinIO 上的 Delta Table 再读回,延迟定义为"消息发出 → Pathway 成功处理"的时间间隔:
- 测试速率范围 10k~250k 条/秒,每个速率跑 5 分钟(默认
--seconds-to-stream 300)。 - 前 30 秒不计入结果:因为 producer 先于 consumer 启动,头部数据需要追赶,会造成延迟虚高;剔除后更能反映稳态表现。
- MinIO 与测试机同网络,排除外部网络延迟;Pathway 以单线程模式运行以保证一致性。
- 复现命令示例:测试 10k/20k/30k 三个速率、各跑 10 分钟:
python benchmark.py --range-start 10000 --range-end 30000 --range-step 10000 --seconds-to-stream 600
结果 CSV 写入 benchmark-results/ 目录(详见 benchmarks/README.md)。
默认参数下的表现
在"批大小较大、提交间隔较长"的默认配置下,高吞吐下延迟仍可控,适合近实时应用:
- 70,000 条/秒:中位延迟约 1.45 秒
- 150,000 条/秒:中位延迟约 1.68 秒
- 250,000 条/秒:中位延迟约 1.98 秒
调整参数实现亚秒级延迟
默认配置下,消息在提交到 Delta Table 前最多可能被攒批约 1 秒(示例/基准脚本中 autocommit_duration_ms 默认取 1000,见 benchmark.py;连接器源码层默认值则为 1500 毫秒)。攒批有利于吞吐,但会抬高延迟。
因此官方针对中低速率负载给出的优化建议是:把 autocommit_duration_ms 与 min_commit_frequency 降到 100 毫秒左右,缩小批次、缩短消息在队列中的等待时间。文档给出的调优后实测数据如下:
| Workload (messages/sec) | p50 (s) | p75 (s) | p85 (s) | p95 (s) | p99 (s) |
|---|---|---|---|---|---|
| 10,000 | 0.26 | 0.35 | 0.36 | 0.53 | 0.67 |
| 30,000 | 0.28 | 0.31 | 0.32 | 0.38 | 0.64 |
| 70,000 | 1.24 | 16.06 | 24.82 | 35.66 | 41.32 |
结论与代价
- 在 1 万~3 万条/秒区间,调参后延迟显著改善且抖动很小(约 50 ms 波动);中位处理时间提升约 4 倍,p99 提升约 3 倍;亚秒级表现可维持到约 6 万条/秒。
- 但超过约 6 万条/秒后系统开始劣化:虽然延迟有所改善,吞吐却明显受损。原因在于更小的批次意味着系统要向存储发起约 10 倍的网络调用与同步操作(如轮询最新 Delta Table 版本),同时网络错误的影响被放大——一旦对 Delta Table 的访问短暂中断,恢复难度会更高。
- 从源码看,
min_commit_frequency的文档也强调了类似权衡:每次提交都会创建新文件并写事务日志,过度频繁的提交会增加下游处理该表的开销(write docstring)。
因此调参的正确姿势是:先明确自身负载的量级,再结合 autocommit_duration_ms 与 min_commit_frequency 做实验,找到延迟与吞吐的平衡点。
给复现者的建议
- 自建基准时务必保留 warm-up 阶段,等系统进入稳态再计延迟。
- 文档强调:观察到的延迟已包含 S3 兼容存储自身的读写开销,换用不同的对象存储会得到不同数值。
- 确保测试环境资源充足(CPU、内存、网络都会直接影响结果),并注意上述数字来自特定环境,你的实际数值应以本地复现为准。
在三大云平台落地同一套架构
由于连接器统一支持 S3 协议,上述无 Kafka 方案可直接迁移到公有云:
- AWS:Pathway 与 AWS S3 原生集成,通过 S3 + Delta Lake 连接器即可读写存储于 S3 的 Delta Tables,省去 Kafka 的复杂度。仓库提供了 examples/projects/aws-fargare-deploy 方向的部署示例(对应目录为 examples/projects/aws-fargate-deploy)。
- Google Cloud:可部署于 Google Kubernetes Engine(GKE),连接 Google Cloud Storage(同样是 S3 兼容访问)处理流式数据,仓库提供 GCP 部署参考 examples/projects/aws-fargate-deploy 同级部署示例目录。具体云上配置请以官方部署文档为准。
- Azure:可基于 Azure Container Instances 部署,对接 Azure 存储服务,仓库中的参考示例位于 examples/projects/azure-aci-deploy。
无论落到哪朵云,"对象存储作为消息载体 + Pathway 提供流式处理能力"的核心模式保持一致,只是把 MinIOSettings 换成对应云厂商 S3 兼容配置并填写相应 region/endpoint 即可。
TL;DR:把 S3 变成一条完整的流式数据管线
- Kafka 功能强大但复杂昂贵;若你已拥有 S3 兼容存储且只需亚秒级延迟,完全可以用 Pathway + Delta Tables on S3 替代 Kafka。
- 该架构将基础设施收敛为对象存储、Delta Tables 与处理引擎三部分,运维负担与成本显著下降,且能用纯 Python 构建管线、学习曲线平缓。
- 官方基准表明:默认参数下可支撑最高约 25 万条/秒的近实时吞吐;将提交/批处理参数下调至 100 ms 后,约 6 万条/秒以内可稳定获得亚秒级延迟(10k/30k 负载下 p99 约 0.64~0.67 秒)。
- 调优核心是权衡:批次越小延迟越低,但网络调用与版本轮询同步操作越多,吞吐上限随之下降。
- 已验证适用场景包括 IoT/物流传感器采集、金融交易与欺诈检测、Web/移动分析与实时竞价等。
想亲手验证,可从仓库的 minio-ETL 示例 开始;想要在不同速率下复现官方数据,可直接运行 benchmarks 下的脚本。Delta Lake 连接器的完整参数说明与实现细节,可进一步阅读 python/pathway/io/deltalake/init.py。
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 StartedRust0625
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