dcf写入机制


难度 中等

写入

dcf提供如下两个写入接口:

  • dcf_write

    int dcf_write(unsigned int stream_id, const char* buffer, unsigned int length, unsigned long long key, unsigned long long *index);

    仅在leader节点调用。

  • dcf_universal_write

    int dcf_universal_write(unsigned int stream_id, const char* buffer, unsigned int length, unsigned long long key, unsigned long long *index);

    可以在任意节点调用。

确认

  • dcf_register_after_writer

    用于注册leader节点写入成功回调函数。

    int dcf_register_after_writer(usr_cb_after_writer_t cb_func)
    {
        return rep_register_after_writer(ENTRY_TYPE_LOG, cb_func);
    }

    最终将回调函数注册到全局变量

    int rep_register_after_writer(entry_type_t type, usr_cb_after_writer_t cb_func)
    {
        g_cb_after_writer[type] = cb_func;
        return CM_SUCCESS;
    }

    调用流程如下:

graph TB rep_common_init-->cm_create_thread--> |启动个线程来执行apply|rep_apply_thread_entry-->rep_apply_proc -->g_cb_after_writer

线程在等待条件变量释放,并执行apply。

// 先查出有多少条流
if (md_get_stream_list(streams, &stream_count) != CM_SUCCESS) {
        LOG_DEBUG_ERR("[REP]md_get_stream_list failed");
        return;
    }
while (!thread->closed) { // 若线程没有关闭则循环等待
        if (!exists_log) { // 首先判断是否存在日志,不存在就休眠等待唤醒
            (void)cm_event_timedwait(&g_apply_cond, CM_SLEEP_500_FIXED);
        }
        LOG_TRACE(g_rep_tracekey, "apply_thread work");
        exists_log = CM_FALSE;
        for (uint32 i = 0; i < stream_count; i++) { // 遍历每一条流
            uint32 stream_id = streams[i];
            bool8 stream_exists_log = CM_FALSE;
            LOG_TIME_BEGIN(rep_apply_proc);
            // 执行apply
            if (rep_apply_proc(stream_id, &stream_exists_log) != CM_SUCCESS) {
                LOG_DEBUG_ERR_EX("[REP]rep_apply_proc failed.");
            }
            LOG_TIME_END(rep_apply_proc);

            exists_log = (exists_log || stream_exists_log);
        }
    }

其中等待的条件变量为,在如下的地方唤醒。

void rep_apply_trigger()
{
  LOG_DEBUG_INF("[REP]rep_apply_trigger");
  LOG_TRACE(g_rep_tracekey, "common:rep_apply_trigger.");
  cm_event_notify(&g_apply_cond);
}

rep_apply_trigger的调用栈如下;

graph TB rep_acceptlog_proc-->rep_leader_acceptlog_proc--> rep_try_commit_log-->|leader|rep_apply_trigger rep_accept_thread_entry-->rep_acceptlog_proc--> rep_follower_acceptlog_proc-->|follower|rep_apply_trigger

经过上面调用流程可以看到,apply线程是由accept线程唤醒的。accept线程与apply线程类似,同样是等待条件变量将线程唤醒。

if (md_get_stream_list(streams, &stream_count) != CM_SUCCESS) {
        LOG_DEBUG_ERR("[REP]md_get_stream_list failed");
        return;
    }

    while (!thread->closed) {
        if (!exists_log) {
            LOG_TRACE(g_rep_tracekey, "accept_thread wait.");
            (void)cm_event_timedwait(&g_accept_cond, CM_SLEEP_500_FIXED);
        }
        LOG_TRACE(g_rep_tracekey, "accept_thread work.");

        exists_log = CM_FALSE;
        for (uint32 i = 0; i < stream_count; i++) {
            uint32 stream_id = streams[i];
            date_t now = g_timer()->now;
            exists_log = (exists_log || g_common_state[stream_id].accept_log);
            if (g_common_state[stream_id].accept_log ||
                now - g_common_state[stream_id].last_accept_time > CM_DEFAULT_HB_INTERVAL*MICROSECS_PER_MILLISEC) {
                LOG_TRACE(g_rep_tracekey, "accept_thread do work.");
                g_common_state[stream_id].accept_log = CM_FALSE;
                g_common_state[stream_id].last_accept_time = now;
                if (rep_acceptlog_proc(stream_id) != CM_SUCCESS) {
                    LOG_DEBUG_ERR("[REP]rep_acceptlog_proc failed.");
                }
            } else {
                LOG_TRACE(g_rep_tracekey, "accept_thread no work.");
            }
        }
    }

accept线程等待的条件变量为g_accept_cond,其唤醒流程如下:

void rep_set_accept_flag(uint32 stream_id)
{
    LOG_DEBUG_INF("rep_set_accept_flag.");
    g_common_state[stream_id].accept_log = CM_TRUE;
    cm_event_notify(&g_accept_cond);
}

rep_set_accept_flag调用堆栈如下:

graph TB stg_register_cb-->rep_accepted_trigger

