You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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()

关键逻辑说明

  1. 样本缓冲区:用列表实时累计所有接收的样本,避免数据丢失,直到凑够70000个再拆分列
  2. 列拆分机制:每次缓冲区样本数达标,就拆分出完整的一组作为新列,保证每列严格对应70000个样本
  3. CSV写入:用Pandas的DataFrame统一管理所有列,每次新增列后直接覆盖写入CSV,确保文件始终包含所有已收集的数据
  4. 中断处理:捕获手动中断信号,保存剩余未凑够数量的样本,避免中途丢失数据

额外提示

  • 如果处理超大规模数据,担心内存占用过高,可以每写入几列就将DataFrame导出并清空,只保留当前缓冲区的样本
  • 若需要实时追加写入而非覆盖,需注意CSV的列对齐问题,这种场景下仍推荐用DataFrame覆盖写入,Pandas会自动处理列的对齐

内容的提问来源于stack exchange,提问作者Yash Aherrao

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 18:10:37