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

nng是否支持单连接并行请求/响应?多线程RPC开发咨询

NNG在多线程双向RPC场景的实现方案

场景支持

NNG完全支持你需求的多线程并行收发类RPC消息、客户端与服务端双向发起请求的场景。你可以借助其套接字的双向通信能力,结合多上下文(nng_ctx)实现全双工RPC交互,无需严格区分传统的客户端/服务端角色,两端均可同时作为请求发起方和响应处理方。

单连接 vs 双连接

可以通过单连接实现双向RPC,利用NNG的全双工连接特性:

  • 若追求原生请求-响应语义,推荐用reqrep模式配合多上下文,单连接即可承载双向请求与响应的传输;
  • 若需更灵活的消息控制,也可使用pair模式,但需自行实现请求ID匹配逻辑来关联请求与响应。
    如果希望逻辑更清晰,也可让两端各自绑定rep套接字并连接对方的rep套接字,形成两个逻辑连接,但单连接方案更节省资源。

线程安全性

NNG针对多线程场景做了针对性设计:

  • 单个nng_socket对象并非完全线程安全,不能同时在多个线程中直接调用nng_send()或nng_recv();
  • 推荐使用**nng_ctx(上下文)**:每个上下文可独立绑定到套接字,不同线程可使用各自的上下文并行收发,上下文之间完全线程安全,这是多线程场景下的标准用法;
  • 连接管理类操作(如nng_listen()、nng_dial())本身是线程安全的。

示例代码

以下是简化的双向RPC示例,两端均可发起请求并处理响应,基于reqrep模式和多上下文实现多线程并发:

通用工具函数

#include <nng/nng.h>
#include <nng/protocol/reqrep0/req.h>
#include <nng/protocol/reqrep0/rep.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <pthread.h>
#include <unistd.h>

// 响应请求的线程逻辑
void *rep_thread(void *arg) {
    nng_socket *sock = (nng_socket *)arg;
    nng_ctx ctx;
    int rv;

    if ((rv = nng_ctx_open(&ctx, *sock)) != 0) {
        fprintf(stderr, "nng_ctx_open failed: %s\n", nng_strerror(rv));
        return NULL;
    }

    while (1) {
        char *buf = NULL;
        size_t sz;
        if ((rv = nng_ctx_recv(ctx, &buf, &sz, 0)) != 0) {
            fprintf(stderr, "nng_ctx_recv failed: %s\n", nng_strerror(rv));
            break;
        }
        printf("Received request: %.*s\n", (int)sz, buf);

        // 构造并发送响应
        char resp[128];
        snprintf(resp, sizeof(resp), "Response to: %.*s", (int)sz, buf);
        if ((rv = nng_ctx_send(ctx, resp, strlen(resp)+1, 0)) != 0) {
            fprintf(stderr, "nng_ctx_send failed: %s\n", nng_strerror(rv));
        }
        nng_free(buf, sz);
    }
    nng_ctx_close(ctx);
    return NULL;
}

// 发起请求的线程逻辑
void *req_thread(void *arg) {
    nng_socket *sock = (nng_socket *)arg;
    nng_ctx ctx;
    int rv;

    if ((rv = nng_ctx_open(&ctx, *sock)) != 0) {
        fprintf(stderr, "nng_ctx_open failed: %s\n", nng_strerror(rv));
        return NULL;
    }

    int count = 0;
    while (count < 5) {
        char req[128];
        snprintf(req, sizeof(req), "Request #%d from thread", ++count);
        printf("Sending request: %s\n", req);

        if ((rv = nng_ctx_send(ctx, req, strlen(req)+1, 0)) != 0) {
            fprintf(stderr, "nng_ctx_send failed: %s\n", nng_strerror(rv));
            sleep(1);
            continue;
        }

        char *resp = NULL;
        size_t sz;
        if ((rv = nng_ctx_recv(ctx, &resp, &sz, 0)) != 0) {
            fprintf(stderr, "nng_ctx_recv failed: %s\n", nng_strerror(rv));
        } else {
            printf("Received response: %.*s\n", (int)sz, resp);
            nng_free(resp, sz);
        }
        sleep(1);
    }
    nng_ctx_close(ctx);
    return NULL;
}

节点A(双向角色)

int main() {
    nng_socket rep_sock, req_sock;
    pthread_t rep_tid, req_tid;
    int rv;

    // 创建响应套接字并监听端口
    if ((rv = nng_rep0_open(&rep_sock)) != 0) {
        fprintf(stderr, "nng_rep0_open failed: %s\n", nng_strerror(rv));
        return 1;
    }
    if ((rv = nng_listen(rep_sock, "tcp://0.0.0.0:5555", NULL, 0)) != 0) {
        fprintf(stderr, "nng_listen failed: %s\n", nng_strerror(rv));
        return 1;
    }

    // 创建请求套接字并连接节点B
    if ((rv = nng_req0_open(&req_sock)) != 0) {
        fprintf(stderr, "nng_req0_open failed: %s\n", nng_strerror(rv));
        return 1;
    }
    if ((rv = nng_dial(req_sock, "tcp://127.0.0.1:5556", NULL, 0)) != 0) {
        fprintf(stderr, "nng_dial failed: %s\n", nng_strerror(rv));
        return 1;
    }

    // 启动响应和请求线程
    pthread_create(&rep_tid, NULL, rep_thread, &rep_sock);
    pthread_create(&req_tid, NULL, req_thread, &req_sock);

    pthread_join(rep_tid, NULL);
    pthread_join(req_tid, NULL);

    nng_close(rep_sock);
    nng_close(req_sock);
    return 0;
}

节点B(双向角色)

int main() {
    nng_socket rep_sock, req_sock;
    pthread_t rep_tid, req_tid;
    int rv;

    // 创建响应套接字并监听端口
    if ((rv = nng_rep0_open(&rep_sock)) != 0) {
        fprintf(stderr, "nng_rep0_open failed: %s\n", nng_strerror(rv));
        return 1;
    }
    if ((rv = nng_listen(rep_sock, "tcp://0.0.0.0:5556", NULL, 0)) != 0) {
        fprintf(stderr, "nng_listen failed: %s\n", nng_strerror(rv));
        return 1;
    }

    // 创建请求套接字并连接节点A
    if ((rv = nng_req0_open(&req_sock)) != 0) {
        fprintf(stderr, "nng_req0_open failed: %s\n", nng_strerror(rv));
        return 1;
    }
    if ((rv = nng_dial(req_sock, "tcp://127.0.0.1:5555", NULL, 0)) != 0) {
        fprintf(stderr, "nng_dial failed: %s\n", nng_strerror(rv));
        return 1;
    }

    // 启动响应和请求线程
    pthread_create(&rep_tid, NULL, rep_thread, &rep_sock);
    pthread_create(&req_tid, NULL, req_thread, &req_sock);

    pthread_join(rep_tid, NULL);
    pthread_join(req_tid, NULL);

    nng_close(rep_sock);
    nng_close(req_sock);
    return 0;
}

该示例中,两个节点各自启动响应线程处理对方请求,同时启动请求线程发起请求,实现了双向RPC的多线程并行交互。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 15:25:23