Pathway 实时流式线性回归实战:从 Kafka 数据流接入到增量参数估计 📅 发布时间:2026/9/8 23:56:42 👁 浏览次数: Pathway 实时流式线性回归实战从 Kafka 数据流接入到增量参数估计【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway导读本文基于 PathwayLive 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.pyPathway 流式回归主程序generating_kafka_stream.pyKafka 数据流生成器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.serversKafka broker 地址host:port形式如server-address:9092或localhost:9092security.protocol传输加密方式托管服务通常为sasl_ssl本地开发可为plaintextsasl.mechanismSASL 认证机制常见SCRAM-SHA-256、SCRAM-SHA-512、PLAINgroup.id消费组 ID决定消息如何被分组消费用你的项目/消费者组命名session.timeout.ms消费组会话超时时间示例给6000sasl.username/sasl.password认证账号密码建议通过环境变量注入勿硬编码1.2 定义 schema 并读取 topicKafka 消息只是字节流要让 Pathway 知道每列是什么类型需要先声明一个pw.Schemaclass InputSchema(pw.Schema): x: float y: float随后一行代码即可建立输入表t pw.io.kafka.read( rdkafka_settings, topiclinear-regression, schemaInputSchema, formatcsv, autocommit_duration_ms1000 )关键参数rdkafka_settings上一步的 librdkafka 风格配置字典topic要订阅的 Kafka topicschema输入表的列结构由pw.Schema子类声明format消息的序列化格式决定字节如何被解析成列autocommit_duration_ms两次提交之间的最大时间间隔引擎会按此节奏把收到的新消息批量推进到计算图中示例为1000毫秒即大约每秒刷新一次结果。在当前仓库的 Kafka 连接器实现中该参数的默认值为1500ms并额外提供modestreaming/static、parallel_readers、with_metadata、start_from_timestamp_ms等选项可在需要控制消费模式、并行度或记录元信息时查阅。关于format的一点版本说明教程撰写时的csv格式在当前仓库的pw.io.kafka.read签名中已演化为显式支持的plaintext、raw、json三种raw把消息原样读入产生一个含data列的表当前实现中 raw/plaintext 还会附带key列见该文件 docstringplaintext把消息按 UTF-8 文本解析同样落在data列json先把 JSON 负载解析出来再按schema中声明的列建表非常适合一条消息 一个点的场景。仓库配套示例 realtime_regression.py 目前使用的正是formatjson。本文为了贴合教程原文主体演示仍采用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即理想直线 yx 幅值 ±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_squaret.x * t.x, x_yt.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( countpw.reducers.count(), sum_xpw.reducers.sum(t.x), sum_ypw.reducers.sum(t.y), sum_x_ypw.reducers.sum(t.x_y), sum_x_squarepw.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( apw.apply(compute_a, **statistics_table), bpw.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)$最后关闭 Producerproducer KafkaProducer( bootstrap_servers[server-address:9092], sasl_mechanismSCRAM-SHA-256, security_protocolSASL_SSL, sasl_plain_usernameusername, sasl_plain_password********, ) producer.send(topic, (x,y).encode(utf-8), partition0) producer.send( linear-regression, (0,0).encode(utf-8), partition0 ) producer.send( linear-regression, (1,1).encode(utf-8), partition0 ) producer.close()为了使回归不那么平凡真实数据几乎都带噪声本例让数据点围绕直线 $yx$ 采样并对每个 $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.pyPathway 处理主程序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, topiclinear-regression, schemaInputSchema, formatcsv, autocommit_duration_ms1000, ) # 2. 把收到的原始输入也旁路写一份便于对照 pw.io.csv.write(t, regression_input.csv) # 3. 扩展列x_square, x_y t t.select( *pw.this, x_squaret.x * t.x, x_yt.x * t.y, ) # 4. 全局累加五个统计量 statistics_table t.reduce( countpw.reducers.count(), sum_xpw.reducers.sum(t.x), sum_ypw.reducers.sum(t.y), sum_x_ypw.reducers.sum(t.x_y), sum_x_squarepw.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( apw.apply(compute_a, **statistics_table), bpw.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.pyKafka 数据流生成器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_mechanismSCRAM-SHA-256, security_protocolSASL_SSL, sasl_plain_usernameusername, sasl_plain_password********, ) producer.send(topic, (x,y).encode(utf-8), partition0) time.sleep(5) for i in range(10): time.sleep(1) producer.send( topic, (str(i) , str(get_value(i))).encode(utf-8), partition0 ) producer.close()它先发送 CSV 列头x,y停 5 秒让 Pathway 连接器就位然后每秒发一个点共 10 个点第 $i$ 个点的 $xi$$y$ 在 $i$ 附近叠加 ±0.1 的噪声。因此理想回归结果应当是 $(a0, b1)$实测会因噪声而略微偏离。运行顺序建议先启动realtime_regression.py等待消费再运行generating_kafka_stream.py开始灌数据这样能观察每次消息到达带来的增量更新。第五步读懂输出——变更流语义本工程的输出有两个 CSV 文件两者都是 Pathway 的tables of changes变更表每次 Kafka 新消息都会触发一次新的计算引擎把这次计算相对上次的变化写出来而不是重写全量结果。文件因此自带两列元数据time这次更新所属的处理时间/批次号不同运行环境中可能体现为批次序号或毫秒级时间戳diff1表示该行新增/更新出现-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数值确实都围绕 $yx$ 上下浮动。仓库中对应的实际运行样例见 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的旧值来替换旧行。比如time2时估计从(0.06888, 0.98271)更新为(0.07725, 0.95761)于是能看到旧行以-1被撤销、新行以1插入。收满 10 个点后估计值收敛到 $a \approx 0.026$、$b \approx 0.996$非常接近真实直线 $yx$即 $a0,b1$。仓库中的完整样例见 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),仅供参考