设计思路

有两种方式,一种按照接收的完整数据进行存放回放,回放的时候再以udp的形式再将数据转发出来,缺点是不好对数据进行二次处理,回放数据不透明,进度条回退不好实现;另外一种是将接收到的数据拆包解析后存储,缺点是要再写一套数据下发逻辑以及读写文件逻辑。本文选用了后者。

第一步:接收维护所有数据

首先有一个DataManager类用于接收所有数据回调,在DataManager类里包含了DataPlayback数据回放类,两个类都做成了单例类的形式,防止多个地方调用数据混乱错改。
以下仅为一个结构体为例,DataManager类主要包含的函数如下:

/**
 * @file data_manager.h
 * @brief 该头文件定义了 data_manager 类,用于管理所有接收到的数据以及下发的数据。
 * @details 该类负责处理所有接收的数据以及下发的数据,并将所有数据存储在一个容器里方便后续回放。
 * @author 奇树谦
 * @date 2025年10月21日
 */
#ifndef DATAMANAGER_H
#define DATAMANAGER_H

#include "data_playback.h"

#include <QObject>
#include <map>
#include <memory>
#include <mutex>
#include <string>
#include <thread>
#include <chrono>
#include <fstream>
#include <sstream>
#include <iostream>

class DataManager : public QObject
{
    Q_OBJECT

public:
    static DataManager *getInstance();
    ~DataManager();

    void startTimer(); // 启动数据回放相关定时器
    void stopTimer();  // 停止数据回放相关定时器

private:
    DataManager();                                        // 私有构造函数
    DataManager(const DataManager &) = delete;            // 禁止拷贝构造
    DataManager &operator=(const DataManager &) = delete; // 禁止赋值操作

    void init(); // 初始化函数
    // 生成当前时间戳
    std::string generateCurrentTimestamp();

    void receiveGetBaseInfo(const Remote::BaseInfo &data, Remote::CommunicationChannel channel);
    static DataManager *instance; // 静态实例指针

    DataPlayback dataPlayback;		// 数据回放类
    bool is_data_playback;	// 是否处于数据回放状态
};

#endif // DATAMANAGER_H

下面是数据回放包含的所有函数:

/**
 * @file data_playback.h
 * @brief 该头文件定义了 data_playback 类,用于处理数据回放数据。
 * @details 该类负责处理所有接收的数据以及下发的数据,并将所有数据存储在本地文件,用于数据回放。 
 *          数据存储格式:[timestamp] [type_id] [version] [serial_num] [子类字段...]
 * @author 奇树谦
 * @date 2025年10月21日
 */
#ifndef DATAPLAYBACK_H
#define DATAPLAYBACK_H

#include "frontend/data_manager/replay_frame.h"
#include <glog/logging.h>
#include <string>
#include <fstream>
#include <mutex>
#include <unordered_map>
#include <memory>
#include <thread>
#include <chrono>
#include <ctime>
#include <iostream>
#include <functional>
#include <atomic>
#include <queue>

class DataPlayback
{
public:
       static DataPlayback *getInstance()
       {
              static DataPlayback instance;
              return &instance;
       }

       DataPlayback();
       ~DataPlayback();

       void startTimer();
       void stopTimer();
       void addData(std::shared_ptr<RemoteType::BaseData> data);
       void readDataFromFile(const std::string &filepath);
       // 获取时间戳范围
       void getTimeRange(std::string &begin, std::string &end) const;
       // 将时间字符串转换为时间戳 (秒级)
       long long timeStringToTimestamp(const std::string &timeStr);
       // 将时间戳转换为时间字符串
       std::string timestampToTimeString(long long timestamp);
       std::map<long long, std::shared_ptr<ReplayFrame>> read_dataMap; // 时间戳 当前时间戳下所有数据值

private:
       std::string generateFileName();
       void writeHeader();
       void writeDataToFilePeriodically();
       void writeDataToFile();
       void createDirectory(const std::string &folderPath);
       std::string getExecutablePath();
       // 把文件缓冲区立即刷到物理磁盘
       void forceFlushToDisk(FILE *fp);
       // 构建派生类函数表实现
       void WriteBaseInfo(std::ostream &outFile, const RemoteType::BaseInfo &data);

