如何用Python DataFrame处理ZMQ接收器接收的实时连续数据
解决ZMQ实时接收峰值数据并按70000样本/列写入CSV的问题
原代码存在的问题
- 未定义变量
data,直接运行会报错 - 同时混用Pandas和原生
csv模块,逻辑混乱,没有统一的样本累计机制 - 完全没实现“累计70000个样本后切换到下一列”的核心需求
修正后的代码
import zmq import numpy as np import pandas as pd # 全局配置 SAMPLE_PER_COLUMN = 70000 # 每列固定存储70000个样本 CSV_FILE_PATH = "peak_data.csv" # 输出CSV文件路径 def consumer(): # 初始化ZMQ PULL套接字,连接数据发送端 context = zmq.Context() consumer_receiver = context.socket(zmq.PULL) consumer_receiver.connect("tcp://127.0.0.1:5557") # 初始化样本缓冲区(临时存未凑够70000的样本)和结果DataFrame sample_buffer = [] df = pd.DataFrame() try: while True: # 接收ZMQ二进制数据,解析为float32格式的数组 buff = consumer_receiver.recv() received_data = np.frombuffer(buff, dtype="float32") # 将新接收的样本添加到缓冲区 sample_buffer.extend(received_data.tolist()) # 循环检查缓冲区是否足够拆分出完整的列样本 while len(sample_buffer) >= SAMPLE_PER_COLUMN: # 取出前70000个样本作为新列 new_column_samples = sample_buffer[:SAMPLE_PER_COLUMN] # 移除缓冲区中已使用的样本 sample_buffer = sample_buffer[SAMPLE_PER_COLUMN:] # 生成新列的名称(比如column_1、column_2) column_name = f"column_{len(df.columns) + 1}" # 将新列加入DataFrame,自动对齐行索引 df[column_name] = pd.Series(new_column_samples) # 将完整的DataFrame写入CSV,覆盖原文件(确保所有列都保存) df.to_csv(CSV_FILE_PATH, index=False) print(f"已完成第{len(df.columns)}列写入,共{len(new_column_samples)}个样本") except KeyboardInterrupt: # 处理手动中断(Ctrl+C),保存缓冲区剩余的样本作为最后一列 if sample_buffer: column_name = f"column_{len(df.columns) + 1}" df[column_name] = pd.Series(sample_buffer) df.to_csv(CSV_FILE_PATH, index=False) print(f"程序中断,已保存剩余{len(sample_buffer)}个样本为最后一列") # 关闭ZMQ资源 consumer_receiver.close() context.term() if __name__ == "__main__": consumer()
关键逻辑说明
- 样本缓冲区:用列表实时累计所有接收的样本,避免数据丢失,直到凑够70000个再拆分列
- 列拆分机制:每次缓冲区样本数达标,就拆分出完整的一组作为新列,保证每列严格对应70000个样本
- CSV写入:用Pandas的DataFrame统一管理所有列,每次新增列后直接覆盖写入CSV,确保文件始终包含所有已收集的数据
- 中断处理:捕获手动中断信号,保存剩余未凑够数量的样本,避免中途丢失数据
额外提示
- 如果处理超大规模数据,担心内存占用过高,可以每写入几列就将DataFrame导出并清空,只保留当前缓冲区的样本
- 若需要实时追加写入而非覆盖,需注意CSV的列对齐问题,这种场景下仍推荐用DataFrame覆盖写入,Pandas会自动处理列的对齐
内容的提问来源于stack exchange,提问作者Yash Aherrao
相关产品推荐
相关产品推荐

