River 未发布版本更新详解:CategoricalNB 新分类器、流式读取加固与 DenStream 聚类修复

发布时间:2026/10/10 20:20:45

River 未发布版本更新详解:CategoricalNB 新分类器、流式读取加固与 DenStream 聚类修复 人工智能机器学习流处理数据分析【免费下载链接】river Online machine learning in Python项目地址https://gitcode.com/gh_mirrors/river12/river点击查看免费下载本篇技术指南以仓库中的未发布版本发布说明 docs/releases/unreleased.md 为核心骨架逐一剖析 RiverPython 在线机器学习库下一个版本在drift、naive_bayes、linear_model、stream、preprocessing、cluster六大子包中的改动既有新增的CategoricalNB分类器与BayesianLinearRegression分布预测接口调整也有流式缓存原子化写入、CSV/SQL 读取资源管理、DenStream 聚类半径与邻域扩展等一批值得关注的修复。读完本文你将清楚这些改动的动机、底层实现机制以及对应的源码位置便于在升级 River 后快速理解行为变化、定位相关代码。从发布说明看 River 的演进脉络该文档是 River 当前主分支上尚未正式发布的变更记录按子包分节记录。整体改动体现出三个明显取向类型与接口规范化river.drift整体通过严格 mypy 检查、stream.cache/iter_csv/iter_sql也进入严格类型检查名单BayesianLinearRegression用predict_dist_one取代predict_one(..., with_distTrue)使其更贴合Regressor接口约定。数据加载与缓存的健壮性stream.Cache改为先写临时文件、流结束后再原子改名落盘避免中断导致截断文件被后续读取当作完整数据集iter_csv、iter_sql补齐了空行处理与资源释放。算法正确性修复cluster.DenStream的微簇半径、邻域扩展与键回收三处缺陷被修正preprocessing.Normalizer不再对零向量抛出ZeroDivisionError。下面按模块逐一展开。drift类型检查收紧与检测器细节修复river.drift达到严格 mypy 标准发布说明指出river.drift子包现在在严格 mypy 模式下保持干净并且pyproject.toml中非严格覆盖non-strict overrides列表里的river.drift.*条目已被移除。也就是说该子包的公开签名与 docstring 虽未变化但其内部实现包括drift/binary/下的检测器、adwin.py、kswin.py、page_hinkley.py等已满足更严格的类型约束。这与仓库中 pyproject.toml 的[tool.mypy]与[[tool.mypy.overrides]]配置第 218 行起相印证——类型检查覆盖范围在持续扩大。_annotations_aggregated的潜在缺陷修复river.drift.datasets.base.ChangePointFileDataset._annotations_aggregated位于 river/drift/datasets/base.py此前存在两处隐患读取了一个并不存在的self._annotations属性正确属性名是构造方法中赋值的self.annotations在intersection分支中用整数0作为 annotator 键去索引而实际键是字符串0。修复后的实现为if annotator_aggregation union: annotations: set[int] set() for annotator in self.annotations: annotations.update(self.annotations[annotator]) elif annotator_aggregation intersection: annotations set(self.annotations[0]) for annotator in self.annotations: annotations.intersection_update(self.annotations[annotator]) elif annotator_aggregation majority: counts: dict[int, int] {} ...该方法按union并集、intersection交集、majority多数表决三种方式聚合多个标注者annotator标注的变点位置。由于当前仓库中该方法没有调用方因此修复后对外行为不变属于先修好、留待后用的防御性改动。DriftRetrainingClassifier的错误指示统一为boolDriftRetrainingClassifierriver/drift/retrain.py是一个包装任意分类器的漂移重训封装它监控模型准确率上的概念漂移与预警检测到预警时启动后台模型训练检测到漂移时用后台模型替换当前模型。其_update_detector中原本以整数0/1形式向漂移检测器传递分类是否正确的指示现在改为显式的boolincorrectly_classifies: bool y_pred ! y self.drift_detector.update(incorrectly_classifies)bool是int的子类因此对已存在的检测器实现完全兼容但语义更清晰也符合严格类型检查下对二元指示变量的约定。结合其learn_one中先更新检测器、再更新模型的调用顺序_update_detector→model.learn_one读者可以理解该包装器的完整工作流。naive_bayes新增 CategoricalNB 与多后端批量支持CategoricalNB面向类别型特征的朴素贝叶斯本次发布的核心新增之一是naive_bayes.CategoricalNBriver/naive_bayes/categorical.py一个专门处理类别型特征的朴素贝叶斯分类器。它为每个特征、每个取值、每个类别分别维护出现频次并支持加性Laplace/Lidstone平滑同时提供在线learn_one/predict_one与小批量learn_many/predict_many两种模式。核心参数参数默认值说明alpha1.0加性平滑参数Laplace/Lidstone 平滑设为0表示不平滑其底层数据结构为class_counts每个类别出现的次数collections.Counter与feature_counts按特征 → 取值 → 类别 → 频次组织的嵌套defaultdict。条件概率计算在p_feature_given_class中完成def p_feature_given_class(self, f: str, v, c: str) - float: num self.feature_counts.get(f, {}).get(v, {}).get(c, 0.0) self.alpha feature_total sum( values.get(c, 0.0) for values in self.feature_counts.get(f, {}).values() ) n_categories max(1, len(self.feature_counts[f])) den feature_total self.alpha * n_categories return num / den即P(特征值 | 类别) (频次 alpha) / (特征总频次 alpha × 类别数)分母中的类别数以该特征实际见过的取值数为准。预测时通过joint_log_likelihood计算各类别的联合对数似然log P(c) Σ log P(x_i | c)再由predict_proba_one用logsumexp归一化为概率。在线模式示例来自 docstring可直接运行 from river import naive_bayes X [ ... {outlook: sunny, temp: hot}, ... {outlook: sunny, temp: hot}, ... {outlook: sunny, temp: mild}, ... {outlook: sunny, temp: cool}, ... {outlook: rainy, temp: cool}, ... {outlook: rainy, temp: cool}, ... {outlook: rainy, temp: mild}, ... {outlook: rainy, temp: hot}, ... {outlook: overcast, temp: cool}, ... {outlook: overcast, temp: mild}, ... {outlook: overcast, temp: mild}, ... {outlook: overcast, temp: hot}, ... ] y [no, no, no, no, yes, yes, yes, yes, ... yes, yes, yes, yes] model naive_bayes.CategoricalNB(alpha1) for x, yi in zip(X, y): ... model.learn_one(x, yi) model.predict_one({outlook: sunny, temp: mild}) no model.predict_proba_one({outlook: sunny, temp: mild}) {no: 0.755..., yes: 0.244...}小批量模式示例CategoricalNB的learn_many通过 narwhals 将输入的DataFrame如 pandas转换为底层 NumPy 数组按列特征与类别分组累加频次joint_log_likelihood_many则用向量化方式逐类别填充似然矩阵最后经to_native_frame还原为与输入同后端的输出保留原有索引。 import pandas as pd df pd.DataFrame(X) y pd.Series(y) batch_model naive_bayes.CategoricalNB(alpha1) batch_model.learn_many(df, y) unseen pd.DataFrame([{outlook: rainy, temp: cool}]) batch_model.predict_many(unseen) 0 yes dtype: object batch_model.predict_proba_many(unseen) no yes 0 0.109900 0.890100参考实现为 scikit-learn 的CategoricalNB见 docstring ReferencesRiver 版额外支持流式在线更新适合类别型特征持续到达的在线场景。MultinomialNB 的多后端支持与BaseNB.predict_many抽取同一发布周期内naive_bayes.MultinomialNBriver/naive_bayes/multinomial.py的learn_many、predict_many、predict_proba_many从仅支持 pandas改为接受任意 narwhals 支持的 eager 后端pandas、polars、pyarrow 等并且在输出时保留输入后端包括 pandas 的索引。其learn_many内部先将目标y做 one-hot 稀疏编码再通过y_one_hot X一次完成按类别的频次聚合最后按列切分稀疏矩阵更新feature_counts。与此配套基类 river/naive_bayes/base.py 新增了后端无关的BaseNB.predict_many它调用各子类实现的predict_proba_many用np.argmax在概率矩阵上取各类别 argmax 作为预测标签从而把概率 → 标签的公共逻辑上移供GaussianNB、MultinomialNB、BernoulliNB、ComplementNB、CategoricalNB等所有变体共享。predict_proba_many本身也以special.logsumexp沿样本轴归一化对数似然。linear_modelBayesianLinearRegression 的分布预测接口调整linear_model.BayesianLinearRegressionriver/linear_model/bayesian_lin_reg.py本次用predict_dist_one取代predict_one(..., with_distTrue)。新的predict_dist_one在返回点估计之外还能给出完整的预测分布同时更好地遵循Regressor接口——不再靠一个可选关键字参数临时切换行为而是提供独立的方法名。从源码看predict_dist_one先调用predict_one得到后验均值Bishop 书公式 3.58再按公式 3.59 计算预测方差1/beta x^T Σ_N x其中Σ_N是维护在_ss_inv_arr中的后验协方差对从未见过的特征则叠加1/alpha的先验方差项最终返回一个proba.Gaussian分布对象 model linear_model.BayesianLinearRegression() x, _ next(iter(datasets.TrumpApproval())) model.predict_one(x) 43.855... model.predict_dist_one(x) (μ43.85..., σ1.00...)其背后的在线更新机制值得一提无平滑时learn_one用 Sherman-Morrison 秩 1 更新维护后验协方差utils.math.sherman_morrison见 river/utils/math.py精度矩阵与自然参数均为精确累加器开启smoothing取值 01用于对抗概念漂移时则退化为标准的 EMA 加权 矩阵求逆。learn_many将逐行秩 1 更新折叠成批量矩阵乘法X^T X与X^T y有平滑时按行权重(1-s)·s^(N-1-k)精确等价于逐行循环。stream缓存原子化与资源管理加固Cache写盘原子化stream.Cacheriver/stream/cache.py用于把一次性流缓存到磁盘后续迭代直接从 pickle 二进制缓存读取从而加速重复遍历例如多次读取 CSV 的场景。此前的问题是第一次遍历被中断break、异常、生成器被遗弃时磁盘上会留下截断文件后续每次遍历都会把这份截断数据当作完整数据集读回。本次修复把写盘流程改为先写临时文件、流耗尽后再原子改名descriptor, partial_path tempfile.mkstemp(dirself.directory, suffix.part) try: with os.fdopen(descriptor, wb) as file: pickler pickle.Pickler(file) for element in stream: pickler.dump(element) yield element except BaseException: with contextlib.suppress(OSError): os.remove(partial_path) raise os.replace(partial_path, path) self.keys.add(key)关键点在于os.replace只在循环正常结束后执行因此最终落盘的{key}.river_cache.pkl一定是完整的任何BaseException包括KeyboardInterrupt与生成器遗弃都会触发临时文件清理绝不会留下半截缓存。缓存文件以.river_cache.pkl后缀标识模块级常量CACHE_SUFFIXCache构造时还会扫描目录中已有的缓存键支持clear(key)与clear_all()清理。使用方式来自 docstring from river import datasets from river import stream dataset datasets.Phishing() cache stream.Cache() for x, y in cache(dataset, keyphishing): # 首次边遍历边缓存 ... pass for x, y in cache(dataset, keyphishing): # 再次直接读缓存更快 ... pass cache.clear(phishing) # 清理单个缓存默认缓存目录按平台推断Linux/Darwin 为/tmpWindows 为C:\TEMP也可通过directory参数指定。iter_csv跳过空行stream.iter_csvriver/stream/iter_csv.py基于csv.DictReader的覆盖类DictReader实现支持target、converters、parse_dates、drop、drop_nones、fraction随机采样比例、compression自动推断.gz/.zip解压、seed、field_size_limit等参数。此前文件中间出现空行时会产出一个空的x字典现在DictReader.__next__在values为空时直接continue跳过与csv.DictReader及iter_arff的行为保持一致def __next__(self) - dict[FeatureName, typing.Any]: while True: values next(self.reader) if not values: continue if self.fraction 1 and self.rng.random() self.fraction: continue return dict(zip(self.names, values))iter_csv资源清理与field_size_limit恢复iter_csv现在即使流未被耗尽例如调用方提前break也会关闭它自己打开的文件并恢复csv.field_size_limit。这得益于contextlib.ExitStack进入函数时先把field_size_limit的回滚与文件的关闭都注册到栈上无论生成器以何种方式退出都会执行清理调用方传入的 buffer 则不会被关闭文件的所有权仍归调用方with contextlib.ExitStack() as stack: if field_size_limit is not None: previous_limit csv.field_size_limit(field_size_limit) stack.callback(csv.field_size_limit, previous_limit) if isinstance(filepath_or_buffer, (str, os.PathLike)): buffer stack.enter_context(utils.open_filepath(filepath_or_buffer, compression)) else: buffer filepath_or_buffer ...iter_sql释放游标stream.iter_sqlriver/stream/iter_sql.py现在会关闭其迭代的 result 对象使底层 cursor 在流耗尽或中断时被释放。实现上使用with conn.execute(statement) as result:上下文管理器包裹整个迭代result退出作用域即自动关闭。该函数接受字符串 SQL 或任意 SQLAlchemy 可执行对象如sqlalchemy.select的产物配合target_name指定目标字段 from river import stream with engine.connect() as conn: ... query sqlalchemy.sql.select(t_sales) ... dataset stream.iter_sql(query, conn, target_nameamount) ... for x, y in dataset: ... print(x, y) {shop: Hema, date: datetime.date(2016, 8, 2)} 20 {shop: Ikea, date: datetime.date(2016, 8, 2)} 18 ...值得留意的是 SQLAlchemy 默认会预取结果若希望真正的逐行流式读取可在连接上设置stream_resultsTrue并非所有数据库引擎都支持。另外iter_sql的query与conn参数现在按 SQLAlchemy 2.0 类型进行类型检查sqlalchemy从忽略检查改为参与检查。preprocessingNormalizer 对零向量的处理preprocessing.Normalizerriver/preprocessing/scale.py将特征向量缩放为单位范数默认order2即 L2 范数常用于feature_extraction.TFIDF之后。此前若输入是零向量范数为 0除零会抛出ZeroDivisionError。修复后借由safe_div辅助函数——分母为 0 时返回0.0——零向量被原样返回def safe_div(a, b): Return a if b is nil, else divides a by b. return a / b if b else 0.0 def transform_one(self, x): norm utils.math.norm(x, orderself.order) return {i: safe_div(xi, norm) for i, xi in x.items()}对零向量而言每个分量都是 0返回0.0在数值上与原样返回等价且不再中断流式处理。这也是safe_div在StandardScaler等缩放器中处理零方差特征时的通用策略。clusterDenStream 三项修复cluster.DenStreamriver/cluster/denstream.py是基于密度的演化数据流聚类算法维护 p-micro-cluster潜在微簇与 o-micro-cluster离群微簇两组结构。本次修复了三处问题微簇半径计算修正微簇半径此前使用了平方和向量的范数norm of the vector of squared sums而非其各分量之和导致离开原点后半径恒为 0。修复后的calc_radius正确地对各维度的平方和累加def calc_radius(self, timestamp): inv_n 1.0 / self.N sum_ls_sq 0.0 sum_ss 0.0 squared_sum self.squared_sum for key, ls_k in self.linear_sum.items(): sum_ls_sq ls_k * ls_k sum_ss squared_sum[key] diff sum_ss * inv_n - sum_ls_sq * inv_n * inv_n return math.sqrt(diff) if diff 0 else 0.0这里sum_ss是各维平方和分量之和Σ s_k而非向量范数√Σ s_k²半径公式对应方差形式的簇半径。相应地docstring 示例改用epsilon1.0默认值为0.02保证示例数据能顺利形成预期聚类。DenStream 的完整默认参数为decaying_factor0.25、beta0.75、mu2、epsilon0.02、n_samples_init1000、stream_speed100其中beta必须落在(0, 1]。predict_one的邻域扩展DenStream.predict_one实现的是对 p-micro-cluster 的 DBSCAN 式密度聚类先对微簇打标签再生成簇。此前邻域扩展存在缺陷邻居的邻居只有在已有标签时才会进入队列导致一串相连的微簇被拆成多个簇。修复后在 BFS 扩展种子队列时只要邻居的邻居尚未打标签就将其入队从而保证密度可达的微簇链归属同一簇while seed_queue: if labels[seed_queue[0]] is not None: seed_queue.popleft() continue if seed_queue: labels[seed_queue[0]] c neighbor_neighbors self._query_neighbor(seed_queue[0]) for neighbor_neighbor in neighbor_neighbors: if labels[neighbor_neighbor] is None: seed_queue.append(neighbor_neighbor)_query_neighbor基于微簇权重大于mu且中心距离小于2*epsilon且不超过calc_radius判断直接密度可达_is_directly_density_reachable。微簇键回收问题新建 p-micro-cluster / o-micro-cluster 时此前用len(...)生成键会复用已删除微簇的键从而覆盖尚存的微簇。修复后统一用当前最大键 1分配新键self.p_micro_clusters[max(self.p_micro_clusters, default-1) 1] closest_omc self.o_micro_clusters[max(self.o_micro_clusters, default-1) 1] mc_from_pmax(..., default-1) 1保证新键严格递增且永不与现有键冲突同时 o-micro-cluster 成长为 p-micro-cluster权重超过mu * beta时会先从o_micro_clusters中删除再插入p_micro_clusters两者键空间各自独立递增避免误覆盖。总结这份未发布版本发布说明覆盖了 River 六个子包的改动可归纳为四类收益接口现代化CategoricalNB补齐了类别型特征的朴素贝叶斯支持含在线与批量双模式BayesianLinearRegression.predict_dist_one让分布预测成为一等公民BaseNB.predict_many使批量预测逻辑在所有朴素贝叶斯变体间共享。后端扩展MultinomialNB及朴素贝叶斯家族的批量方法全面拥抱 narwhals 多后端pandas、polars、pyarrow 等输出保留输入后端与索引。健壮性提升Cache原子化落盘杜绝截断缓存iter_csv跳过空行并可靠释放自开文件、恢复field_size_limititer_sql自动关闭 result 释放游标Normalizer容忍零向量。算法与类型正确性DenStream 的半径、邻域扩展与键回收三项修复以及drift、stream子包严格 mypy 覆盖的落实。若想深入验证或跟进这些改动可在仓库中直接阅读对应实现river/naive_bayes/categorical.py、river/naive_bayes/multinomial.py、river/naive_bayes/base.py、river/linear_model/bayesian_lin_reg.py、river/stream/cache.py、river/stream/iter_csv.py、river/stream/iter_sql.py、river/preprocessing/scale.py、river/cluster/denstream.py、river/drift/retrain.py以及类型检查配置 pyproject.toml。赞分享人工智能机器学习流处理数据分析【免费下载链接】river Online machine learning in Python项目地址https://gitcode.com/gh_mirrors/river12/river点击查看免费下载相关推荐River 0.21.2 版本解读Polars 依赖解耦、ODAC 时序聚类与 ARF 修复River 0.21.2 版本解读Polars 依赖解耦、ODAC 时序聚类与 ARF 修复 本篇文章围绕 River Online machine l人工智能机器学习流处理数据分析Presto 0.144.3 版本发布详解Planner 类型计算修复、TRY 编译器修复与跨 HDFS 实例 Symlink 读取Presto 0.144.3 版本发布详解Planner 类型计算修复、TRY 编译器修复与跨 HDFS 实例 Symlink 读取 本篇文章围绕 Prest大数据数据库后端OpenSearch 3.8.0 版本发布详解数据流管理、聚合新能力与安全加固全览OpenSearch 3.8.0 版本发布详解数据流管理、聚合新能力与安全加固全览 导读 本文基于 opensearch.release notes 3.8.搜索引擎全文检索可观测性数据分析上一篇readthedocs.org 团队结构与支持体系深度解析下一篇QQ空间历史说说一键备份免费开源GetQzonehistory把你的十年回忆完整搬回本地创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/10/10 20:20:45