可以看已看到accept线程的唤醒是通过注册回调函数来触发的。

    if (stg_register_cb(ENTRY_TYPE_LOG, rep_accepted_trigger) != CM_SUCCESS) {
    LOG_DEBUG_ERR("[REP]rep register stg callback failed");
    return CM_ERROR;
}
status_t stg_register_cb(entry_type_t type, void *func)
{
    switch (type) {
        case ENTRY_TYPE_CONF:
            g_write_conf_func = (write_conf_func_t)func;
            break;
        case ENTRY_TYPE_LOG:
            g_notify_rep_func = (notify_rep_func_t)func;
            break;
        default:
            LOG_RUN_ERR("[STG]Register callback failed");
            return CM_ERROR;
    }
    return CM_SUCCESS;
}

可以看到回调函数最终被注册到了g_notify_rep_func。

graph TB disk_thread_entry-->|append线程|process_append_action-->stream_append_entry_impl -->callback_rep_func-->g_notify_rep_func stream_batcher_flush-->callback_rep_func

可以看到其中一个回调函数是由append线程触发的,append线程会等待stream->disk_event条件变量被唤醒。stream->disk_event在stream_append_entry中被唤醒。

graph TB rep_appendlog_req_proc-->rep_follower_process -->rep_follower_appendlog-->stg_append_entry -->stream_append_entry rep_write-->stg_append_entry
register_msg_process(MEC_CMD_APPEND_LOG_RPC_REQ, rep_appendlog_req_proc, PRIV_LOW);

MEC_CMD_APPEND_LOG_RPC_REQ在rep_appendlog_node中发送。

graph TB rep_appendlog_thread_entry-->rep_appendlog_stream-->rep_appendlog_node

rep_appendlog_thread_entry有条件变量g_appendlog_cond唤醒。g_appendlog_cond又通过rep_appendlog_trigger唤醒。

graph TB rep_write-->rep_appendlog_trigger rep_appendlog_ack_proc-->rep_appendlog_trigger rep_rematch_proc-->rep_appendlog_trigger

rep_appendlog_trigger由MEC_CMD_APPEND_LOG_RPC_ACK消息触发。

register_msg_process(MEC_CMD_APPEND_LOG_RPC_ACK, rep_appendlog_ack_proc, PRIV_LOW);

dd

  • dcf_register_consensus_notify

    用于注册follower节点写入数据成功的回调函数。


文章作者: growdu
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 growdu !
  目录
分类导航
随笔2 AI27 算法1 计算机基础13 博客搭建7 ChatGPT2 集群63 计算机通信1 数据库34 数据库深入80 DPDK26 Docker11 Elasticsearch4 编辑工具4 FAQ1 Go Web1 hometown2 编程语言16 网络9 OPC1 Linux38 openGauss4 页面12 PostgreSQL54 程序员自我修养1 协议11 成长之路1 stock1 存储5 工具20 VPP18 视频作品1 Vue13 Web1 代码示例11 数据库15 BenchmarkSQL1 PostgreSQL 源码修炼之路14
最热文章
1
13 逻辑复制深入
数据库深入🔥 1570
2
0 Postgresql存储、索引及系统优化、主备切换
PostgreSQL🔥 1495
3
一文读懂openguass dcf网络模块
集群🔥 1420
4
逻辑复制源码分析
数据库深入🔥 1327
5
PostgreSQL 分区表:从一行 `PARTITION BY` 到路由热路径的全链路拆解
数据库🔥 1094
6
applyparallelworker.c 之 LA 端源码深度解析:Leader Apply Worker 的指挥中枢
数据库深入🔥 1082
7
PostgreSQL Background Worker 全解:从 `RegisterBackgroundWorker` 到逻辑复制 4 类 worker 的全生命周期
数据库🔥 1078
8
PostgreSQL的后台进程walsender分析 - 关系型数据库 - 亿速云
PostgreSQL🔥 1033
9
PostgreSQL 逻辑复制的监控:六张视图 + 一组可执行 SQL,把 publisher/subscriber 的速率与健康度彻底看透
数据库🔥 1032
10
PostgreSQL 逻辑复制支持 DDL 之后:DDL 与 DML 的时序难题(重点:分区表)
数据库🔥 999
11
reorderbuffer.c 源码深度解析:PostgreSQL 逻辑复制的"事务重组引擎
数据库深入🔥 953
12
PostgreSQL 内核开发:读取一张表的 9 步标准流程与缓存全景
数据库🔥 938
13
从 `postgres` 二进制到生产级守护 —— PostgreSQL 最外层模块与启动全流程拆解
数据库🔥 936
14
支持逻辑复制同步 DDL 适配 SQL Server 方案
数据库深入🔥 934
15
PostgreSQL 逻辑复制的 ReorderBuffer 与事务机制:从一行 WAL 到一致性变更流的全链路绑定
数据库🔥 913
16
DDL同步架构(美化版)
数据库深入🔥 908
17
PostgreSQL Latch 机制详解:从一行 SetLatch 到 epoll 的内核之旅
数据库🔥 871
18
pgbench 源码全解:一个 C 文件如何撑起 PostgreSQL 官方压测工具
数据库🔥 860
19
PostgreSQL libpq 机制与缓冲区详解
数据库🔥 850
20
PostgreSQL 逻辑复制 spill 文件深度剖析:从 `xid-*.spill` 到 TPC-C 的增长方程
数据库🔥 845