
使用 Pandas 与 Flower 实现联邦分析Iris 直方图聚合快速入门指南【免费下载链接】flowerFlower: A Friendly Federated AI Framework项目地址: https://gitcode.com/GitHub_Trending/flo/flower本文以 Flower 官方快速入门示例 examples/quickstart-pandas 为核心讲解如何用 Pandas 与 Flower 构建一种**联邦分析Federated Analytics**应用节点不训练模型而是在各自持有的数据上运行查询计算直方图与统计量再将结果发送给ServerApp聚合。读完本文你将掌握低层Message/RecordDictAPI 的完整调用链、ClientApp与ServerApp的编写方式以及如何用flwr run在模拟与部署两种模式下运行同一套代码。[!CAUTION] 本示例使用 Flower 的低层 API该 API 目前仍是预览特性未来可能发生变化。ClientApp与ServerApp直接操作 Message 和 RecordDict 对象而非高层Client/Strategy抽象。示例要解决的问题联邦分析而非联邦训练传统联邦学习Federated Learning通过聚合本地训练的模型权重来协作训练一个共享模型。而本示例实现的是一种联邦分析数据不出本地节点节点各自对拥有的数据执行查询仅把查询得到的统计结果上传聚合。这样做既保护了数据隐私又能获得全局视角的数据洞察。具体到本示例查询任务是对 Iris 数据集特定列计算直方图客户端在本地数据上计算SepalLengthCm、SepalWidthCm两列的直方图频次、加权均值与样本计数这些指标被打包进MetricRecord并通过Message发送给ServerApp服务端将各节点的部分直方图逐 bin 累加得到全局直方图。数据方面示例使用 Flower Datasets仓库中对应的 Python 包位于 datasets/flwr_datasets下载、分区并预处理 scikit-learn/iris 数据集运行起来非常简单。项目结构与核心文件执行flwr new flwrlabs/quickstart-pandas后会生成如下结构的目录quickstart-pandas ├── pandas_example │ ├── __init__.py │ ├── client_app.py # 定义 ClientApp │ └── server_app.py # 定义 ServerApp ├── pyproject.toml # 项目元数据依赖与应用配置 └── README.md在本文所在仓库中该示例位于 examples/quickstart-pandas其中pandas_example/client_app.py客户端逻辑负责加载本地分区数据、计算直方图pandas_example/server_app.py服务端逻辑负责下发查询、聚合结果pyproject.toml声明依赖并通过[tool.flwr.app]段注册应用组件与运行配置。环境准备与获取应用首先安装 Flowerpip install flwr然后从 Flower Hub 拉取应用模板flwr new flwrlabs/quickstart-pandas接着安装pyproject.toml中声明的依赖并以可编辑模式安装pandas_example包pip install -e .从本示例的 pyproject.toml 可以看到实际依赖为[project] name quickstart-pandas version 1.1.11 dependencies [ flwr[simulation]1.36.0, flwr-datasets[vision]0.6.1, numpy2.0.2, pandas2.2.3, ]其中flwr-datasets[vision]提供 Iris 数据集的下载、分区与预处理能力flwr[simulation]为模拟引擎提供支持。客户端实现数据本地查询打开 client_app.py可以看到客户端全部逻辑由app.query()装饰的函数承载# Flower ClientApp app ClientApp() app.query() def query(msg: Message, context: Context): Construct histogram of local dataset and report to ServerApp. # 从 node_config 读取本节点对应的数据分区 partition_id context.node_config[partition-id] num_partitions context.node_config[num-partitions] dataset get_clientapp_dataset(partition_id, num_partitions) metrics {} # 对 DataFrame 的每一列计算统计量 for feature_name in dataset.columns: # 计算直方图 freqs, _ np.histogram(dataset[feature_name], binsnp.linspace(2.0, 10.0, 10)) metrics[feature_name] freqs.tolist() # 计算加权平均均值 × 样本数 metrics[f{feature_name}_avg] dataset[feature_name].mean() * len(dataset) metrics[f{feature_name}_count] len(dataset) reply_content RecordDict({query_results: MetricRecord(metrics)}) return Message(reply_content, reply_tomsg)数据分区加载get_clientapp_dataset借助FederatedDataset与IidPartitioner实现同一份公开数据集按节点切分的模拟def get_clientapp_dataset(partition_id: int, num_partitions: int): # 只初始化一次 FederatedDataset全局缓存 global fds if fds is None: partitioner IidPartitioner(num_partitionsnum_partitions) fds FederatedDataset( datasetscikit-learn/iris, partitioners{train: partitioner}, ) dataset fds.load_partition(partition_id, train).with_format(pandas)[:] # 仅保留本示例关心的两列 return dataset[[SepalLengthCm, SepalWidthCm]]其中partition_id与num_partitions来自context.node_config由 Flower 运行环境自动注入IidPartitioner(num_partitions...)对 Iris 做独立同分布切分模拟各客户端持有不同子集with_format(pandas)将数据集转为 Pandas DataFrame 格式。直方图与统计量计算对两列特征各做三件事指标键计算方式用途SepalLengthCm/SepalWidthCmnp.histogram(..., binsnp.linspace(2.0, 10.0, 10))得到 9 个 bin 的频次数组各节点部分直方图服务端按 bin 求和{feature}_avgcolumn.mean() * len(column)加权均值供服务端推导全局均值{feature}_countlen(column)样本数用于验证聚合正确性注意{feature}_avg故意以均值 × 样本数形式上传这正是加权聚合的前提——服务端拿到各节点(sum, count)后可恢复全局均值。返回结构结果通过RecordDictMetricRecord封装reply_content RecordDict({query_results: MetricRecord(metrics)}) return Message(reply_content, reply_tomsg)从框架源码 recorddict.py 可见RecordDict是存储数组、指标与配置的统一载体值的类型必须是ArrayRecord、MetricRecord或ConfigRecord三者之一对应_check_value的类型校验并作为Message内容在ClientApp与ServerApp之间传输。MetricRecord正是用于承载标量型指标此处为直方图频次列表与统计量。服务端实现查询下发与聚合打开 server_app.py服务端逻辑集中在app.main()装饰的函数中app ServerApp() app.main() def main(grid: Grid, context: Context) - None: num_rounds context.run_config[num-server-rounds] min_nodes 2 fraction_sample context.run_config[fraction-sample] for server_round in range(num_rounds): log(INFO, ) # 增加空行便于阅读日志 log(INFO, Starting round %s/%s, server_round 1, num_rounds) # 轮询等待足够数量的节点上线 all_node_ids: list[int] [] while len(all_node_ids) min_nodes: all_node_ids list(grid.get_node_ids()) if len(all_node_ids) min_nodes: # 按比例采样节点 num_to_sample int(len(all_node_ids) * fraction_sample) node_ids random.sample(all_node_ids, num_to_sample) break log(INFO, Waiting for nodes to connect...) time.sleep(2) log(INFO, Sampled %s nodes (out of %s), len(node_ids), len(all_node_ids)) # 为每个采样节点构造一条 QUERY 消息 recorddict RecordDict() messages [] for node_id in node_ids: message Message( contentrecorddict, message_typeMessageType.QUERY, # 对应 ClientApp 的 query 方法 dst_node_idnode_id, group_idstr(server_round), ) messages.append(message) # 发送并等待全部回复 replies grid.send_and_receive(messages) log(INFO, Received %s/%s results, len(replies), len(messages)) # 聚合部分直方图 aggregated_hist aggregate_partial_histograms(replies) log(INFO, Aggregated histogram: %s, aggregated_hist)Grid服务端与节点的通信抽象grid是 Grid 抽象基类 的实例其核心方法get_node_ids()返回当前所有已连接节点的 ID 列表见 grid.pysend_and_receive(messages, timeoutNone)将消息推送给dst_node_id指定的节点并阻塞等待全部回复或超时见 grid.py。在框架实现中send_and_receive本质上是推送 拉取的组合push_messages发送消息后返回消息 IDpull_messages依据消息 ID 从 SuperLink 取回回复。这里使用单条send_and_receive调用使代码在每一轮中先等待所有节点返回再进行聚合。消息类型与分发MessageType.QUERY是消息路由的关键框架的ClientApp类见 client_app.py定义了train、evaluate、query三种动作方法分别对应MessageType.TRAIN、MessageType.EVALUATE、MessageType.QUERY其中query方法在 client_app.py 定义。因此服务端发出的QUERY消息会自动命中客户端的app.query()处理器——这正是本示例无需训练/评估代码即可运行的机制。此外group_idstr(server_round)把同一轮的消息归入一个组便于按轮次跟踪。部分直方图的聚合def aggregate_partial_histograms(messages: Iterable[Message]): aggregated_hist {} total_count 0 for rep in messages: if rep.has_error(): continue query_results rep.content[query_results] # 逐 bin 累加直方图 for k, v in query_results.items(): if k in [SepalLengthCm, SepalWidthCm]: if k in aggregated_hist: aggregated_hist[k] np.array(v) else: aggregated_hist[k] np.array(v) if _count in k: total_count v # 校验聚合直方图的总频次应等于上报的样本总数 assert total_count sum([sum(v) for v in aggregated_hist.values()]) return aggregated_hist这段代码展示了三个关键细节容错rep.has_error()跳过执行失败的节点回复聚合对部分失败是鲁棒的向量化累加aggregated_hist[k] np.array(v)将各节点的频次数组按 bin 位置相加得到全局直方图正确性校验断言各列直方图频次之和 全部节点上报的_count之和从数据一致性层面验证了聚合结果。运行配置pyproject.toml 中的 App 注册示例的运行配置全部集中在 pyproject.toml 的[tool.flwr.app]段[tool.flwr.app] publisher flwrlabs fab-format-version 1 flwr-version-target 1.37.0 [tool.flwr.app.components] serverapp pandas_example.server_app:app clientapp pandas_example.client_app:app [tool.flwr.app.config] num-server-rounds 3 fraction-sample 1.0各配置项含义配置项默认值说明serverapppandas_example.server_app:app指向ServerApp实例的导入路径flwr run据此加载服务端clientapppandas_example.client_app:app指向ClientApp实例的导入路径num-server-rounds3联邦分析的轮数对应context.run_config[num-server-rounds]fraction-sample1.0每轮采样的节点比例1.0表示全部节点参与对应context.run_config[fraction-sample]注意run_config与node_config的区别run_config由pyproject.toml或命令行注入、对全部节点统一生效node_config由运行环境按节点注入如本示例中的partition-id、num-partitions每个节点各不相同。运行项目同一套代码两种引擎flwr run默认使用模拟引擎Simulation Engine它把所有客户端进程放在本机无需手动启动任何组件是入门 Flower 的首选方式。模拟模式运行flwr run . --stream--stream让日志实时流式输出到终端便于观察每轮节点采样 → 查询下发 → 回复聚合的过程。也可以覆盖pyproject.toml中定义的运行配置例如把轮数改为 5flwr run . --run-config num-server-rounds5 --stream--run-config的键值会合并进context.run_config优先级高于pyproject.toml中的默认值fraction-sample同样可覆盖例如--run-config fraction-sample0.5将只采样半数节点。部署模式运行同一个应用也可以在不改动任何代码的前提下改用**部署引擎Deployment Engine**运行此时需要分别启动 SuperLink 与若干 SuperNode客户端运行在真实或远程节点上。如需在生产化环境中使用还可以进一步为联邦启用基于 TLS 的安全通信以及 SuperNode 认证若已有部署引擎使用经验也可通过 Docker 容器化部署整个联邦。端到端数据流回顾将客户端、服务端与框架机制串联起来一次完整查询的数据流如下ServerApp.main()通过grid.get_node_ids()等待至少 2 个节点上线按fraction-sample采样节点为每个节点构造MessageType.QUERY消息通过grid.send_and_receive()下发各节点的ClientApp.query()被触发从context.node_config取分区号、加载本地 Iris 子集计算两列直方图与统计量封装为RecordDict(MetricRecord(...))的Message回复ServerApp收到全部回复后跳过失败节点将各节点频次数组逐 bin 累加并用_count总和做一致性断言聚合直方图通过日志输出完成一轮联邦分析。这种查询下发—本地执行—统计上传—服务端聚合的模式正是联邦分析区别于联邦训练的典型形态节点永远不共享原始数据只共享统计结果。小结与延伸本示例用不到百行代码完整演示了 Flower 低层 API 在联邦分析场景的落地联邦分析范式以查询代替训练app.query()MessageType.QUERY构成完整调用链数据分区FederatedDatasetIidPartitioner让客户端在本地模拟持有私有子集通信载体RecordDict/MetricRecord/Message作为统一的跨节点数据结构其类型约束与校验逻辑见 recorddict.py服务端编排Grid抽象了节点发现与消息收发send_and_receive让多轮查询的编写保持简洁双模式运行flwr run在模拟与部署两种引擎间无缝切换--run-config支持运行时参数覆盖。想深入了解低层 API 的其他动作如app.train()、app.evaluate()及其与高层Client/Strategy的映射关系可以继续阅读框架源码 client_app.py 与 grid.py或在仓库的 examples 目录下对照其他快速入门示例进行实践。【免费下载链接】flowerFlower: A Friendly Federated AI Framework项目地址: https://gitcode.com/GitHub_Trending/flo/flower创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考