跳过正文

为什么多线程会出现问题的总结

作者
杨全烨
系统软件:操作系统、网络与分布式系统。
目录

多线程下RDMA的崩溃问题 多线程下的RDMA的崩溃问题(无锁队列IO模型) 上述的两个问题都和RDMA通知模型相关. 对 socket/unix,ConnectionType 里 update_state 为 NULL,connUpdateState 基本是空操作。 对 RDMA,updateRdmaState 会 直接调用数据面:

static void updateRdmaState(struct connection *conn) {
	rdma_connection *rdma_conn = (rdma_connection *)conn;
	connRdmaSetRwHandler(conn);
	connRdmaEventHandler(NULL, -1, rdma_conn, 0);
}

而 connRdmaEventHandler 会 poll CQ,并在未 postpone 时 循环调用 read_handler(把已 DMA 到 rx.addr 的数据交给上层):

static void connRdmaEventHandler(struct aeEventLoop * el, int fd, void * clientData, int mask) {
        rdma_connection * rdma_conn = (rdma_connection * ) clientData;
        connection * conn = & rdma_conn - > c;
        struct rdma_cm_id * cm_id = rdma_conn - > cm_id;
        RdmaContext * ctx = cm_id - > context;
        int ret = 0;

        UNUSED(el);
        UNUSED(fd);
        UNUSED(mask);

        ret = connRdmaHandleCq(rdma_conn);
        if (ret == C_ERR) {
            conn - > state = CONN_STATE_ERROR;
            return;
        }

        /* uplayer should read all */
        // 认为上层应该读取所有的数据,每次更新状态都可能重入readQueryFromClient.
        while (!(rdma_conn - > flags & RDMA_CONN_FLAG_POSTPONE_UPDATE_STATE) && ctx - > rx.pos < ctx - > rx.offset) {
            if (conn - > read_handler && (callHandler(conn, conn - > read_handler) == C_ERR)) {
                return;
            }
        }

问题是怎么被RDMA + IO多线程引入的: 用因果链说清楚:

  1. IO 线程读完数据后,主线程在 processClientIOReadsDone 里会 connSetPostponeUpdateState(c->conn, 0) 然后 connUpdateState(c->conn)(你当前 unstable 上仍是这样,见下段引用)。
  2. 对 RDMA,connUpdateState → updateRdmaState → connRdmaEventHandler → connRdmaHandleCq + 可能多次 read_handler。也就是说:主线程在处理「某次 IO 线程刚读完」的收尾时,会立刻再跑一轮传输层读逻辑。
  3. 这与 「IO 线程刚把该 client 的读标成 COMPLETED、主线程正在做 parse / batch / 状态机」 叠在一起,容易出现:
    • 重入:再次触发读路径、再次尝试 把同一 client 丢进 IO 线程(trySendReadToIOThreads),或
    • 状态与假设不一致:例如某处 serverAssert(c->io_read_state == ...) 认为「此时只应处于 IDLE/COMPLETED 之一」,但 RDMA 已在同栈帧里又改了一轮读状态。
  4. TCP 不太会这样:update_state 为 NULL,不会在 processClientIOReadsDone 中间泵出一轮新读;RDMA 则 把「泵 CQ + 调 read_handler」绑在了 connUpdateState 上,所以 同样的 networking 代码路径,在 RDMA 上语义更重,断言或重入问题就暴露出来。

# [BUG] RDMA: Assertion ‘c->cmd_queue.len == 0’ failed with high pipelining under new I/O queue model
#

为什么一直出问题的都是这个断言?

/* Parse one or more commands from the query buf.
 *
 * This function may be called from the main thread or from the I/O thread. 主线程或者IO线程都可以进行调用.
 *
 * Sets the client's read_flags to indicate the parsing outcome. If multiple
 * commands could be parsed, additional parsed commands are stored in the
 * client's command queue. */
// 这里的意思应该是解析buffer.
void parseInputBuffer(client *c) {
    /* The command queue must be emptied before parsing. */
    serverAssert(c->cmd_queue.len == 0);

一次解析前,命令队列必须是空的。 不允许在「上一轮解析留在队列里的命令还没被 processInputBuffer 消费路径处理完」时再开一轮 parseInputBuffer,否则两套解析状态会叠在一起,队列语义也不清晰。 主线程里消费顺序在 processInputBuffer 的循环里:先尽量 consumeCommandQueue 弹出已解析命令执行,只有队列弹空了才会再 parseInputBuffer 去解析 querybuf 里剩余字节。 主线程在循环内部不断解析指令:

int processInputBuffer(client *c) {
    /* Parse the query buffer and/or execute already parsed commands. */
    while ((c->querybuf && c->qb_pos < sdslen(c->querybuf)) ||
           c->cmd_queue.off < c->cmd_queue.len) {
        if (!canParseCommand(c)) {
            break; 
        }
        // ...
        /* If commands are queued up, pop from the queue first */
        // 也就是当队列为空的时候才会进行调用.
        if (!consumeCommandQueue(c)) {
            parseInputBuffer(c);
            prepareCommandQueue(c);
        }
        // ...

可以看这个函数:

/* Pops a command from the command queue and sets it as the client's current
* command. Returns true on success and false if the queue was empty. */
static bool consumeCommandQueue(client * c) {
    cmdQueue * queue = & c - > cmd_queue;
    // 这里就是当队列为空的时候才会返回false.
    if (queue - > off >= queue - > len) return false;
    parsedCommand * p = & queue - > cmds[queue - > off++];
    
    /* Combine the command's read flags with the client's read flags. Some read
    * flags describe the client state (AUTH_REQUIRED) while others describe the
    * command parsing outcome (PARSING_COMPLETED). */
    c - > read_flags |= p - > read_flags;
    c - > argc = p - > argc;
    c - > argv = p - > argv;
    c - > argv_len = p - > argv_len;
    c - > argv_len_sum = p - > argv_len_sum;
    c - > net_input_bytes_curr_cmd = p - > input_bytes;
    c - > parsed_cmd = p - > cmd;
    c - > slot = p - > slot;
    if (queue - > off == queue - > len) {
        /* The queue is empty. Don't free it here, because if parsing is done in
        * I/O threads, we want to free it in I/O threads too, to avoid
        * fragmentation. */
        queue - > off = queue - > len = 0;
    }
    return true;
}

I/O 线程路径不走这个循环,而是直接解析(见下节),因此同样受 parseInputBuffer 入口 len == 0 的约束. ioThreadReadQueryFromClient 在 I/O 线程里读完 querybuf 后,直接调用 parseInputBuffer(没有先走 processInputBuffer 那套「先消费队列」逻辑):

void ioThreadReadQueryFromClient(client *c) {
    serverAssert(c->io_read_state == CLIENT_PENDING_IO);

    /* Read */
    readToQueryBuf(c);
    // 这里就是直接进行解析的操作.
    parseInputBuffer(c);
    trimCommandQueue(c);
    prepareCommandQueue(c);

高 pipeline(例如 -P 256)时,一次 parseInputBuffer 往往会在 cmd_queue 里塞很多条已解析命令(len 很大),同时 querybuf 里可能还有未解析的尾巴(或后续 RDMA 又写入新数据)。主线程稍后要配合 addCommandToBatchAndProcessIfFull / processClientsCommandsBatch 等慢慢消费这些队列项。

  1. 第一次 I/O 线程:parseInputBuffer 成功,cmd_queue.len 已是 256(举例),主线程还没开始消费。
  2. 主线程 processClientIOReadsDone 过早 connUpdateState → 嵌套 readQueryFromClient → postponeClientRead 再次成功 → 又投递 第二次读任务。
  3. 第二次 I/O 线程再次进入 ioThreadReadQueryFromClient → 再次 parseInputBuffer(c) → 入口 serverAssert(c->cmd_queue.len == 0),此时 len 仍为上一次的 256 → 断言失败。 也就是说:不是「解析算法写错了」,而是 「在仍有未消费 cmd_queue 时,又安排了一次 I/O 线程解析」;高 pipeline 让第一次解析后 cmd_queue.len 长期不为 0,所以必现。 简单来说就是: c->cmd_queue.len == 0 断言失败,是因为 I/O 线程第一次 parseInputBuffer 已把大量 pipeline 命令放进 cmd_queue,而主线程在 processClientIOReadsDone 里过早执行 RDMA 的 connUpdateState,嵌套触发 readQueryFromClient → 再次 trySendReadToIOThreads,第二次 I/O 线程又调用 parseInputBuffer,此时队列非空,直接违反「解析前队列必须空」的不变量。