【问题标题】:Possibility of using circular buffer in a single thread在单个线程中使用循环缓冲区的可能性
【发布时间】:2016-03-14 17:33:54
【问题描述】:

我有一个 UDP 线程,它通过 recvmmsg 系统调用从不同的多路复用流中读取多个数据报,并将它们推送到不同的循环/环形缓冲区中。这些环形缓冲区是 Stream 结构的一部分。每个流每 20 毫秒发送一个语音帧。所以 UDP 数据包可能看起来像这样:F1S1 F1S2 F1S3 F2S1 等等,或者在突发的情况下它可能看起来像这样:F1S1 F2S1 F3S1 F1S2 等等。接收后,这些数据包将由一个按照ITP 原理工作的库并行处理。 UDP 线程必须分派这些并行任务以及要处理的数据包列表。这里的限制是任务不能并行处理来自相同流的两个帧,并且任务必须有自己的独立内存来处理帧。所以我需要确保这些帧的 FIFO 执行顺序,这将在生成这些任务之前在 UDP 线程中完成。目前,当我收到这些数据包时,我会查找 streamId 并将帧放在循环缓冲区中,该缓冲区是带有 for_loop 的 Stream Strctures 的一部分。

这是显示 UDP 线程中发生的事情的代码。

while (!(*thread_stop))
    {
        int nr_datagrams = recvmmsg(socket_handle->fd_udp, datagramS, VLEN, 0,
                NULL);
        .....
        for (int i = 0; i < nr_datagrams; i++)
        {
        ....
            hash_table_search(codec_buffers, _stream_id, (void **) &codecPtr))
        .....
            // Circular Buffer for speech frames
            // store the incoming sequence number for the newest packet
            codecPtr->circBuff.seqNum[codecPtr->circBuff.newestIdx] = _seq_num;

            // Update the entry pointer to point to the newest frame
            codecPtr->circBuff.entries = codecPtr->circBuff.entries
                    + codecPtr->circBuff.newestIdx * codecPtr->frameLength;
            // Copy the contents of the frame in entry buffer
            memcpy(codecPtr->circBuff.entries,
                    2 * sizeof(uint16_t) + datagramBuff[i],
                    codecPtr->frameLength);
            // Update the newest Index
            codecPtr->circBuff.newestIdx =
                    (codecPtr->circBuff.newestIdx + 1) &
                    codecPtr->circBuffSize;
          }

我的程序现在应该从最近收到数据但不是最新数据的不同流的环形缓冲区中弹出帧,因为所有最近收到的数据包可能属于同一流,以防突发。现在我应该如何前进是我面临的困境?

【问题讨论】:

  • 你问的是环形缓冲区的实现吗?
  • 我认为这里最重要的问题是您的处理逻辑是如何构建的。您当然可以在单个线程中使用循环缓冲区,但它会有用吗?你甚至需要它吗?
  • 是的,它是环形缓冲区。我确实需要它,因为发往特定环形缓冲区的数据报与来自其他缓冲区的其他数据报交错,并且这些数据报的处理需要严格的 FIFO 逻辑。
  • 让我发布我的代码的一部分,这将清除我所追求的要求。
  • 那么,您正在做的是处理乱序的多路复用消息,是吗?

标签: c linux circular-buffer


【解决方案1】:

(这是在 OP 用更多细节澄清问题后的完整重写。)

我建议对接收到的数据报使用无序缓冲区,每个缓冲区都有一个流标识符、流计数器和接收计数器;以及每个流的最新调度计数器:

#define _GNU_SOURCE
#include <stdlib.h>
#include <unistd.h>
#include <limits.h>
#include <sys/types.h>
#include <sys/socket.h>
#include <netinet/in.h>

/* Maximum number of buffered datagrams */
#define  MAX_DATAGRAMS  16

/* Maximum size of each datagram */
#define  MAX_DATAGRAM_SIZE  4096

/* Maximum number of streams */
#define  MAX_STREAMS  4

typedef struct {
    int             stream;                     /* -1 for none */
    unsigned int    counter;                    /* Per stream counter */
    unsigned int    order;                      /* Global counter */
    unsigned int    size;                       /* Bytes of data in data[] */
    char            data[MAX_DATAGRAM_SIZE];
} datagram;


