首页
/ Pathway 实现无 Kafka 的亚秒级实时流处理:基于 S3/MinIO 与 Delta Tables 构建 Python 数据管线

Pathway 实现无 Kafka 的亚秒级实时流处理:基于 S3/MinIO 与 Delta Tables 构建 Python 数据管线

2026-09-07 14:13:11作者:邬祺芯Juliet

本文围绕官方 ETL 模板文章的核心方案展开:以 Pathway(Rust 内核的 Python 流处理框架)+ Delta Lake + S3 兼容对象存储(如 MinIO) 取代 Apache Kafka,作为消息队列与流处理层的轻量替代。你将掌握该架构的三大组件如何协同、如何在本地复现 minio-ETL 示例管线(从 MinIO 写入、读取、时区统一转换到回写 Delta Tables),以及如何通过 autocommit_duration_msmin_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 架构的三大核心组件

该架构把基础设施压缩为三个可清晰分工的组件:

  1. S3 兼容对象存储(如 MinIO):消息以对象形式存放在可扩展、低成本的对象存储上。S3(Simple Storage Service)将数据以"bucket 内带唯一 key 的对象"方式管理,本质是一个 key-value 存储:按 key 存取、无目录层级/顺序概念——这与传统文件系统的差异正是对象存储能支撑海量流式数据的原因。
  2. Delta Tables:基于 Delta Lake 实现的 ACID 存储层,让对象存储"表现得像数据库表":记录数据变更、保证一致性、支持 schema 演进与高性能分析,并提供追加式(append-only)写入与读取机制。它可以独立于 Kafka 充当"数据流落地点"。
  3. 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(计算批次)与 diff1 表示新增、-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。本节先复现前者——一个跨时区消息的时间戳统一管线。

项目结构与文件职责

官方模板将该管线描述为如下结构(.envREADME.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.pyetl.pyproducer.pyread-results.py),凭据通过 python-dotenvload_dotenv() 从环境变量注入,下文会逐一说明。

第 1 步:准备 MinIO 实例与 Bucket

  • 启动一个 MinIO 实例(自托管或托管均可),准备好凭据 MINIO_S3_ACCESS_KEYMINIO_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.pystorage_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 Tabletimezone_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 步:运行管线四部曲

  1. 运行 Producer:启动 producer.py 生成数据流并写入两个 Delta Table。
  2. 启动 ETL 管线:另开一个终端运行 etl.py,Pathway 将持续消费 Delta Table 上的新数据并完成转换与回写。
  3. 停止管线producer.py 结束后等待约 10 秒,按 Ctrl+C 停止 Pathway。
  4. 读取并保存结果:运行 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_msmin_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_msmin_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

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