Django学籍管理系统开发实战:数据模型、权限与Excel导入导出

1. 学籍管理系统到底在管什么:先别急着写代码看到这个标题,估计不少人第一反应是"又是一个CRUD"——没错,学籍管理系统本质上就是学生信息的新增、删除、修改、查询,但真要做过的人才知道,这套CRUD背后牵扯的…

2026/10/10 20:20:45

区间次方和刷题笔记:从暴力到前缀和与树状数组

这篇牛客刷题记录2,记录的是我最近两周在牛客网上集中刷"区间次方和"的完整过程。本来只想随便找几道题保持手感,结果这个看起来不起眼的考点,把前缀和、快速幂、取模、数据结构更新全都串了起来。如果你也在牛客刷题,或…

2026/10/10 20:20:44

ICPC杭州站五题复盘:Trie离线计数、分组背包与树哈希实战

2022ICPC杭州站打完到现在,每次复盘我还是会翻K、A、C、G、M这五道题的提交记录。这篇是个人复盘向的题解,不是官方标程汇编,核心是把每道题从“读题”到“建模”再到“写代码”的完整链路重新走一遍。K题是字符串加Trie离线计数,…

2026/10/10 21:10:50

好消息与坏消息:如何建立不被情绪绑架的消息处理机制

1. 好消息与坏消息的真相:先别急着高兴,也别急着崩溃你肯定有过这种时刻:手机一震,屏幕上弹出一条消息,你心跳加速,点开之后要么想唱歌要么想砸手机。但过了一个星期回头看,当初那个让你兴奋得整…