void process(const int socketfd)
{
    /* Per-stream counters for latest dispatched message */
    unsigned int newest_dispatched[MAX_STREAMS] = { 0U };

    /* Packet buffer */
    datagram     buffer[MAX_DATAGRAMS];

    /* Sender IPv4 addresses */
    struct sockaddr_in  from[MAX_DATAGRAMS];

    /* Vectors to refer to packet buffer .data member */
    struct iovec        iov[MAX_DATAGRAMS];

    /* Message headers */
    struct mmsghdr      hdr[MAX_DATAGRAMS];

    /* Buffer index; hdr[i], iov[i], from[i] all refer
     * to buffer[buf[i]]. */
    unsigned int        buf[MAX_DATAGRAMS];

    /* Temporary array indicating which buffer contains
     * the next datagram to be dispatched for each stream */
    int                 next[MAX_STREAMS];

    /* Receive counter (not stream specific) */
    unsigned int        order = 0U;

    int                 i, n;

    /* Mark all buffers unused. */
    for (i = 0; i < MAX_DATAGRAMS; i++) {
        buffer[i].stream = -1;
        buffer[i].size = 0U;
    }

    /* Clear stream dispatch counters. */
    for (i = 0; i < MAX_STREAMS; i++)
        newest_dispatched[i] = 0U;

    while (1) {

        /* Discard datagrams received too much out of order. */
        for (i = 0; i < MAX_DATAGRAMS; i++)
            if (buffer[i].stream >= 0)
                if (buffer[i].counter - newest_dispatched[buffer[i].stream] >= UINT_MAX/2) {
                    /* Either in the past, or too far into the future */
                    buffer[i].stream = -1;
                    buffer[i].size = 0U;
                }

        /* Prepare for receiving new messages.
         * Stream -1 indicates unused/processed message. */
        for (n = 0, i = 0; i < MAX_DATAGRAMS; i++)
            if (buffer[i].stream == -1) {

                /* Prep the buffer. */
                buffer[i].stream = -1;
                buffer[i].counter = 0U;
                buffer[i].order = order + n;
                buffer[i].size = 0U;

                /* Local index n refers to buffer i. */
                buf[n] = i;

                /* Local index n refers to buffer i data. */
                iov[n].iov_base = buffer[i].data;
                iov[n].iov_len  = sizeof buffer[i].data;

                /* Clear received bytes counter. */
                hdr[n].msg_len = 0U;

                /* Source address to from[] array. */
                hdr[n].msg_hdr.msg_name = from + i;
                hdr[n].msg_hdr.msg_namelen = sizeof from[i];

                /* Payload per iov[n]. */
                hdr[n].msg_hdr.msg_iov = iov + n;
                hdr[n].msg_hdr.msg_iovlen = 1;

                /* No ancillary data. */
                hdr[n].msg_hdr.msg_control = NULL;
                hdr[n].msg_hdr.msg_controllen = 0;

                /* Clear received message flags */
                hdr[n].msg_hdr.msg_flags = 0;

                /* Prepared one more hdr[], from[], iov[], buf[]. */
                n++;
            }

        if (n < 1) {
            /* Buffer is full. Find oldest received datagram. */
            unsigned int max_age = 0U;
            int          oldest = -1;

            for (i = 0; i < MAX_DATAGRAMS; i++) {
                const unsigned int age = order - buffer[i].order;
                if (age >= max_age) {
                    max_age = age;
                    oldest = i;
                }
            }

            /* TODO: Dispatch the oldest received datagram:
             * Stream  buffer[oldest].stream
             * Data    buffer[oldest].data, buffer[oldest].size bytes
            */

            /* Update stream counters. */
            newest_dispatched[buffer[oldest].stream] = buffer[oldest].counter;             

            /* Remove buffer. */
            buffer[oldest].stream = -1;
            buffer[oldest].size   =  0;

            /* Need more datagrams. */
            continue;
        }

        n = recvmmsg(socketfd, hdr, n, 0, NULL);
        if (n < 1) {
            /* TODO: Check for errors. */
            continue;
        }

        /* Update buffer description for each received message. */
        for (i = 0; i < n; i++) {
            const int b = buf[i];

            buffer[b].order = order;          /* Already set, actually */
            buffer[b].size  = hdr[i].msg_len;

            /* TODO: determine stream and counter,
             *       based on from[i] and buffer[b].data.
             *       This assigns them in round-robin fashion. */
            buffer[b].stream  = order % MAX_STREAMS;
            buffer[b].counter = order / MAX_STREAMS;

            /* Account for the message received. */
            order++;
        }

    while (1) {

            /* Clear next-to-be-dispatched index list. */
            for (i = 0; i < MAX_STREAMS; i++)
                next[i] = -1;

            /* Find next messages to be dispatched. */
            for (i = 0; i < MAX_DATAGRAMS; i++)
                if (buffer[i].stream >= 0 && buffer[i].counter == newest_dispatched[buffer[i].stream] + 1U)
                    next[buffer[i].stream] = i;

            /* Note: This is one point where you will wish to 
             *       ensure all pending dispatches are complete,
             *       before issuing new ones. */

            /* Count (n) and dispatch the messages. */
            for (n = 0, i = 0; i < MAX_STREAMS; i++)
                if (next[i] != -1) {
                    const int b = next[i];
                    const int s = buffer[b].stream;

                    /* TODO: Dispatch buffer b, stream s. */

                    /* Update dispatch counters. */
                    newest_dispatched[s]++;
                    n++;
                }

            /* Nothing dispatched? */
            if (n < 1)
                break;

            /* Remove dispatched messages from the buffer. Also remove duplicates. */
            for (i = 0; i < MAX_DATAGRAMS; i++)
                if (buffer[i].stream >= 0 && buffer[i].counter == newest_dispatched[buffer[i].stream]) {
                    buffer[i].stream = -1;
                    buffer[i].size = 0U;
                }

        }
    }
}