       // 解析文件中数据到结构体中
       bool parseBaseInfo(std::istringstream &iss, RemoteType::BaseInfo &data);

       std::queue<std::shared_ptr<RemoteType::BaseData>> dataQueue;
       std::shared_ptr<ReplayFrame> current_frame_; // 始终存活,持续累积
       std::string current_second_;                 // 当前正在累积的秒级时间戳
       std::mutex dataQueueMutex;
       std::atomic<bool> timerRunning{false};
       std::condition_variable cv; // 用于唤醒线程
       std::thread writeThread;
       FILE *file;

       // 定义写入函数的类型
       using WriteFunction = std::function<void(std::ostream &, const RemoteType::BaseData &)>;

       // 写入函数表
       std::unordered_map<DataType, WriteFunction, RemoteType::DataTypeHash> writeFunctionTable;

       std::string folderPath;
       std::string filePath;

       std::string begin_time_;
       std::string end_time_;
};

#endif // DATAPLAYBACK_H

下面按着数据流对每个函数进行详细的介绍解析

开启数据回放数据录入定时器

首先默认建立连接后开启数据连接,在connect函数里添加

// 开启数据回放记录
DataManager::getInstance()->startTimer();

DataManager类中仅对数据回放类开启进行了调用,所有数据回放类仅在DataManager里进行调用,防止环境污染。

void DataManager::startTimer()
{
    dataPlayback.startTimer();
}

实际DataPlayback实现的定时器如下,我在初始化的时候设置了数据回放文件创建路径,可以自定义解决,也可以使用数据库管理,由于是自定义类型的文件,开头可以根据需求自定义一些字段用来判断是否为合法的数据回放文件,最后开启新线程用于将数据写入文件内,调用winapi和定时器每分钟落一次盘,防止软件崩溃导致的数据丢失,完整代码见文末git链接。


void DataPlayback::startTimer()
{
    LOG(INFO) << "DataPlayback startTimer";
    if (timerRunning.load())
    {
        LOG(WARNING) << "Timer already running.";
        return;
    }

    // 等上次线程彻底结束
    if (writeThread.joinable())
    {
        writeThread.join();
    }

    timerRunning = true;

    // 生成文件名
    std::string fileName = generateFileName();
    filePath = folderPath + "/" + fileName;

    file = fopen(filePath.c_str(), "wb");
    if (!file)
    {
        LOG(ERROR) << "Failed to open file: " << filePath;
        timerRunning = false;
        return;
    }

    writeHeader();
    writeThread = std::thread(&DataPlayback::writeDataToFilePeriodically, this);
}

接收数据

开启数据回放定时器后,每当有数据进来,调用:
第一个参数为类型,第二个参数为数据来源通道,默认只有一个数据通道也可不做定义。

void DataManager::receiveGetBaseInfo(const Remote::BaseInfo &data, Remote::CommunicationChannel channel)
{
    std::string timestamp = generateCurrentTimestamp();
    auto data_ptr = std::make_shared<RemoteType::BaseInfo>(data);
    data_ptr->timestamp = timestamp;
    dataPlayback.addData(data_ptr);

    std::stringstream get_info;
    get_info << analyticCommunicationChannel(channel) << "基本信息" << "\n"
             << "---名称:" << data.name << "\n"
             << "---类型:" << data.type << "\n";

    MessageBoxDialog::GetInstance()->appendOperationLog(QString::fromStdString(get_info.str()));
}

向容器里添加数据

void DataPlayback::addData(std::shared_ptr<RemoteType::BaseData> data)
{
    std::lock_guard<std::mutex> lock(dataQueueMutex);
    dataQueue.push(std::move(data));
}

将容器内的数据定时写入文件中

该函数实现就是开启定时器时开启的写入线程调用的函数,实现每分钟写入一次磁盘

void DataPlayback::writeDataToFilePeriodically()
{
    while (timerRunning.load())
    {
        std::this_thread::sleep_for(std::chrono::minutes(1));
        if (!timerRunning.load())
            break;
        writeDataToFile();
    }
    writeDataToFile();
}

数据缓存使用的是std::queue先进先出,实际写入函数如下:

