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

基于libwebsocket多线程接收WebSocket数据的优化方案咨询

问题描述

我需要基于libwebsocket创建两个WebSocket连接,要求能在不同线程接收数据。已知只需要创建一个lws_context,官方说明支持在n个线程中运行n个事件循环,但我找不到对应的示例代码。

目前我用独立线程处理数据解析来绕过这个问题,但这个方案并不理想——事件循环调度回调时会产生额外耗时,现在观测到两个连接的解析时间经常有0.2ms的明显差异(多数情况下接收的是同一条消息)。想请教有哪些优化方向或者更优的实现方案?

最小复现代码
#include <stdio.h>
#include <stdlib.h>
#include <libwebsockets.h>
#include <string.h>
#include <signal.h>
#include <pthread.h>

static struct lws *fclient_wsi;
static struct lws *client_wsi;
pthread_t* threads;
static int interrupted = 0;
static volatile bool ready1 = false;
static volatile bool ready2 = false;
static int counter = 0;

// Use this to get parse times
double get_time()
{
    LARGE_INTEGER t, f;
    QueryPerformanceCounter(&t);
    QueryPerformanceFrequency(&f);
    return (double)t.QuadPart/(double)f.QuadPart;
}
static void* AwaitParse(void* data){
    bool* ready = (bool *) data;
    while(!interrupted){
        if(!*ready)
            continue;
        double done = get_time();
        printf("Recieved at %f\n", done);

        *ready = false;
    }
}
static int ParseData1(struct lws *wsi, enum lws_callback_reasons reason,
                      void *user, void *in, size_t len)
{
    int res;
    switch (reason) {
        case LWS_CALLBACK_CLIENT_RECEIVE:
            if(!ready1){
                // Do something here
                counter++;
                if(counter > 10)
                    interrupted = 1;
                ready1 = true;
                lwsl_user("Data: %s\n", (const char*)in);
            }
            break;
        case LWS_CALLBACK_CLIENT_CONNECTION_ERROR:
            lwsl_err("CLIENT_CONNECTION_ERROR: %s\n",
                     in ? (char *)in : "(null)");
            client_wsi = NULL;
            break;

        case LWS_CALLBACK_CLIENT_ESTABLISHED:
            lwsl_user("%s: established\n", __func__);
            break;


        case LWS_CALLBACK_CLIENT_CLOSED:
            client_wsi = NULL;
            break;

        default:
            break;
    }

    return lws_callback_http_dummy(wsi, reason, user, in, len);
}
static int ParseData2(struct lws *wsi, enum lws_callback_reasons reason,
                      void *user, void *in, size_t len)
{
    int res;
    switch (reason) {
        case LWS_CALLBACK_CLIENT_RECEIVE:
            if(!ready2) {
                ready2 = true;
                lwsl_user("Data: %s\n", (const char*)in);
            }
            break;
        case LWS_CALLBACK_CLIENT_CONNECTION_ERROR:
            lwsl_err("CLIENT_CONNECTION_ERROR: %s\n",
                     in ? (char *)in : "(null)");
            client_wsi = NULL;
            break;

        case LWS_CALLBACK_CLIENT_ESTABLISHED:
            lwsl_user("%s: established\n", __func__);
            break;


        case LWS_CALLBACK_CLIENT_CLOSED:
            client_wsi = NULL;
            break;

        default:
            break;
    }

    return lws_callback_http_dummy(wsi, reason, user, in, len);
}
static const struct lws_extension extensions[] = {
        {
                "permessage-deflate",
                lws_extension_callback_pm_deflate,
                "permessage-deflate"
                "; client_no_context_takeover"
                "; client_max_window_bits"
        },
        { NULL, NULL, NULL /* terminator */ }
};

static const struct lws_protocols protocols[] = {
        {
                "data1-ws",
                ParseData1,
                      0,
                         0
        },
        {
                "data2-ws",
                ParseData2,
                      0,
                         0,
        },
        { NULL, NULL, 0, 0 }
};

static void sigint_handler(int sig)
{
    interrupted = 1;
}
int main() {
    threads = malloc(2 * sizeof(pthread_t));
    struct lws_context_creation_info info;
    struct lws_client_connect_info i;
    struct lws_context *context;
    int n = 0, logs = LLL_USER | LLL_WARN;

    signal(SIGINT, sigint_handler);

    lws_set_log_level(logs, NULL);
    memset(&info, 0, sizeof info); /* otherwise uninitialized garbage */
    info.options = LWS_SERVER_OPTION_DO_SSL_GLOBAL_INIT;
    info.port = CONTEXT_PORT_NO_LISTEN; /* we do not run any server */
    info.protocols = protocols;
    info.extensions = extensions;
    info.fd_limit_per_thread = 1 + 1 + 1 + 1;
    context = lws_create_context(&info);
    if (!context) {
        lwsl_err("Private WS init failed\n");
        return -1;
    }
    for(int j = 0; j < 2; j++){
        memset(&i, 0, sizeof i); /* otherwise uninitialized garbage */
        i.context = context;
        i.port = 443;
        i.address = "fstream.binance.com";
        i.path = "/ws/btcusdt@bookTicker";
        i.host = i.address;
        i.origin = i.address;
        i.ssl_connection = LCCSCF_USE_SSL;
        i.protocol = protocols[j].name;
        i.pwsi = (j == 0 ? &fclient_wsi : &client_wsi);
        lws_client_connect_via_info(&i);
    }

    pthread_create(&threads[0], NULL, AwaitParse, &ready1);
    pthread_create(&threads[1], NULL, AwaitParse, &ready2);

    while (n >= 0 && !interrupted)
        n = lws_service(context, 0);
    
    pthread_join(threads[0], NULL);
    pthread_join(threads[1], NULL);

    lws_context_destroy(context);
    lwsl_user("Completed Combined Streams\n");
    return 0;
}
优化方案与实现建议