请注意,我省略了您应该等待发送消息完成的点(因为有多个选项,具体取决于您发送的方式以及您是否希望同时进行“工作”)。此外,此代码仅经过编译测试,因此可能包含逻辑错误。

循环结构如下:

  1. 丢弃过去或未来太远而无用的缓冲消息。

    计数器是循环的。我添加了计数器包装逻辑here的描述。

  2. 为每个空闲缓冲区槽构造recvmmsg() 的标头。

  3. 如果没有可用的缓冲槽,则查找并分派或丢弃最旧的缓冲槽,然后从步骤 1 开始重复。

  4. 接收一条或多条消息。

  5. 根据收到的消息,更新缓冲槽。

    主要是确定流,流计数器,以及接收到的消息中的字节数。

  6. 调度循环。

    这是一个循环,因为如果我们收到乱序的消​​息,但稍后完成它们,我们将需要一次发送多组消息。

    在循环内,首先清除流索引数组(next[])。

    然后,我们检查缓冲区是否有接下来要分派的消息。为此,我们需要每个流的计数器。这是在单独的步骤中完成的,以防我们收到重复的 UDP 数据报。

    如果没有一个流已经缓冲了它们的下一条消息,我们退出这个循环,并等待新的数据报到达。

    接下来发送消息。循环最多为每个流分派一条消息。

    发送后,我们删除发送的消息。我们不是遍历每个流并删除与该流对应的缓冲区,而是遍历整个缓冲区,以便我们也捕获重复的 UDP 消息。

请注意,缓冲区根本没有按上述顺序复制。

如果消息是压缩或未压缩的音频,您确实需要额外的(循环)缓冲区用于未压缩的音频流。为所有 UDP 消息拥有一个共享的无序缓冲区的好处是,您始终可以选择下一个要推进的音频流(如果您已收到该数据报),并且不会意外地花费太多时间推进一个流以致其他流可能会用完数据,导致音频故障。

每个音频流的循环缓冲区的大小应至少是数据报最大大小的三倍。这使您可以对每个样本使用包装逻辑(((later % LIMIT) + LIMIT - (earlier % LIMIT)) % LIMIT,结果 > LIMIT/2 表示逆序),并且即使在播放/解压缩期间也可以附加新数据。 (调度程序更新一个索引,音频播放另一个。只要确保它们是原子访问的。)更大的音频流缓冲区可能会导致更大的延迟。

总而言之,假设音频流解复用和分派已在手边,则需要使用两个完全独立的缓冲区结构。对于 UDP 数据报,使用一组无序的缓冲槽。缓冲槽需要一些记账(如上面的代码所示),但是为了许多不同的流而分派它们非常简单。但是,每个音频流都需要一个循环缓冲区(至少是(解压缩的)数据报最大大小的三倍)。

不幸的是,我认为在这里使用独立任务并行性没有任何好处(例如wool C library)。

事实上,为每个流添加一个结构来描述解压缩器状态可能更简单,并根据哪个循环音频缓冲区剩余缓冲数据最少来确定它们的优先级。典型的解压器会报告它们是否需要更多数据,因此为每个流添加一个临时工作区(两个压缩数据报),将允许解压器消耗整个数据包,但仅在绝对必要时才复制内存。