void DataPlayback::writeDataToFile()
{
    if (!file)
    {
        LOG(ERROR) << "writeDataToFile: file is null!";
        return;
    }

    std::queue<std::shared_ptr<RemoteType::BaseData>> localQueue;
    {
        std::lock_guard<std::mutex> lock(dataQueueMutex);
        localQueue.swap(dataQueue); // 瞬间完成
    } // 马上解锁,让 addData 无阻塞

    LOG(INFO) << "Writing data to file, size: " << localQueue.size();

    while (!localQueue.empty())
    {
        auto base = localQueue.front();
        localQueue.pop();

        std::ostringstream oss;

        // 写入通用字段
        oss << base->timestamp << ' '
            << static_cast<int>(base->type_id) << ' '
            << base->version_num << ' '
            << base->serial_num << ' ';

        // 写入子类字段
        auto it = writeFunctionTable.find(base->type_id);
        if (it != writeFunctionTable.end())
        {
            std::ostringstream content;
            it->second(content, *base);
            oss << content.str();
            oss << '\n';
        }
        else
        {
            LOG(WARNING) << "Unknown data type: " << static_cast<int>(base->type_id);
            continue;
        }

        const std::string &line = oss.str();
        fwrite(line.c_str(), 1, line.size(), file);
    }

    fflush(file);
    forceFlushToDisk(file);
}

fflush在stdio.h文件中实现,forceFlushToDisk如下:
FlushFileBuffers是winapi里的函数


void DataPlayback::forceFlushToDisk(FILE *fp)
{
#ifdef _WIN32
    HANDLE hFile = (HANDLE)_get_osfhandle(_fileno(fp));
    if (hFile != INVALID_HANDLE_VALUE)
    {
        BOOL result = FlushFileBuffers(hFile);
        LOG(INFO) << "FlushFileBuffers result: " << result;
    }
#else
    int fd = fileno(fp);
    if (fd != -1)
    {
        fsync(fd);
    }
#endif
}

writeFunctionTable为写入文件的函数表,里面注册了所有需要记录的数据,注册代码以一个结构体为例如下所示

writeFunctionTable.emplace(DataType::BaseInfo, [this](std::ostream &outFile, const RemoteType::BaseData &data){ WriteBaseInfo(outFile, static_cast<const RemoteType::BaseInfo &>(data)); });

如何读取文件回放

加载文件

dataReplay_为数据回放时间轴UI,如何调用:

/* 后台线程加载,避免界面冻结 */
QtConcurrent::run([=]()
{
	DataPlayback::getInstance()->readDataFromFile(filePath.toStdString());
	std::string b, e;
	// 为了方便获取时间轴左右值
	DataPlayback::getInstance()->getTimeRange(b, e);
	dataReplay_->SetTimeRange(QString::fromStdString(b), QString::fromStdString(e));
});

解析文件

先解析文件头是否符合类型要求,然后将解析后的数据放到容器中,全部放完以后按时间戳进行排序,将当前时刻的所有数据存入std::shared_ptr current_frame_;,每次循环当前时刻的数据存入std::map<long long, std::shared_ptr> read_dataMap; // 时间戳 当前时间戳下所有数据值,具体实现:


void DataPlayback::readDataFromFile(const std::string &filepath)
{
    begin_time_.clear();
    end_time_.clear();
    current_frame_ = nullptr;

    LOG(INFO) << "readDataFromFile";
    std::ifstream inFile(filepath);
    if (!inFile.is_open())
    {
        std::cerr << "The file cannot be opened " << filepath << std::endl;
        return;
    }

    std::string line;
    // ====== 文件头部格式检查可以根据需求自定义 ======
    std::getline(inFile, line);
    if (line != "adp")
    {
        std::cerr << "Invalid file format: missing 'adp'" << std::endl;
        return;
    }

    std::getline(inFile, line);
    if (line.find("format binary_little_endian") == std::string::npos)
    {
        std::cerr << "Invalid file format: expected 'format binary_little_endian ...'" << std::endl;
        return;
    }

    std::getline(inFile, line);
    if (line.find("comment") == std::string::npos)
    {
        std::cerr << "Invalid file format: missing 'comment' line" << std::endl;
        return;
    }

    std::getline(inFile, line);
    if (line != "end_header")
    {
        std::cerr << "Invalid file format: missing 'end_header'" << std::endl;
        return;
    }

    std::vector<std::pair<long long, std::string>> allLines;
    // ====== 正式开始解析数据行 ======
    while (std::getline(inFile, line))
    {
        std::istringstream iss(line);
        std::string date, hms;
        if (!(iss >> date >> hms))
            continue; // 时间解析失败跳过
        std::string full_ts = date + " " + hms;
        long long ts = timeStringToTimestamp(full_ts);
        allLines.emplace_back(ts, line);
    }
    // 按时间戳排序
    std::sort(allLines.begin(), allLines.end(),
              [](const auto &a, const auto &b)
              { return a.first < b.first; });

    current_frame_ = std::make_shared<ReplayFrame>();
    // 解析每行
    for (const auto &[ts, rawLine] : allLines)
    {
        std::istringstream iss(rawLine);
        std::string date, hms;
        iss >> date >> hms;
        int type_id, version;
        std::string serial;
        iss >> type_id >> version >> serial;
        std::string full_ts = date + " " + hms;
        /* 记录首帧 */
        if (begin_time_.empty())
            begin_time_ = full_ts;

        /* 记录末帧(每次覆盖即可) */
        end_time_ = full_ts;

        std::cout << "Parsed timestamp: " << full_ts << "-----Line: " << line << std::endl;

        if (version != 1)
        {
            std::cerr << "The version numbers are incompatible. Skip this line\n";
            continue;
        }

        current_frame_->base.timestamp = full_ts;
        current_frame_->base.serial_num = serial;
        current_frame_->base.type_id = static_cast<DataType>(type_id);
        current_frame_->base.version_num = version;

        switch (current_frame_->base.type_id)
        {
        case DataType::BaseInfo:
        {
            RemoteType::BaseInfo st;
            if (!parseBaseInfo(iss, st))
                continue;
            current_frame_->fill(st);
            break;
        }
        // TODO: 其他类型也加入 switch 分支
        default:
            std::cerr << "Unknown type ID " << type_id << std::endl;
            continue;
        }

        std::cout << "read_dataMap ts: " << ts << std::endl;
        read_dataMap[ts] = std::make_shared<ReplayFrame>(*current_frame_);
    }
}

与UI联动

如此我们就得到了每个时刻的所有数据,配合数据回放UI对数据进行回放,UI头文件如下

#ifndef FRONTEND_DATA_REPLAY_H
#define FRONTEND_DATA_REPLAY_H

#include "frontend/universal/base_dialog.h"
#include "frontend/data_manager/replay_controller.h"
#include "frontend/universal/timeprogressslider.h"
#include "glog/logging.h"

#include <QMutex>
#include <QLineEdit>
#include <QPushButton>
#include <QLabel>
#include <QHBoxLayout>
#include <QVBoxLayout>
#include <QMessageBox>

/**
 * @brief 数据回放主界面UI
 * @details 单例模式,提供:
 *  1. 加载回放文件并设置时间范围
 *  2. 时间轴、倍速、跳转、播放/暂停/停止控制
 *  3. 实时接收底层帧回调驱动滑块
 */
class DataReplay : public BaseDialog
{
    Q_OBJECT

public:
    /**
     * @brief 线程安全单例获取
     */
    static DataReplay *GetInstance(QWidget *parent = nullptr);

    /**
     * @brief 静态析构器,程序退出时自动释放单例
     */
    class Deletor
    {
    public:
        ~Deletor();
    };

    /* 构造/析构 */
    explicit DataReplay(QWidget *parent = nullptr);
    ~DataReplay();

    /**
     * @brief 由外部调用,传入回放文件起止时间并加载数据
     * @param begin 起始时间戳 yyyy-MM-dd hh:mm:ss
     * @param end   终止时间戳 yyyy-MM-dd hh:mm:ss
     */
    void SetTimeRange(QString begin, QString end);

protected:
    void closeEvent(QCloseEvent *event);

signals:
    void signalClearUI();
    void signalCloseEvent();

private:
    void buildUi();

private slots:
    /*----------- 控制按钮槽函数 -----------*/
    void jumpToTimestamp(const QString &ts);  ///< 通用跳转逻辑
    void onFrameArrived(const std::string &); ///< 底层帧回调:驱动滑块
    void onSliderSeek(const QString &ts);     ///< 滑块拖动/点击跳转
    void OnJumpButtonClicked();               ///< “跳转”按钮
    void onPlayPauseToggled();                ///< 播放/暂停 二合一按钮
    void OnStopButtonClicked();               ///< 停止按钮
    void OnLeftSpeedButtonClicked();          ///< 减速
    void OnRightSpeedButtonClicked();         ///< 加速