1. 实现官方推荐的多线程事件循环

这是最贴合libwebsocket设计的方案,能从根源消除回调调度的耗时差异:

  • 核心逻辑:创建lws_context时指定线程服务索引(TSI)数量,每个线程独立运行对应的事件循环,连接绑定到指定TSI后,其所有事件(包括接收回调)直接在对应线程处理。
  • 关键代码修改:
    // 创建context时指定TSI数量为2
    info.tsi_count = 2;
    info.options |= LWS_SERVER_OPTION_EXPLICIT_VHOSTS;
    
    // 线程服务函数
    static void* thread_service(void* arg) {
        int tsi = *(int*)arg;
        while (!interrupted) {
            // 每个线程处理对应TSI的事件循环
            lws_service_tsi(context, 0, tsi);
        }
        return NULL;
    }
    
    // 创建连接时指定归属的TSI
    for(int j = 0; j < 2; j++){
        // ... 其他连接参数设置 ...
        i.tsi = j; // 将第j个连接绑定到第j个线程的事件循环
        lws_client_connect_via_info(&i);
    }
    
    // 启动两个事件循环线程
    int tsi0 = 0, tsi1 = 1;
    pthread_create(&threads[0], NULL, thread_service, &tsi0);
    pthread_create(&threads[1], NULL, thread_service, &tsi1);
    
  • 优势:无需额外线程同步,每个连接的事件处理完全在独立线程中完成,彻底消除回调调度的耗时差异。

2. 优化现有线程同步逻辑

如果暂时无法重构多线程事件循环,可先修复当前方案的低效同步问题:

  • 原方案用忙等待轮询ready变量,会占用CPU资源且导致响应延迟,替换为条件变量+互斥锁:
    // 定义同步变量
    pthread_mutex_t mutex1 = PTHREAD_MUTEX_INITIALIZER;
    pthread_cond_t cond1 = PTHREAD_COND_INITIALIZER;
    static volatile bool ready1 = false;
    
    // 回调中通知解析线程
    case LWS_CALLBACK_CLIENT_RECEIVE:
        if(!ready1){
            // ... 原有逻辑 ...
            pthread_mutex_lock(&mutex1);
            ready1 = true;
            pthread_cond_signal(&cond1);
            pthread_mutex_unlock(&mutex1);
        }
        break;
    
    // 解析线程逻辑
    static void* AwaitParse(void* data){
        bool* ready = (bool *) data;
        pthread_mutex_t* mutex = &mutex1; // 对应每个连接的锁
        pthread_cond_t* cond = &cond1;
        while(!interrupted){
            pthread_mutex_lock(mutex);
            while(!*ready && !interrupted){
                pthread_cond_wait(cond, mutex);
            }
            if(interrupted){
                pthread_mutex_unlock(mutex);
                break;
            }
            double done = get_time();
            printf("Received at %f\n", done);
            *ready = false;
            pthread_mutex_unlock(mutex);
        }
        return NULL;
    }
    
  • 效果:消除忙等待的CPU占用,让解析线程及时被唤醒,显著缩小解析时间差异。

3. 减少回调内的耗时操作

回调函数运行在事件循环线程中,耗时操作会阻塞整个事件处理流程:

  • 回调中仅做必要操作:比如拷贝接收到的数据,把解析、日志打印等耗时操作转移到独立线程。
  • 避免在回调中修改全局变量:改用线程局部存储(__thread关键字)或无锁数据结构,减少锁竞争。
  • 预分配缓冲区:使用内存池或固定大小的缓冲区,避免频繁内存分配的耗时。

4. 调整libwebsocket配置参数

  • 增大info.fd_limit_per_thread:原代码设置为4,可适当增大到8或16,避免文件描述符不足导致的调度延迟。
  • 启用第三方事件库:在创建context时添加LWS_SERVER_OPTION_LIBUV或LWS_SERVER_OPTION_LIBEV选项,替换默认的select/poll,提升事件处理效率。
  • 调整事件循环超时:将lws_service的第二个参数从0改为10(单位ms),平衡响应性和CPU占用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 08:37:06