已编辑以添加有关循环缓冲区的详细信息:

有两种主要的方法可以跟踪循环缓冲区的状态,另外还有第三种派生方法,我认为这里可能有用:

  1. 使用单独的索引来添加 (head) 和删除 (tail) 数据

    如果有一个生产者和一个消费者,循环缓冲区可以无锁维护,因为生产者只增加head,消费者增加tail

    head == tail 时缓冲区为空。如果缓冲区有 SIZE 条目、head = head % SIZEtail = tail % SIZE,则有 (head + SIZE - tail) % SIZE 缓冲条目。

    缺点是一个简单的实现总是在缓冲区中至少有一个空闲条目,因为上面的简单模算术无法区分所有使用的条目和没有使用的条目。对于稍微复杂的代码有一些变通方法。

    在简单的情况下,缓冲区有SIZE - 1 - (head + SIZE - tail) % SIZE 空闲条目。

  2. 缓冲的数据从索引 start 开始,并以缓冲的 length 条目。

    缓冲区内容在内存中总是连续的,或者在内存中分成两部分(第一部分在缓冲区空间的末尾结束,第二部分从缓冲区空间的开头开始)。生产者和消费者都需要修改startlength,因此无锁使用需要比较和交换原子操作(通常将两者打包成一个整数)。

    在任何时候,都有length 条目被使用,而size - length 条目在循环缓冲区中是空闲的。

    当生产者附加 n 数据条目时,它会从索引 (start + length) % SIZE 开始复制数据,最终在索引 (start + length + n - 1) % SIZE 处复制数据,并将 length 递增 n。如前所述,要复制的数据可能是连续的,也可能分为两部分。

    当消费者消费 n 数据条目时,它会从索引 start 开始复制数据,在索引 (start + n) % SIZE 处复制最终条目,并更新 start = (start + n) % SIZE;length = length - n;。同样,消耗的数据可能会在内存中分成两部分(如果它会跨越缓冲区的末尾)。

  3. 衍生物。

    如果只有一个生产者线程/任务和一个消费者,我们可以将缓冲区状态变量加倍,以允许通过 DMA 或异步 I/异步从缓冲区添加或消费数据哦。

    1. 使用 headtailhead_pendingtail_pending 索引

      head != head_pending 时,正在消耗从 headhead_pending-1 的数据,包括在内。完成后,消费者设置head = head_pending % SIZE

      tail != tail_pending 时,在索引tailtail_pending-1(含)处添加更多数据。传输完成后,生产者设置tail = tail_pending % SIZE

      请注意,使用 DMA 时,通常最好使用内存中的连续块。在微控制器中,通常使用中断将下一个 DMA'ble 块加载到 DMA 寄存器中,在这种情况下,您实际上有 headhead_pendinghead_next,或 tailtail_pending , 和tail_next,选择每个 DMA 块的大小,这样您就不会在分割点附近(在缓冲区的物理端)对非常短的块进行 DMA,但保持中断率可接受。

      在任何时候,缓冲区中都有(head + SIZE - tail) % SIZE 条目可供使用。使用简单的模运算,缓冲区中至少有一个条目始终未被使用,因此可以添加的最大条目数为SIZE - 1 - (head + SIZE - tail) % SIZE

    2. 使用startlengthincomingoutgoing

      这里,startlength 必须原子修改,这样对方就无法观察到旧的start 和新的length,反之亦然。如上所述,这可以无锁地完成,但必须小心,因为这是问题的常见来源。

      在任何时候,缓冲区都包含 length 条目,其中添加了 incoming 条目(在 (start + length) % SIZE(start + length + incoming - 1) % SIZE,包括在内,如果 incoming &gt; 0),并且正在消耗 outgoing 条目(在 @ 987654390@ 到 (start + outgoing - 1) % SIZE,如果是 outgoing &gt; 0,则包括在内。

      当传入传输完成时,生产者将length 增加incoming

      当传出传输完成时,消费者更新start = (start + outgoing) % SIZElength = length - outgoing

关于原子处理:

支持 C11 的 C 编译器提供了一系列原子函数,可用于以原子方式更新上述变量。使用弱版本可以最大程度地兼容不同类型的硬件。 对于startlength

    uint64_t buffer_state; /* Actual atomic variable */
    uint64_t old_state;    /* Temporary variable */

    temp_state = atomic_load(&buffer_state);
    do {
        uint32_t start = temp_state >> 32;
        uint32_t length = (uint32_t)temp_state;
        uint64_t new_state;

        /* Update start and length as necessary */

        new_state = (uint64_t)length | ((uint64_t)state << 32);
    } while (!atomic_compare_exchange_weak(&buffer_state, &old_state, new_state));

为了将一些缓冲区状态变量state 增加amount,缓冲区大小为size,假设所有都是size_t 类型:

    size_t old; /* Temporary variable */

    old = atomic_load(&state) % size;
    do {
        size_t new = (old + amount) % size;
    } while (!atomic_compare_exchange_weak(&state, &old, new));

请注意,如果atomic_compare_exchange_weak() 失败,它会将state 的当前值复制到old。这就是为什么只需要一个初始原子负载的原因。

许多 C 编译器提供非标准的 C11 之前的原子内置,只是许多 C 编译器提供的通用扩展。例如,startlength 可以使用

进行原子修改
    uint64_t buffer_state;         /* Actual atomic variable */
    uint64_t old_state, new_state; /* Temporary variables */

    do {
        uint32_t start, length;

        old_state = buffer_state; /* Non-atomic access */

        start = old_state >> 32;
        length = (uint32_t)old_state;

        /* Update start and/or length */

        new_state = (uint64_t)length | ((uint64_t)start << 32);
    } while (!__sync_bool_compare_and_swap(&buffer_state, old_state, new_state));

要在许多 C11 之前的编译器上将一些缓冲区状态变量 state 增加 amount,缓冲区大小为 size,假设所有类型均为 size_t,您可以使用:

    size_t old_state, new_state; /* Temporary variables */

    do {
        old_state = state;
        new_state = (old_state + amount) % size;
    } while (!__sync_bool_compare_and_swap(&state, old_state, new_state));

所有这些原子本质上都是旋转的,直到修改原子地成功。虽然看起来两个或多个并发内核可能会无休止地战斗,但当前的缓存架构使得一个内核总是会获胜(第一)。因此,在实践中,只要每个核心在执行此类原子更新循环之一之间还有其他工作要做,这些都可以正常工作。 (事实上​​,在无锁 C 代码中无处不在。)

我想提到的最后一部分是允许部分分派数据报。这基本上意味着每个数据报缓冲区槽不仅有size(表示该槽中的字节数),还有start。当接收到新数据报时,start 设置为零。如果无法完全分派数据报(复制到每个流缓冲区),则更新缓冲区槽startsize,但流分派计数器不递增。这样,在下一轮,这个数据报的其余部分就被缓冲了。

我可以编写一个完整的示例,展示如何使用我在上一段中提到的部分缓冲数据报方案将传入的数据报从无序数据报缓冲区解压缩为多个流,但具体实现在很大程度上取决于编程接口解压器库有。

特别是,我个人更喜欢使用的界面,例如POSIX iconv() 函数——但可能返回状态代码而不是转换的字符数。各种音频和语音库有不同的接口,甚至可能无法将它们转换成这样的接口。 (作为来自不同领域的示例,大多数用于保护套接字通信的 SSL/TLS 库没有这样的接口,因为它们总是期望直接直接访问套接字描述符;这使得单线程多套接字异步 SSL/TLS 实现“困难”。嗯,更像是“如果你想要的话,从头开始写”。)对于解压缩数据的音频分析,比如使用 FFTW 库进行快速傅里叶变换(或 DCT,或 Hartley,或其他变换之一出色的库性能,尤其是在该窗口大小的转换的优化智慧可用时),通常需要固定大小的块中的解压缩数据。这也会影响确切的实施。

【讨论】:

  • 我在问题中明确提到,我只有一个 UDP 线程,它负责接收和调度。您的整个解决方案基于“独立”处理线程的前提,我不能在我的设计中随意使用。
  • 可能会出现“然后由谁进行处理”的问题。它的独立任务谁进行处理。 UDP 使用需要处理的数据来分派这些任务,并且它不能分派来自同一流的两个连续数据报中包含的两个连续语音帧。
  • UDP 但是可能会在一个 recvmmsg 调用中从同一流中接收两个连续的数据报(因此是连续的语音帧)。那么它们应该如何以先进先出的方式处理,是我所追求的解决方案。
  • 我在这里提出了类似的问题,涉及整个处理链。问题在这里停留了将近 3 周没有答案,所以我删除了它并将问题分成几部分,以便专家可以理解。但似乎也不起作用。
猜你喜欢
  • 2012-04-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-06-11
相关资源
最近更新 更多