news 2026/9/8 23:56:40

Pathway 实时流式线性回归实战:从 Kafka 数据流接入到增量参数估计

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Pathway 实时流式线性回归实战:从 Kafka 数据流接入到增量参数估计

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聚合语义带来的好处(下文的计算实现会展开)。

整个工程涉及两条数据通路:

  1. Producer 侧:用kafka-pythonKafkaProducer往 Kafka topic 写入带噪声的数据点;
  2. 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.pygenerating_kafka_stream.py中通过环境变量UPSTASH_KAFKA_USERUPSTASH_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.serversKafka broker 地址,host:port形式server-address:9092localhost:9092
security.protocol传输加密方式托管服务通常为sasl_ssl,本地开发可为plaintext
sasl.mechanismSASL 认证机制常见SCRAM-SHA-256SCRAM-SHA-512PLAIN
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_readerswith_metadatastart_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(列xy)之后,计算分三步走。

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.xt.x * t.y新增x_squarex_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_xsum_ysum_x_ysum_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 消息时,需要遵守两条规则:

  1. 第一条消息必须是列头,例如"x,y",否则连接器无法完成列映射;
  2. 列头只能发送一次——若重复发送,第二条列头会被当作普通数据行解析。

下面用kafka-pythonKafkaProducer演示:先发列头,再发两个点 $(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 实时数据 → 增量统计 → 模型参数实时刷新"的完整链路,这比传统"攒批 + 重跑回归"的方案更适合对时效性敏感的场景。作为下一步,你可以尝试:

  1. t.reduce里的聚合列拓展为多个自变量的累积量($\sum x_1^2$、$\sum x_1x_2$ 等),即可把本例推广成多元线性回归,计算骨架无需改变;
  2. 把输出从 CSV 换成 Pathway 的 Kafka 写连接器(pw.io.kafka.write),把ab推回消息队列供下游消费;
  3. 阅读仓库内更多 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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/8 23:55:35

PyTorch实现对偶GAN图像去雾:从原理到工程实战

简介:基于PyTorch实现图像去雾的对偶生成对抗网络,是一个包含完整Python源码、项目说明及详细代码注释的毕业设计项目。项目针对雾气导致图像对比度下降、细节丢失等问题,利用生成器与判别器相互对抗的方式恢复清晰无雾图像,适合计…

作者头像 李华
网站建设 2026/9/8 23:55:14

AI画板实测:GPT-6 Astra在原理图与PCB设计中的能力与局限

把同一个电源域的电容分两排放在芯片两侧,结果回流路径被拉得很长,纹波指标差了30%。这种问题AI不一定能看出来,但要靠它把所有细节都安排到位,现阶段还不现实。哪些可以放心交给AI适合让GPT-6 Astra处理的,是那些“规…

作者头像 李华
网站建设 2026/9/8 23:52:26

RS-485收发器选型实测:从MAX485到THVD1550,8款芯片横向对比

去年秋天,我接手的一批无刷电机控制器在客户车间的配电柜合闸瞬间,总线上挂着的半双工RS-485收发器一片接一片被打穿。上位机一直报通信超时,现场测A、B线对地电阻,好几块板子只有几十欧姆,拆下来看,清一色…

作者头像 李华
网站建设 2026/9/8 23:51:25

PyTorch实战:基于CRNN+CTC的车牌识别全流程解析

1. 为什么把车牌识别当作 PyTorch 实战项目1.1 车牌识别看起来简单,实际上卡在哪儿先说一个反直觉的现象:车牌识别在工程里看起来特别成熟,门口停车场、高速收费口都在用,似乎是个“老掉牙”的需求。但当我自己把任务拆开&#xf…

作者头像 李华
网站建设 2026/9/8 23:51:11

AI网关全面测评:从API网关到MAI Gateway的七类方案对比

最近大半年,只要聊到 AI 应用落地,迟早会碰到一个绕不开的基础设施话题——AI 网关。模型越来越多,调用协议五花八门,团队既要接 OpenAI 又要兼容国产模型,还要控成本、管权限、做审计,这时候靠业务代码一层…

作者头像 李华