    void onMinimizeClicked(); // 最小化
    void onRestoreClicked();  // 从悬浮按钮恢复

private:
    static DataReplay *m_pInstance; ///< 单例指针
    static Deletor deleter;         ///< 静态析构器

    /*----------- UI 控件 -----------*/
    QVBoxLayout *main_layout_ = nullptr;   ///< 主布局
    TimeProgressSlider *slider_ = nullptr; ///< 时间轴滑块
    QLineEdit *time_edit_ = nullptr;       ///< 时间戳输入框
    QPushButton *play_btn_ = nullptr;      ///< 播放/暂停按钮
    QPushButton *stop_btn_ = nullptr;      ///< 停止按钮
    QPushButton *left_spd_btn_ = nullptr;  ///< 减速按钮
    QPushButton *right_spd_btn_ = nullptr; ///< 加速按钮
    QLabel *speed_label_ = nullptr;        ///< 倍速显示
    QPushButton *float_btn_ = nullptr;     // 悬浮按钮

    /*----------- 回放控制器 -----------*/
    ReplayController rc;      ///< 底层回放引擎
    bool is_playing_ = false; ///< 当前播放状态
    int m_currentSpeed = 1;   ///< 当前倍速
    bool is_load_ = false;    ///< 是否已加载回放文件

    /*----------- 工具函数 -----------*/
    QString parseTimeString(const QString &text); ///< hh:mm:ss → 完整时间戳
};

#endif // FRONTEND_DATA_REPLAY_H

回放数据定时器

除了最上层的UI类,以及读写文件解析结构体的类,还有最后一个和数据回放相关的类,就是有关控制所有数据如何下发的类,头文件如下

#ifndef REPLAY_CONTROLLER_H
#define REPLAY_CONTROLLER_H

#include "frontend/data_manager/replay_frame.h"
#include <functional>
#include <memory>
#include <mutex>
#include <atomic>
#include <thread>
#include <vector>

/**
 * @brief 离线数据回放控制器
 * @details 支持从 DataPlayback 加载数据,按每秒钟回放,支持暂停、停止、跳转和倍速播放(1x、2x、4x)
 */
class ReplayController
{
public:
    using FrameCallback = std::function<void(const std::string &)>;

    ReplayController();
    ~ReplayController();

    bool isRunning() const { return running_.load(); }
    double getSpeed() const { return speed_.load(); }

    // 控制接口
    void loadFromDataPlayback();      // 加载数据
    void play(FrameCallback cb);      // 播放
    void pause();                     // 暂停
    void stop();                      // 停止并重置游标
    void setSpeed(double speed);      // 设置倍速(支持 1.0 / 2.0 / 4.0)
    void seek(const std::string &ts); // 跳转到指定时间

    // 查询状态
    long long currentTime() const;
    long long totalStartTime() const;
    long long totalEndTime() const;

private:
    void loop(); // 回放线程主循环
    std::atomic<long long> cursor_{0};
    std::atomic<double> speed_{1.0};
    std::atomic<bool> running_{false};
    FrameCallback cb_;
    std::thread worker_;
    mutable std::mutex mtx_;
    std::map<long long, std::shared_ptr<ReplayFrame>> local_map_;
    long long start_time_;
    long long end_time_;
};

#endif // REPLAY_CONTROLLER_H

回放数据下发线程

几个关键函数,回放线程主循环