2026/10/10 21:10:50

WorkBuddy FDE 90天路径:从一句话需求到上线App的实战指南

一句话需求丢过来,三周后要看到能装进手机里的东西,这种场景在不少小团队里反复上演。WorkBuddy FDE 这套打法,就是冲着这种"需求模糊、时间紧、人手少"的处境来的。它把从一句话到上线 App 的全过程拆成可执行的阶段,核…

2026/10/10 21:10:50

LL(1)分析法实现IF-ELSE翻译程序:四元式与真假链回填

简介:一份面向编译原理学习者的IF-ELSE条件语句翻译程序设计资料,基于LL(1)预测分析法,完成词法分析、语法分析并输出四元式中间代码。资料以Visual Studio工程形式组织,共17个文件,包含C源码、头文件、工程配置文件&a…

2026/10/10 21:10:50

Vibe Coding实战:用Cursor+SDD+Claude Code建立可控AI开发链路

1. Vibe Coding不是让AI写代码,是在和需求反复博弈先说个大家可能都有的经历:拿到Cursor第一周,感觉很爽,让它生成个函数、写个页面,几乎都是秒出。但两周之后,项目越做越乱,AI生成的代码散落各…

2026/10/10 21:05:50

AnyPS5:跨平台异构硬件通用运行环境的设计与实现

