Pathway 实时流式线性回归实战:从 Kafka 数据流接入到增量参数估计
【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway
导读
本文基于 Pathway(Live Data Framework)官方模板文档,完整讲解如何对一个持续到达的 Kafka 数据流做流式简单线性回归:每收到一个新数据点 $(x_i, y_i)$,系统就基于"至今为止的全部数据点"重新估计回归参数 $(a,b)$,使 $y_i \approx a + b \times x_i$,并把最新估计结果实时输出。读完本文,你将掌握 Pathway 的 Kafka/CSV 连接器配置、select/reduce/pw.apply增量计算写法、pw.run()引擎启动方式,以及"变更流(tables of changes)"输出文件的解读方法,并可直接用仓库附带的完整示例工程跑通端到端流程。
本文是 Pathway 首个实时应用教程(实时求和)的进阶延伸:从"实时求和"升级到"实时做机器学习"。
整体思路:把最小二乘"改造成"增量统计
对一组已知点做简单线性回归,等价于求使残差平方和最小的截距 $a$ 与斜率 $b$。其闭式解只需五个统计量:数据点个数 $n$、$\sum x$、$\sum y$、$\sum x^2$、$\sum xy$。基于它们可以写出:
$$ d = n\sum x^2 - (\sum x)^2,\quad a = \frac{\sum y \cdot \sum x^2 - \sum x \cdot \sum xy}{d},\quad b = \frac{n \cdot \sum xy - \sum x \cdot \sum y}{d} $$
流式化的关键在于:这五个量全部可以用累加(aggregation)表达。每当 Kafka 里到达一个新数据点,Pathway 引擎就会增量地更新这些累加值,再顺势重算一次 $a$、$b$。于是"批式最小二乘"就自然变成了"实时回归",无须手动维护滑动窗口或增量状态——这正是 Pathwayreduce聚合语义带来的好处(下文的计算实现会展开)。
整个工程涉及两条数据通路:
- Producer 侧:用
kafka-python的KafkaProducer往 Kafka topic 写入带噪声的数据点; - Pathway 侧:
realtime_regression.py通过 Kafka 连接器消费这些消息,做增量回归并把结果写入 CSV 输出。
对应代码、样例输入与样例输出都放在仓库的示例工程目录下:
- realtime_regression.py:Pathway 流式回归主程序;
- generating_kafka_stream.py:Kafka 数据流生成器;
- example_regression_input.csv 与 example_regression_output_stream.csv:跑通后的真实样例输出。
环境与前置条件
本文假设你已准备好:
- Python 环境:装有 Pathway(本仓库的 Python 包源码位于 python/),以及生成流用的
kafka-python(代码里以from kafka import KafkaProducer导入)。 - 一个可用的 Kafka:既可以使用托管 Kafka 服务(如 Confluent Cloud、Upstash),也可以使用兼容 Kafka 的 Redpanda,或直接用 Docker / Docker Compose 在本地跑一个 Kafka 镜像做实验。仓库示例
realtime_regression.py与generating_kafka_stream.py中通过环境变量UPSTASH_KAFKA_USER、UPSTASH_KAFKA_PASS注入凭证,属于使用托管服务的写法,可参考。 - 一个 topic:本文统一使用
"linear-regression"。
第一步:用 Kafka 连接器读取输入流
在 Pathway 中,读写外部数据靠的是连接器(connectors)。这里我们只关心两类:Kafka 输入连接器,以及把结果写盘用的 CSV 输出连接器。
1.1 rdkafka 连接参数
Pathway 的 Kafka 连接器底层对接 librdkafka,因此所有 Kafka 连接参数都放进一个 Python 字典里,键名遵循 librdkafka 的配置规范。以下是一个使用SASL-SSL + SCRAM-SHA-256认证的典型配置(务必把服务地址、用户名、密码替换成你自己的):
rdkafka_settings = { "bootstrap.servers": "server-address:9092", "security.protocol": "sasl_ssl", "sasl.mechanism": "SCRAM-SHA-256", "group.id": "$GROUP_NAME", "session.timeout.ms": "6000", "sasl.username": "username", "sasl.password": "********", }各参数作用如下:
| 参数 | 含义 | 取值建议 |
|---|---|---|
bootstrap.servers | Kafka broker 地址,host:port形式 | 如server-address:9092或localhost:9092 |
security.protocol | 传输加密方式 | 托管服务通常为sasl_ssl,本地开发可为plaintext |
sasl.mechanism | SASL 认证机制 | 常见SCRAM-SHA-256、SCRAM-SHA-512、PLAIN |
group.id | 消费组 ID,决定消息如何被分组消费 | 用你的项目/消费者组命名 |
session.timeout.ms | 消费组会话超时时间 | 示例给6000 |
sasl.username/sasl.password | 认证账号密码 | 建议通过环境变量注入,勿硬编码 |
1.2 定义 schema 并读取 topic
Kafka 消息只是字节流,要让 Pathway 知道每列是什么类型,需要先声明一个pw.Schema:
class InputSchema(pw.Schema): x: float y: float随后一行代码即可建立输入表:
t = pw.io.kafka.read( rdkafka_settings, topic="linear-regression", schema=InputSchema, format="csv", autocommit_duration_ms=1000 )关键参数:
rdkafka_settings:上一步的 librdkafka 风格配置字典;topic:要订阅的 Kafka topic;schema:输入表的列结构,由pw.Schema子类声明;format:消息的序列化格式,决定字节如何被解析成列;autocommit_duration_ms:两次提交之间的最大时间间隔,引擎会按此节奏把收到的新消息批量推进到计算图中(示例为1000毫秒,即大约每秒刷新一次结果)。在当前仓库的 Kafka 连接器实现中,该参数的默认值为1500ms,并额外提供mode("streaming"/"static")、parallel_readers、with_metadata、start_from_timestamp_ms等选项,可在需要控制消费模式、并行度或记录元信息时查阅。
关于format的一点版本说明:教程撰写时的csv格式,在当前仓库的pw.io.kafka.read签名中已演化为显式支持的"plaintext"、"raw"、"json"三种:
raw:把消息原样读入,产生一个含data列的表(当前实现中 raw/plaintext 还会附带key列,见该文件 docstring);plaintext:把消息按 UTF-8 文本解析,同样落在data列;json:先把 JSON 负载解析出来,再按schema中声明的列建表,非常适合"一条消息 = 一个点"的场景。
仓库配套示例 realtime_regression.py 目前使用的正是format="json"。本文为了贴合教程原文,主体演示仍采用csv的写法;若你安装的 Pathway 版本较新,可直接切到json,两种写法在下面的代码中都能套用(差异点我们会在生成数据流中专门对比)。
💡只想快速验证算法、不想搭 Kafka?可以跳过连接器,直接用 Pathway 自带的流生成器:
t = pw.demo.noisy_linear_stream()该函数定义在 python/pathway/demo/init.py:签名是
noisy_linear_stream(nb_rows: int = 10, input_rate: float = 1.0),内部固定random.seed(0),生成两列数据——x为从 0 递增的整数(并被标记为主键),y = x + (2·r − 1)/10,即"理想直线 y=x + 幅值 ±0.1 的均匀噪声",随后经pw.io.python.read按 JSON 格式逐行喂入引擎(python/pathway/demo/init.py 的generate_custom_stream)。它的测试覆盖见 python/pathway/tests/test_demo.py。更全面的pw.demoAPI 说明可参考人工数据流文档。
1.3 CSV 输出连接器
结果要落盘观察,CSV 连接器一行即可:
pw.io.csv.write(t, "regression_output_stream.csv")CSV 输出连接器会把 Pathway 表的每次更新(而不是最终快照)追加写进文件,因此该连接器也适合把"中间输入"旁路存档。本文我们会同时用它导出原始输入和回归结果两个文件。更完整的 CSV 连接器讲解见实时应用教程与连接器总览文档。
第二步:做流式线性回归计算
拿到流式输入表t(列x、y)之后,计算分三步走。
2.1 用select扩展出 $x^2$ 与 $xy$ 两列
t = t.select( *pw.this, x_square=t.x * t.x, x_y=t.x * t.y )*pw.this表示保留原表所有列,再按表达式t.x * t.x、t.x * t.y新增x_square、x_y两列,得到每行含有x, y, x_square, x_y的中间表。
2.2 用reduce累加五个统计量
statistics_table = t.reduce( count=pw.reducers.count(), sum_x=pw.reducers.sum(t.x), sum_y=pw.reducers.sum(t.y), sum_x_y=pw.reducers.sum(t.x_y), sum_x_square=pw.reducers.sum(t.x_square), )reduce把整张表折叠成一行:count是累计到达的数据点个数,sum_x、sum_y、sum_x_y、sum_x_square分别是对应列的累计和。Pathway 的引擎会对这行做增量维护——每来一个新点,五个聚合值原地更新,statistics_table永远是"到当前时刻为止"的统计。
2.3 用pw.apply逐点计算回归参数
根据前面给出的闭式解公式,把分母记为 $d = n\sum x^2 - (\sum x)^2$,注意退化情形:当所有点的 $x$ 完全相同(例如只收到 1 个点)时 $d = 0$,公式不可用,代码里直接返回0兜底:
def compute_a(sum_x, sum_y, sum_x_square, sum_x_y, count): d = count * sum_x_square - sum_x * sum_x if d == 0: return 0 else: return (sum_y * sum_x_square - sum_x * sum_x_y) / d def compute_b(sum_x, sum_y, sum_x_square, sum_x_y, count): d = count * sum_x_square - sum_x * sum_x if d == 0: return 0 else: return (count * sum_x_y - sum_x * sum_y) / d results_table = statistics_table.select( a=pw.apply(compute_a, **statistics_table), b=pw.apply(compute_b, **statistics_table), )**statistics_table把单行统计表中的列按名字解包成关键字参数,喂给pw.apply(compute_a, ...);pw.apply让纯 Python 函数(而非 Pathway 表达式)得以对行内各列执行标量计算。最终results_table只含两列:估计出的截距a与斜率b。
第三步:生成输入数据流(Producer 侧)
本节面向用 Kafka 连接器的读者;如果使用
pw.demo.noisy_linear_stream()生成器可直接跳过。
用 CSV 格式消费 Kafka 消息时,需要遵守两条规则:
- 第一条消息必须是列头,例如
"x,y",否则连接器无法完成列映射; - 列头只能发送一次——若重复发送,第二条列头会被当作普通数据行解析。
下面用kafka-python的KafkaProducer演示:先发列头,再发两个点 $(0,0)$、$(1,1)$,最后关闭 Producer:
producer = KafkaProducer( bootstrap_servers=["server-address:9092"], sasl_mechanism="SCRAM-SHA-256", security_protocol="SASL_SSL", sasl_plain_username="username", sasl_plain_password="********", ) producer.send(topic, ("x,y").encode("utf-8"), partition=0) producer.send( "linear-regression", ("0,0").encode("utf-8"), partition=0 ) producer.send( "linear-regression", ("1,1").encode("utf-8"), partition=0 ) producer.close()为了使回归不那么"平凡"(真实数据几乎都带噪声),本例让数据点围绕直线 $y=x$ 采样,并对每个 $y$ 加上小幅随机扰动。
提示:视你的 Kafka 版本,Producer 可能还需要显式指定协议版本
api_version=(0,10,2)才能正常工作。
如果使用json格式(即当前仓库配套示例的写法),Producer 端不需要"列头消息",只需把每个数据点序列化成 JSON 字典即可,例如{"x": i, "y": get_value(i)}(见 generating_kafka_stream.py)。这正是"csv 要首条发列头、json 不用"的核心差别。
第四步:组装完整工程并运行
最终工程由两个文件组成,均可在仓库的 examples/projects/kafka-linear-regression 目录下找到。
4.1realtime_regression.py:Pathway 处理主程序
import pathway as pw rdkafka_settings = { "bootstrap.servers": "server-address:9092", "security.protocol": "sasl_ssl", "sasl.mechanism": "SCRAM-SHA-256", "group.id": "$GROUP_NAME", "session.timeout.ms": "6000", "sasl.username": "username", "sasl.password": "********", } class InputSchema(pw.Schema): x: float y: float # 1. 读取 Kafka 流 t = pw.io.kafka.read( rdkafka_settings, topic="linear-regression", schema=InputSchema, format="csv", autocommit_duration_ms=1000, ) # 2. 把收到的原始输入也旁路写一份,便于对照 pw.io.csv.write(t, "regression_input.csv") # 3. 扩展列:x_square, x_y t = t.select( *pw.this, x_square=t.x * t.x, x_y=t.x * t.y, ) # 4. 全局累加五个统计量 statistics_table = t.reduce( count=pw.reducers.count(), sum_x=pw.reducers.sum(t.x), sum_y=pw.reducers.sum(t.y), sum_x_y=pw.reducers.sum(t.x_y), sum_x_square=pw.reducers.sum(t.x_square), ) # 5. 由统计量估计回归参数 a、b def compute_a(sum_x, sum_y, sum_x_square, sum_x_y, count): d = count * sum_x_square - sum_x * sum_x if d == 0: return 0 else: return (sum_y * sum_x_square - sum_x * sum_x_y) / d def compute_b(sum_x, sum_y, sum_x_square, sum_x_y, count): d = count * sum_x_square - sum_x * sum_x if d == 0: return 0 else: return (count * sum_x_y - sum_x * sum_y) / d results_table = statistics_table.select( a=pw.apply(compute_a, **statistics_table), b=pw.apply(compute_b, **statistics_table), ) # 6. 实时结果写 CSV pw.io.csv.write(results_table, "regression_output_stream.csv") # 7. 启动引擎!没有它,一切都不会运行 pw.run()两点重要提醒:
- 不要忘记
pw.run():Pathway 是声明式框架,此前所有代码只是搭建计算图;只有调用pw.run()才会真正启动引擎去消费 Kafka 消息并执行计算。 pw.run()不会自行退出:一旦启动,它就会持续监听新消息并增量更新结果,直到进程被外部终止。这正是"常驻流处理"与"一次性批处理"在运行形态上的区别。
4.2generating_kafka_stream.py:Kafka 数据流生成器
from kafka import KafkaProducer import time import random topic = "linear-regression" random.seed(0) def get_value(i): return i + (2 * random.random() - 1)/10 producer = KafkaProducer( bootstrap_servers=["server-address:9092"], sasl_mechanism="SCRAM-SHA-256", security_protocol="SASL_SSL", sasl_plain_username="username", sasl_plain_password="********", ) producer.send(topic, ("x,y").encode("utf-8"), partition=0) time.sleep(5) for i in range(10): time.sleep(1) producer.send( topic, (str(i) + "," + str(get_value(i))).encode("utf-8"), partition=0 ) producer.close()它先发送 CSV 列头"x,y",停 5 秒让 Pathway 连接器就位,然后每秒发一个点共 10 个点:第 $i$ 个点的 $x=i$,$y$ 在 $i$ 附近叠加 ±0.1 的噪声。因此理想回归结果应当是 $(a=0, b=1)$,实测会因噪声而略微偏离。
运行顺序建议:先启动
realtime_regression.py(等待消费),再运行generating_kafka_stream.py开始灌数据,这样能观察每次消息到达带来的增量更新。
第五步:读懂输出——"变更流"语义
本工程的输出有两个 CSV 文件,两者都是 Pathway 的tables of changes(变更表):每次 Kafka 新消息都会触发一次新的计算,引擎把这次计算相对上次的变化写出来,而不是重写全量结果。文件因此自带两列元数据:
time:这次更新所属的处理时间/批次号(不同运行环境中可能体现为批次序号或毫秒级时间戳);diff:+1表示该行新增/更新出现,-1表示此前某行的旧值被撤销。
先看regression_input.csv(收到的原始点,可对照上面生成器的 10 个点):
x,y,time,diff "0","0.06888437030500963",0,1 "1","1.0515908805880605",1,1 "2","1.984114316166169",2,1 "3","2.9517833500585926",3,1 "4","4.002254944273722",4,1 "5","4.980986827490083",5,1 "6","6.056759717806955",6,1 "7","6.9606625452157855",7,1 "8","7.995319390830471",8,1 "9","9.016676407891007",9,1由于输入是只追加的,每行diff恒为+1,数值确实都围绕 $y=x$ 上下浮动。仓库中对应的实际运行样例见 example_regression_input.csv(其time列以毫秒级 Unix 时间戳呈现)。
再看regression_output_stream.csv(回归参数随新点到达的迭代更新过程):
a,b,time,diff 0,0,0,1 0,0,1,-1 0.06888437030500971,0.9827065102830508,1,1 0.06888437030500971,0.9827065102830508,2,-1 0.07724821608916699,0.9576149729305795,2,1 0.0769101730536299,0.9581220374838857,3,1 0.07724821608916699,0.9576149729305795,3,-1 0.05833884879671927,0.9766933617407955,4,1 0.0769101730536299,0.9581220374838857,4,-1 0.05087576879874134,0.9822906717392795,5,1 0.05833884879671927,0.9766933617407955,5,-1 0.03085078333935821,0.9943056630149089,6,1 0.05087576879874134,0.9822906717392795,6,-1 0.03085078333935821,0.9943056630149089,7,-1 0.03590542987734715,0.9917783397459139,7,1 0.03198741430177742,0.9934574892783012,8,1 0.03590542987734715,0.9917783397459139,8,-1 0.025649728471303895,0.9958341214647295,9,1 0.03198741430177742,0.9934574892783012,9,-1这份输出能很直观地看出增量语义:results_table永远只有一行"当前最佳估计",每当新点到达、估计值变化时,引擎就输出"一行+1的新值 + 一行-1的旧值"来替换旧行。比如time=2时估计从(0.06888, 0.98271)更新为(0.07725, 0.95761),于是能看到旧行以-1被撤销、新行以+1插入。收满 10 个点后,估计值收敛到 $a \approx 0.026$、$b \approx 0.996$,非常接近真实直线 $y=x$(即 $a=0,b=1$)。仓库中的完整样例见 example_regression_output_stream.csv。
拿到这套程序后,你可以自由调整数据生成器的参数(样本数量、噪声幅度、目标直线等),观察回归参数如何随数据流实时收敛——这正是理解"增量流计算 vs 重复批计算"差异的最佳实验台。
结语与进阶方向
至此,你已经可以用 Pathway 打通"Kafka 实时数据 → 增量统计 → 模型参数实时刷新"的完整链路,这比传统"攒批 + 重跑回归"的方案更适合对时效性敏感的场景。作为下一步,你可以尝试:
- 把
t.reduce里的聚合列拓展为多个自变量的累积量($\sum x_1^2$、$\sum x_1x_2$ 等),即可把本例推广成多元线性回归,计算骨架无需改变; - 把输出从 CSV 换成 Pathway 的 Kafka 写连接器(
pw.io.kafka.write),把a、b推回消息队列供下游消费; - 阅读仓库内更多 ETL 模板与演示数据流用法,例如人工数据流演示 API 与首个实时应用教程,进一步熟悉连接器矩阵。
本模板的原始文档位于 docs/2.developers/7.templates/ETL/5.linear_regression_with_kafka.md,配套可运行的完整代码则在 examples/projects/kafka-linear-regression,欢迎对照源码逐行验证本文描述。
【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考