/* ---------- 回放线程逻辑 ---------- */
void ReplayController::loop()
{
    LOG(INFO) << "ReplayController loop start";
    const int base_interval_ms = 1000;
    auto next_tick = std::chrono::steady_clock::now();
    if (local_map_.empty())
    {
        running_ = false;
        return;
    }
    while (true)
    {
        {
            std::lock_guard<std::mutex> lock(mtx_);
            if (!running_)
                break;
        }

        long long now;
        now = cursor_;
        if (now > end_time_)
            break;

        std::shared_ptr<ReplayFrame> frame = nullptr;
        {
            auto it = local_map_.find(now);
            if (it != local_map_.end())
                frame = it->second;
        }
        
				// 将数据下发到界面
        if (frame)
            Remote::getInstance()->dispatchReplayFrame(*frame);

        // 回调时间戳,无论是否有数据
        if (cb_)
            cb_(DataPlayback::getInstance()->timestampToTimeString(now));

        // 计算下一个 tick
        int delay_ms = static_cast<int>(base_interval_ms / speed_.load());
        next_tick += std::chrono::milliseconds(delay_ms);
        std::this_thread::sleep_until(next_tick);

        cursor_.fetch_add(1, std::memory_order_relaxed);
    }

    running_ = false;
    LOG(INFO) << "ReplayController loop end";
}

关键控制函数

播放控制相关的几个函数:播放、暂停、停止并重置游标、设置倍速(支持 1.0 / 2.0 / 4.0)、跳转到指定时间


/* ---------- 播放控制 ---------- */
void ReplayController::play(FrameCallback cb)
{
    LOG(INFO) << "ReplayController::play()";

    {
        std::unique_lock<std::mutex> lock(mtx_);
        if (worker_.joinable())
        {
            // 通知线程退出
            running_ = false;
        }
    }
    // 必须在解锁后 join,避免死锁
    if (worker_.joinable())
    {
        worker_.join();
    }

    {
        std::unique_lock<std::mutex> lock(mtx_);
        cb_ = cb;
        running_ = true;
        worker_ = std::thread(&ReplayController::loop, this);
    }
}

void ReplayController::pause()
{
    {
        std::lock_guard<std::mutex> lock(mtx_);
        running_ = false;
    }
    if (worker_.joinable())
        worker_.join();
    LOG(INFO) << "ReplayController::pause()";
}

void ReplayController::stop()
{
    pause();
    cursor_.store(start_time_, std::memory_order_relaxed);
    LOG(INFO) << "ReplayController::stop()";
}

void ReplayController::setSpeed(double speed)
{
    double clamped = std::clamp(speed, 1.0, 4.0);
    speed_ = clamped;
    LOG(INFO) << "ReplayController setSpeed: " << clamped << "x";
}

void ReplayController::seek(const std::string &ts)
{
    std::lock_guard<std::mutex> lock(mtx_);
    long long ts_ll = DataPlayback::getInstance()->timeStringToTimestamp(ts);
    cursor_.store(ts_ll, std::memory_order_relaxed);
    // 找出不晚于 cursor_ 的最近帧
    std::shared_ptr<ReplayFrame> nearest;
    for (auto it = local_map_.rbegin(); it != local_map_.rend(); ++it)
    {
        if (it->first <= ts_ll)
        {
            nearest = it->second;
            break;
        }
    }

    if (nearest)
        Remote::getInstance()->dispatchReplayFrame(*nearest);

    LOG(INFO) << "ReplayController seek to: " << ts;
}

其它

switch case和函数表注册两种形式对比,那种形式更好为什么

维度switch-case函数表
编译期检查枚举值不完整时编译器会提醒(-Wswitch无,漏写也不会报错
性能编译器可优化成跳转表;枚举连续时 O(1)理论上也是 O(1),但多一次哈希/查找
代码长度每加一个类型要改两处:定义 + switch只改一处:在表里加 lambda
可读性直观,一眼能看到所有分支如果 lambda 太长会臃肿,可拆成独立函数
动态扩展无法运行时增删可以运行时插入/删除元素
调试断点直接落在 case 内断点落在 lambda 里,可读性稍差

结论

  • 当前项目只有 2~3 个类型,且编译期完整性更重要保留 switch-case
    因为一旦漏写 case,编译器会报警;而函数表不会。

  • 类型非常多(十几、几十种)或需要插件式动态注册函数表才划算。

因此,在少量的代码规模下,switch-case 更优;等类型膨胀到“一眼看不完”时再重构即可。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