1. 项目缘起与核心定位AnyPS5 这个名字第一次出现在我视野里的时候,我正蹲在一堆拆机件中间,手里攥着一块从旧设备上拆下来的定制主板,琢磨着怎么把它的算力榨干。当时脑子里冒出来的念头很直接:能不能做一个足够通用的软硬件框架…

2026/10/10 7:31:36

Jev+Agent接管浏览器:browser-use实战与jev-ultrafast性能优化

1. 从“Jev”说起:为什么我要把Agent接进浏览器“Jev”这个词最近在圈子里出现的频率越来越高,很多人第一次听到会以为是某个新模型的名字,其实它更像是一种思路——把Jev模型的能力当作底座,通过Agent的方式去接管浏览器&#xf…

2026/10/9 20:15:56

多智能体集群实战:DeepAgents编排、MCP与A2A协议及Skills体系

1. 从"单兵作战"到"集群协同":多智能体编排到底在解决什么问题如果你最近在折腾 Agent 相关的东西,大概率会有一种感觉:单个 Agent 能做的事情,其实很快就摸到天花板了。你给它一个提示词,挂几个工…

2026/10/8 6:05:44

无源低通滤波器设计实战:从RC到LC,手把手教你避开那些坑

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/10 0:04:53

从逻辑门到计算机:数字电路核心原理与全加器搭建实战

如果你拆过一台旧电脑的主板,盯着那些黑乎乎的小芯片看上一会儿,可能会冒出同一个疑问:这堆引脚密集的元件,到底是怎么“变”出那么复杂的应用的?答案并不在某个神秘的部件里,而是在所有芯片内部都在反复使…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
☎咨询二维码 ☎ ↑