所以我有N个异步的,时间戳的数据流。每条溪流都有固定的速率.我想处理所有的数据,但问题是我必须处理数据,以便尽可能接近数据到达的时间(这是一个实时流应用程序)。
到目前为止,我的实现是创建一个固定的K消息窗口,我使用优先级队列按时间戳对其进行排序。然后,我按照顺序处理整个队列,然后再转到下一个窗口。这是可以的,但这并不理想,因为它造成的滞后与缓冲区的大小成正比,而且有时还会导致消息在缓冲区结束后才到达时丢弃的消息。看起来是这样的:
// Priority queue keeping track of the data in timestamp order.
ThreadSafeProrityQueue<Data> q;
// Fixed buffer size
int K = 10;
// The last successfully processed data timestamp
time_t lastTimestamp = -1;
// Called for each of the N data streams asyncronously
void receiveAsyncData(const Data& dat) {
q.push(dat.timestamp, dat);
if (q.size() > K) {
processQueue();
}
}
// Process all the data in the queue.
void processQueue() {
while (!q.empty()) {
const auto& data = q.top();
// If the data is too old, drop it.
if (data.timestamp < lastTimestamp) {
LOG("Dropping message. Too old.");
q.pop();
continue;
}
// Otherwise, process it.
processData(data);
lastTimestamp = data.timestamp;
q.pop();
}
}有关数据的信息:它们保证在自己的流中进行排序。他们的频率在5到30赫兹之间。它们由图像和其他数据组成。
一些例子说明了为什么这比看上去更难。假设我有两个流,A和B都以1Hz运行,并按以下顺序得到数据:
(stream, time)
(A, 2)
(B, 1.5)
(A, 3)
(B, 2.5)
(A, 4)
(B, 3.5)
(A, 5)看看如果我按接收数据的顺序处理数据,B总是会被丢弃的吗?这就是我在算法中想要的avoid.Now,B每10帧就会被删除一次,我会以10帧的滞后处理数据,直到过去。
发布于 2017-07-07 13:44:28
我建议建立一个生产者/消费者结构。让每个流将数据放入队列,并让一个单独的线程读取队列。这就是:
// your asynchronous update:
void receiveAsyncData(const Data& dat) {
q.push(dat.timestamp, dat);
}
// separate thread that processes the queue
void processQueue()
{
while (!stopRequested)
{
data = q.pop();
if (data.timestamp >= lastTimestamp)
{
processData(data);
lastTimestamp = data.timestamp;
}
}
}这防止了在处理批处理时在当前实现中看到的“滞后”。
processQueue函数在一个单独的持久线程中运行。stopRequested是程序在关闭时设置的标志--迫使线程退出。有些人会为此使用volatile标志。我更喜欢使用类似手动重置事件之类的方法。
要实现这个功能,您需要一个允许并发更新的优先级队列实现,或者需要用同步锁包装队列。特别是,您希望确保当队列为空时,q.pop()等待下一项。或者,当队列为空时,您永远不会调用q.pop()。我不知道你的ThreadSafePriorityQueue的具体内容,所以我不能确切地说你是怎么写的。
时间戳检查仍然是必要的,因为可以在之前的项之前处理后一项。例如:
processQueue函数从队列中删除数据流2中的事件。这并不稀奇,只是不常见。时差通常是微秒级的。
如果你经常得到的更新不正常,那么你可以引入人为的延迟。例如,在您更新的问题中,您显示的消息以500毫秒的顺序出现。让我们假设500毫秒是您要支持的最大容限。也就是说,如果一条信息晚到500毫秒以上,那么它就会被丢弃。
当您将东西添加到优先级队列时,您要做的是在时间戳中添加500 ms。这就是:
q.push(AddMs(dat.timestamp, 500), dat);在处理事物的循环中,您不会在它的时间戳之前排成队列。类似于:
while (true)
{
if (q.peek().timestamp <= currentTime)
{
data = q.pop();
if (data.timestamp >= lastTimestamp)
{
processData(data);
lastTimestamp = data.timestamp;
}
}
}这在所有项目的处理中引入了500 ms的延迟,但它防止了删除在500 ms阈值范围内的“延迟”更新。您必须平衡您的愿望“实时”更新和您的愿望,以防止删除更新。
发布于 2017-07-06 17:03:06
总是有滞后的,这种滞后将取决于你愿意等待你最慢的“固定利率”流的时间。
建议:
不是完全可靠的(每个缓冲区将被排序,但从一个缓冲区到另一个缓冲区,您可能有时间戳反转),但也许足够好?
使用“满意”标志的计数来触发处理(在步骤3)可以用来使延迟更小,但存在更多缓冲区间时间戳倒置的风险。在极端情况下,接受只有一个满意标志的处理意味着“一收到就推送一个帧,时间戳排序就会被诅咒”。
我提到这一点是为了支持我的感觉,即延迟/时间戳反转平衡是你的问题固有的--除非是绝对平等的框架,否则将有一个不牺牲一方的完美解决方案。
由于“解决方案”将是一种平衡的行为,任何解决方案都需要收集/使用额外的信息来帮助决策(例如,“一系列标志”)。如果我给您的建议听起来很愚蠢(很可能,您选择分享的细节并不太多),那么就开始考虑哪些指标将与您的“体验质量”目标级别相关,并使用其他数据结构来帮助收集/处理/使用这些度量标准。
https://stackoverflow.com/questions/44949899
复制相似问题