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

如何用Streamz.Dask、Matplotlib和Tkinter实时显示图表与直方图?

问题背景

现有一套基于线程池、Tkinter和Matplotlib的代码,用于处理由另一进程写入文件的信号,进程间通过读写同一文件实现同步。当前使用忙等待函数逐字符读取直到换行:

def readline_char_by_char(file,fig):
    line = ''
    while True:
        char = file.read(1)
        if not char:
            fig.canvas.flush_events()
            #time.sleep(0.01)
            continue
        line += char
        if char == '\n':
            return line

核心处理逻辑为逐行读取数据,积累1000个信号点后执行FFT滤波、CFD检测、电荷积分、能量/PSD计算,同时更新Tkinter标签和Matplotlib图表。现需用Streamz+Dask替代现有方案,保留Matplotlib和Tkinter的实时显示功能。


替代方案实现步骤

1. 替换文件读取逻辑:用Streamz监控文件

Streamz提供from_textfile方法,可实时监控文件新增内容并自动按行读取,替代手动实现的忙等待循环,无需处理字符级读取逻辑:

from streamz import Stream

# 初始化文件流,监控目标文件,按行读取新增内容
source = Stream.from_textfile("your_signal_file.txt", mode="r", chunk_size=1)

2. 构建数据流管道:分阶段处理信号

将原循环中的步骤拆分为数据流节点,用Dask实现计算密集型任务的并行处理:

  • 解析行数据:将每行字符串转为(x,y)数值对
  • 批量积累数据:每1000个点为一批,对应原逻辑中len(signaly) == 1000的触发条件
  • 并行计算信号指标:用Dask处理FFT滤波、CFD检测、电荷积分等任务
  • 更新UI和图表:将计算结果回调到Tkinter主线程更新标签与图表

示例管道代码:

import dask.bag as db
import numpy as np

# 1. 解析每行数据
def parse_line(line):
    if not line.strip():
        return None
    x, y = line.strip().split()
    return {"x": float(x), "y": float(y)}

parsed = source.map(parse_line).filter(lambda x: x is not None)

# 2. 每1000个点为一批
batched = parsed.batch(1000)

# 3. 用Dask并行处理每批信号
def process_batch(batch):
    # 提取x、y数组
    signalx = [d["x"] for d in batch]
    signaly = [d["y"] for d in batch]
    
    # 执行原逻辑中的计算步骤
    filtered_signal_y = filter_fft(signaly)
    success, istart, fintercept, y_sub_bipolar, y_baseline_restored, y_shifted, y_attenuated = find_true_cfd(np.multiply(-1, signaly), d)
    
    if not success or istart < 0:
        return None
    
    i_delayed, i_total, y_filtered_baselineCorrected = charge_comparison_method(signalx, filtered_signal_y, istart, delayed_start=istart+l2, end=istart+l3)
    energy = slope * i_total + intercept
    psd = i_delayed / i_total
    
    # 返回需要保存和绘图的数据
    return {
        "signalx": signalx, "signaly": signaly, "filtered_y": filtered_signal_y,
        "istart": istart, "fintercept": fintercept, "i_delayed": i_delayed, "i_total": i_total,
        "energy": energy, "psd": psd, "y_shifted": y_shifted, "y_sub_bipolar": y_sub_bipolar
    }

# 用Dask并行处理批次数据
processed = batched.map(lambda batch: db.from_sequence([batch]).map(process_batch).compute()).filter(lambda x: x is not None)

# 4. 初始化全局统计容器(对应原逻辑中的列表)
I_Delayed = []
I_Total = []
PSD = []
Energy = []
TOF = []
num_signal = 0

# 5. 定义UI与图表更新函数(必须在Tkinter主线程执行)
def update_ui_and_plots(result):
    global num_signal
    num_signal += 1
    
    # 更新Tkinter标签
    pulse_label.config(text=f"Number of Pulses: {num_signal}")
    
    # 更新Matplotlib实时信号图
    first_plot_line_removal(axs[0][0], lines)
    lines[0].set_data(result["signalx"], result["signaly"])
    lines[1].set_data(result["signalx"], np.multiply(-1, result["y_shifted"]))
    lines[2].set_data(result["signalx"], result["y_sub_bipolar"])
    first_plot_vertical_line_drawing(axs[0][0], result["istart"], l1, d, l3)
    
    # 更新全局统计数据并写入文件
    I_Delayed.append(result["i_delayed"])
    I_Total.append(result["i_total"])
    PSD.append(result["psd"])
    Energy.append(result["energy"])
    TOF.append(result["fintercept"])
    
    true_cfd_out_file.write(f"{result['fintercept']}\n")
    integration_file.write(f"{result['i_delayed']} {result['i_total']}\n")
    energy_file.write(f"{result['energy']}\n")
    psd_file.write(f"{result['psd']}\n")
    energy_psd_file.write(f"{result['energy']} {result['psd']}\n")
    
    # 每1000批更新汇总图表
    if len(PSD) % 1000 == 0:
        update_plots(I_Delayed, I_Total, Energy, PSD, TOF, axs, fig)
        fig.canvas.draw()
        fig.canvas.flush_events()

3. 处理Tkinter主线程同步问题

Streamz与Dask的处理在后台线程执行,而Tkinter的UI操作必须在主线程进行,需用after方法切换到主线程执行更新:

# 安全回调:将更新操作切换到Tkinter主线程
def safe_update(result):
    app.after(0, update_ui_and_plots, result)

# 将处理结果绑定到安全回调
processed.sink(safe_update)

4. 启动数据流与Tkinter主循环

最后启动Streamz数据流并运行Tkinter主循环:

# 启动Streamz数据流监控
source.start()

# 运行Tkinter主循环
app.mainloop()

# 主循环退出后停止数据流
source.stop()

关键优势

  • 替代手动忙等待:Streamz的文件监控更高效,自动处理文件新增内容
  • 并行计算提速:Dask将计算密集型任务分配到多核心,提升处理效率
  • 逻辑解耦:数据流管道将读取、解析、计算、UI更新拆分为独立节点,代码更易维护

内容的提问来源于stack exchange,提问作者Ayan Banerjee

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 16:44:51