opengauss项目代码注释 #26

Open
nuoya wants to merge 55 commits from nuoya/openGauss-server:master into master
1 changed files with 111 additions and 52 deletions
Showing only changes of commit 844c260a3b - Show all commits

View File

@ -15,10 +15,9 @@
*
* threadpool_listener.cpp
*
* There are multiple tasks for listener thread:
* 1. Listen to all connections from client or other componets of this cluster
* (like connections from other cn).
* 2. Dispatch session to available woker thread.
* 线
* 1.
* 2. 线
*
* IDENTIFICATION
* src/gausskernel/process/threadpool/threadpool_listener.cpp
@ -74,7 +73,7 @@ void TpoolListenerMain(ThreadPoolListener* listener)
(void)gspqsignal(SIGHUP, SIG_IGN);
(void)gspqsignal(SIGINT, SIG_IGN);
// die with pm
// Postmaster进程退出时终止
(void)gspqsignal(SIGTERM, SIG_IGN);
(void)gspqsignal(SIGQUIT, SIG_IGN);
(void)gspqsignal(SIGPIPE, SIG_IGN);
@ -107,31 +106,49 @@ void ThreadPoolListenerIAm()
ThreadPoolListener::ThreadPoolListener(ThreadPoolGroup* group)
{
// 将传入的线程池组对象保存到成员变量中
m_group = group;
// 初始化一些成员变量,设置为初始值
m_tid = InvalidTid;
m_epollFd = INVALID_FD;
m_epollEvents = NULL;
m_reaperAllSession = false;
m_getKilled = false;
// 创建用于管理空闲工作线程的链表
m_freeWorkerList = New(CurrentMemoryContext) DllistWithLock();
// 创建用于管理准备就绪的会话的链表
m_readySessionList = New(CurrentMemoryContext) DllistWithLock();
// 创建用于管理空闲会话的链表
m_idleSessionList = New(CurrentMemoryContext) DllistWithLock();
// 检查是否启用本地系统缓存,根据情况进行初始化
if (EnableLocalSysCache()) {
/* see HASH_INDEX, Since the hash table must contain a power-of-2 number of elements */
// 根据配置设置会话哈希表的桶数通常为2的幂
#ifdef ENABLE_LITE_MODE
m_session_nbucket = 128;
#else
m_session_nbucket = MAX_THREAD_POOL_SIZE;
#endif
// 分配会话桶的内存空间
m_session_bucket = (Dllist*)palloc0(m_session_nbucket * sizeof(Dllist));
// 分配会话读写锁的内存空间
m_session_rw_locks = (pthread_rwlock_t *)palloc0(m_session_nbucket * sizeof(pthread_rwlock_t));
// 初始化每个会话桶的读写锁
for (int i = 0; i < m_session_nbucket; i++) {
PthreadRwLockInit(&m_session_rw_locks[i], NULL);
}
// 初始化用于匹配搜索的标志
m_match_search = 0;
} else {
// 如果未启用本地系统缓存,则将相关成员变量设置为初始值
m_session_nbucket = 0;
m_session_bucket = NULL;
m_session_rw_locks = NULL;
@ -141,21 +158,27 @@ ThreadPoolListener::ThreadPoolListener(ThreadPoolGroup* group)
ThreadPoolListener::~ThreadPoolListener()
{
// 根据宏定义选择不同的关闭方法
if (ENABLE_THREAD_POOL_DN_LOGICCONN) {
CommEpollClose(m_epollFd);
} else {
comm_close(m_epollFd); /* CommProxy support */
comm_close(m_epollFd); /* CommProxy 支持 */
}
// 将相关成员变量设置为 NULL以释放对资源的引用
m_group = NULL;
m_epollEvents = NULL;
m_freeWorkerList = NULL;
m_readySessionList = NULL;
m_idleSessionList = NULL;
// 如果启用了本地系统缓存,释放相应的资源
if (EnableLocalSysCache()) {
pfree_ext(m_session_bucket);
pfree_ext(m_session_rw_locks);
}
// 将与本地系统缓存相关的成员变量设置为初始值
m_session_nbucket = 0;
m_session_bucket = NULL;
}
@ -173,27 +196,32 @@ void ThreadPoolListener::NotifyReady()
void ThreadPoolListener::CreateEpoll()
{
/* MAX_LISTEN_SESSIONS for epool_create is ignored Since Linux 2.6.8 */
// 根据宏定义选择不同的 epoll 创建函数
if (ENABLE_THREAD_POOL_DN_LOGICCONN) {
// 使用 CommEpollCreate 函数创建 epoll传入最大会话数
m_epollFd = CommEpollCreate(GLOBAL_MAX_SESSION_NUM);
} else {
// 使用 comm_epoll_create 函数创建 epoll传入最大会话数
CommSetEpollOption(CommEpollThreadPoolListener);
m_epollFd = comm_epoll_create(GLOBAL_MAX_SESSION_NUM);
}
// 检查 epoll 创建是否成功
if (m_epollFd == INVALID_FD) {
ereport(LOG,
(errmsg("Fail to create epoll for thread pool listener, "
"check if the system is out of memory or "
"limit on the total number of open files has been reached.")));
proc_exit(0);
proc_exit(0); // 退出当前进程
}
// 分配用于存储 epoll 事件的内存
m_epollEvents = (struct epoll_event*)palloc0_noexcept(sizeof(struct epoll_event) * GLOBAL_MAX_SESSION_NUM);
// 检查内存分配是否成功
if (m_epollEvents == NULL) {
elog(LOG, "Not enough memory for listener epoll");
proc_exit(0);
proc_exit(0); // 退出当前进程
}
}
@ -207,20 +235,19 @@ void ThreadPoolListener::AddEpoll(knl_session_context* session)
(errmodule(MOD_THREAD_POOL),
errmsg("Add a session:%lu to idleSessionList ", session->session_id)));
/*
* Because we will dispatch the socket to worker thread once
* we find an input event of the socket, so we use one_shot mode.
* 线使 "one_shot"
* 线线
* 线
*/
ev.events = EPOLLRDHUP | EPOLLIN | EPOLLET | EPOLLONESHOT;
ev.data.ptr = (void*)session;
if (session->status != KNL_SESS_UNINIT) {
/* CommProxy Support */
if (ENABLE_THREAD_POOL_DN_LOGICCONN) {
res = CommEpollCtl(m_epollFd, EPOLL_CTL_MOD, session->proc_cxt.MyProcPort->sock, &ev);
} else {
res = comm_epoll_ctl(m_epollFd, EPOLL_CTL_MOD, session->proc_cxt.MyProcPort->sock, &ev);
}
} else {
/* CommProxy Support */
if (ENABLE_THREAD_POOL_DN_LOGICCONN) {
res = CommEpollCtl(m_epollFd, EPOLL_CTL_ADD, session->proc_cxt.MyProcPort->sock, &ev);
} else {
@ -244,10 +271,14 @@ void ThreadPoolListener::AddEpoll(knl_session_context* session)
bool ThreadPoolListener::TryFeedWorker(ThreadPoolWorker* worker)
{
Dlelem* sc = GetReadySession(worker);
//获取一个准备就绪的会话
if (sc != NULL) {
worker->SetSession((knl_session_context*)sc->dle_val);
//将该会话设置给工作线程
pg_atomic_fetch_sub_u32((volatile uint32*)&m_group->m_waitServeSessionCount, 1);
//减少等待服务会话的计数,表示该会话已被分配
pg_atomic_fetch_add_u32((volatile uint32*)&m_group->m_processTaskCount, 1);
//增加正在处理任务的计数,表示工作线程正在处理任务。
return true;
} else {
if (EnableLocalSysCache()) {
@ -271,9 +302,10 @@ bool ThreadPoolListener::TryFeedWorker(ThreadPoolWorker* worker)
void ThreadPoolListener::AddNewSession(knl_session_context* session)
{
AddEpoll(session);
AddEpoll(session);//将指定的会话 session 添加到 epoll 监听中
(void)pg_atomic_fetch_add_u32((volatile uint32*)&m_group->m_sessionCount, 1);
ereport(DEBUG2,
//使用原子操作增加线程池组的会话计数,表示成功添加了一个新的会话
ereport(DEBUG2,
(errmodule(MOD_THREAD_POOL),
errmsg("This group add a session, and now sessionCount is %d, ",
m_group->m_sessionCount)));
@ -298,9 +330,8 @@ void ThreadPoolListener::ReaperAllSession()
while (m_group->m_sessionCount > 0) {
/*
* There is a very rare case that all thread pool workers happen to
* encounter FATAL and exit before close session.
* Under such scenarios, we choose to exit directly.
* 线FATAL退
* 退
*/
if (m_group->m_workerNum <= 0 && m_group->m_sessionCount > 0) {
ereport(WARNING,
@ -309,8 +340,7 @@ void ThreadPoolListener::ReaperAllSession()
" encounter FATAL problems before session close.")));
abort();
}
/* m_sessionCount should be sum of the list length of m_idleSessionList and m_readySessionList
and worker's attached session */
/* m_sessionCount 应该是 m_idleSessionList 和 m_readySessionList 的列表长度之和,以及与工作线程关联的会话的数量之和 */
pg_memory_barrier();
if (m_idleSessionList->IsEmpty() && m_readySessionList->IsEmpty() &&
m_group->m_workerNum - m_group->m_idleWorkerNum == 0) {
@ -348,10 +378,10 @@ void ThreadPoolListener::WaitTask()
while (true) {
if (unlikely(m_getKilled)) {
m_getKilled = false;
proc_exit(0);
proc_exit(0); // 如果需要退出,直接退出进程
}
if (unlikely(m_reaperAllSession)) {
ReaperAllSession();
ReaperAllSession(); // 如果需要清理所有会话,执行清理操作
}
/* as we specify timeout -1, so 0 will not be return, either > 0 or < 0 */
@ -360,6 +390,8 @@ void ThreadPoolListener::WaitTask()
} else {
nevents = comm_epoll_wait(m_epollFd, m_epollEvents, GLOBAL_MAX_SESSION_NUM, -1); /* CommProxy Support */
}
// 处理收到的事件
if (nevents > 0 && nevents <= GLOBAL_MAX_SESSION_NUM) {
HandleConnEvent(nevents);
continue;
@ -367,9 +399,9 @@ void ThreadPoolListener::WaitTask()
ereport(PANIC,
(errmsg("epoll receive %d events which exceed the limitation %d", nevents, GLOBAL_MAX_SESSION_NUM)));
} else if (nevents == -1 && errno == EINTR) {
continue;
continue; // 如果是中断信号,继续等待事件
} else {
ereport(LOG, (errmsg("listener wait event encounter some error :%d", errno)));
ereport(LOG, (errmsg("listener wait event encounter some error :%d", errno))); // 其他错误情况下输出日志
}
}
}
@ -384,10 +416,10 @@ void ThreadPoolListener::HandleConnEvent(int nevets)
session = GetSessionBaseOnEvent(tmp_event);
if (session == NULL) {
continue;
continue; // 如果没有获取到会话,继续处理下一个事件
}
DispatchSession(session);
DispatchSession(session); // 处理会话的分派
}
}
@ -402,20 +434,18 @@ knl_session_context* ThreadPoolListener::GetSessionBaseOnEvent(struct epoll_even
} else {
session->status = KNL_SESS_CLOSE;
}
return session;
return session; // 如果发生错误或会话关闭事件,则返回该会话
} else if (ev->events & EPOLLIN) {
return session;
return session; // 如果是输入事件,返回该会话
}
return NULL;
return NULL; // 其他情况下返回 NULL表示没有需要处理的会话
}
void ThreadPoolListener::DispatchSession(knl_session_context* session)
{
m_idleSessionList->Remove(&session->elem);
/*
* If the sock, idx, and streamid parameters of the current session
* do not meet the requirements for logical connection parameters,
* skip this dispatch operation.
* sockidx streamid
*/
if (session->proc_cxt.MyProcPort->sock == NO_SOCKET &&
session->proc_cxt.MyProcPort->gs_sock.idx == 0 &&
@ -443,7 +473,7 @@ void ThreadPoolListener::DispatchSession(knl_session_context* session)
__func__, session->session_id)));
INSTR_TIME_SET_CURRENT(session->last_access_time);
/* Add new session to the head so the connection request can be quickly processed. */
/* 将新会话添加到头部,以便可以快速处理连接请求 */
if (session->status == KNL_SESS_UNINIT) {
AddIdleSessionToHead(session);
} else {
@ -457,10 +487,15 @@ void ThreadPoolListener::DispatchSession(knl_session_context* session)
void ThreadPoolListener::DelSessionFromEpoll(knl_session_context* session)
{
//检查是否启用了线程池的逻辑连接支持
if (ENABLE_THREAD_POOL_DN_LOGICCONN) {
struct epoll_event ev = {0};
//创建一个名为 ev 的 epoll 事件结构,并初始化所有字段为零
struct epoll_event ev = {0};
//设置 epoll 事件的关注事件
ev.events = EPOLLRDHUP | EPOLLIN | EPOLLET | EPOLLONESHOT;
ev.data.ptr = (void*)session;
//将 ev 事件的数据指针设置为指向当前会话 (session) 的指针
ev.data.ptr = (void*)session;
//使用 CommEpollCtl 函数从 epoll 中删除套接字事件
CommEpollCtl(m_epollFd, EPOLL_CTL_DEL, session->proc_cxt.MyProcPort->sock, &ev);
#ifdef ENABLE_MULTIPLE_NODES
} else {
@ -469,6 +504,7 @@ void ThreadPoolListener::DelSessionFromEpoll(knl_session_context* session)
#endif
comm_epoll_ctl(m_epollFd, EPOLL_CTL_DEL, session->proc_cxt.MyProcPort->sock, NULL);
}
//使用原子操作将会话计数减少 1
(void)pg_atomic_fetch_sub_u32((volatile uint32*)&m_group->m_sessionCount, 1);
}
@ -479,37 +515,52 @@ void ThreadPoolListener::RemoveWorkerFromList(ThreadPoolWorker* worker)
bool ThreadPoolListener::GetSessIshang(instr_time* current_time, uint64* sessionId)
{
// 初始化 ishang 为 true
bool ishang = true;
// 获取就绪会话列表的锁
m_readySessionList->GetLock();
// 获取就绪会话列表的头元素
Dlelem* elem = m_readySessionList->GetHead();
// 如果头元素为空,释放锁并返回 false
if (elem == NULL) {
m_readySessionList->ReleaseLock();
return false;
}
// 将头元素转换为 knl_session_context 指针
knl_session_context* head_sess = (knl_session_context *)(elem->dle_val);
// 检查就绪会话的时间戳和会话ID是否与传入的值匹配
if (INSTR_TIME_GET_MICROSEC(head_sess->last_access_time) == INSTR_TIME_GET_MICROSEC(*current_time) &&
head_sess->session_id == *sessionId) {
ishang = true;
ishang = true; // 如果匹配ishang 保持 true
} else {
// 更新传入的时间戳和会话ID
*current_time = head_sess->last_access_time;
*sessionId = head_sess->session_id;
ishang = false;
ishang = false; // ishang 设为 false 表示不再挂起状态
}
// 释放就绪会话列表的锁
m_readySessionList->ReleaseLock();
// 返回 ishang指示监听器是否挂起
return ishang;
}
Dlelem *ThreadPoolListener::GetFreeWorker(knl_session_context* session)
{
/* only lite mode need find right threadworker,
* otherwise since there are so many requests, we dont have any freeworkers. so optimization is not necessary */
/* 只有在 "轻量级模式" 下才需要找到正确的线程工作线程,否则由于有如此
线*/
#ifdef ENABLE_LITE_MODE
if (!EnableLocalSysCache()) {
return m_freeWorkerList->RemoveHead();
}
/* sess is not init, we dont know how to hit the cache */
/* 如果会话未初始化,我们不知道如何命中缓存 */
if (session->status != KNL_SESS_ATTACH && session->status != KNL_SESS_DETACH) {
return m_freeWorkerList->RemoveTail();
}
@ -518,16 +569,16 @@ Dlelem *ThreadPoolListener::GetFreeWorker(knl_session_context* session)
return m_freeWorkerList->RemoveTail();
}
/* for lite_mode, threadworkers are a small amount, so it is quickly to traverse the list */
/* 对于轻量级模式,线程工作者数量较少,因此迅速遍历列表是可行的。 */
m_freeWorkerList->GetLock();
for (Dlelem *elt = m_freeWorkerList->GetHead(); elt != NULL; elt = DLGetSucc(elt)) {
ThreadPoolWorker *worker = (ThreadPoolWorker *)DLE_VAL(elt);
LocalSysDBCache *lsc = worker->GetThreadContextPtr()->lsc_cxt.lsc;
/* uninited lsc are addtotail of the list, so when see one uninited, the follow all are uninited. just break */
/* 未初始化的本地系统缓存lsc被添加到列表的末尾因此当遇到一个未初始化的时候后续的所有也都是未初始化的。因此可以直接中断break */
if (unlikely(lsc == NULL || lsc->my_database_id == InvalidOid)) {
break;
}
/* cache hit */
/* 缓存命中 */
if (likely(lsc->my_database_id == session->proc_cxt.MyDatabaseId)) {
m_freeWorkerList->Remove(elt);
m_freeWorkerList->ReleaseLock();
@ -535,7 +586,8 @@ Dlelem *ThreadPoolListener::GetFreeWorker(knl_session_context* session)
}
}
m_freeWorkerList->ReleaseLock();
/* dont find, use tail instead head, because head of the list has syscache of other db */
/* 建议不要从链表的头部查找,而是从尾部开始查找,因为链表的头部可能包含了
syscache */
return m_freeWorkerList->RemoveTail();
#else
return m_freeWorkerList->RemoveHead();
@ -544,12 +596,12 @@ Dlelem *ThreadPoolListener::GetFreeWorker(knl_session_context* session)
static Dlelem *GetHeadUnInitSession(DllistWithLock* m_readySessionList)
{
/* uninit session needs be replied first */
/* 未初始化的会话应该首先得到回复 */
m_readySessionList->GetLock();
Dlelem *head = m_readySessionList->GetHead();
if (likely(head != NULL)) {
if (((knl_session_context *)DLE_VAL(head))->status != KNL_SESS_UNINIT) {
/* go cache hit branch, set it null */
/* 程序在执行过程中进入了“缓存命中分支”,并且将某个值设置为了 null */
head = NULL;
} else {
head = m_readySessionList->RemoveHeadNoLock();
@ -574,11 +626,11 @@ Dlelem *ThreadPoolListener::GetSessFromReadySessionList(ThreadPoolWorker *worker
break;
}
LocalSysDBCache *lsc = worker->GetThreadContextPtr()->lsc_cxt.lsc;
// worker not init, any session is matched
// 如果工作线程尚未初始化,那么任何会话都不会被匹配或关联
if (unlikely(lsc == NULL || lsc->my_database_id == InvalidOid)) {
break;
}
// now we try to reuse workers syscache
// 现在我们尝试重用工作线程的系统缓存
Index hash_index = HASH_INDEX(lsc->my_database_id, (uint32)m_session_nbucket);
ResourceOwner owner = LOCAL_SYSDB_RESOWNER;
PthreadRWlockRdlock(owner, &m_session_rw_locks[hash_index]);
@ -588,7 +640,7 @@ Dlelem *ThreadPoolListener::GetSessFromReadySessionList(ThreadPoolWorker *worker
break;
}
if (!m_readySessionList->RemoveConfirm(&((knl_session_context *)DLE_VAL(elt))->elem)) {
// someone remove it already
// 已经将他移除了
PthreadRWlockUnlock(owner, &m_session_rw_locks[hash_index]);
break;
}
@ -606,14 +658,17 @@ Dlelem *ThreadPoolListener::GetReadySession(ThreadPoolWorker *worker)
if (!EnableLocalSysCache()) {
return m_readySessionList->RemoveHead();
}
// 如果不启用本地系统缓存,直接从就绪会话列表的头部移除并返回一个会话。
Dlelem *elt = GetSessFromReadySessionList(worker);
if (elt == NULL) {
return NULL;
}
// 从本地系统缓存中获取一个会话。
knl_session_context *session = (knl_session_context *)DLE_VAL(elt);
Oid cur_dbid = session->proc_cxt.MyDatabaseId;
Index hash_index = HASH_INDEX(cur_dbid, (uint32)m_session_nbucket);
ResourceOwner owner = LOCAL_SYSDB_RESOWNER;
// 获取本地系统缓存的资源锁,并从就绪会话列表中移除该会话。
PthreadRWlockWrlock(owner, &m_session_rw_locks[hash_index]);
DLRemove(&session->elem2);
PthreadRWlockUnlock(owner, &m_session_rw_locks[hash_index]);
@ -626,8 +681,10 @@ void ThreadPoolListener::AddIdleSessionToTail(knl_session_context* session)
m_readySessionList->AddTail(&session->elem);
return;
}
// 如果不启用本地系统缓存,将会话添加到就绪会话列表的尾部。
Index hash_index = HASH_INDEX(session->proc_cxt.MyDatabaseId, (uint32)m_session_nbucket);
ResourceOwner owner = LOCAL_SYSDB_RESOWNER;
// 获取本地系统缓存的资源锁,并将会话添加到本地系统缓存和就绪会话列表的尾部。
PthreadRWlockWrlock(owner, &m_session_rw_locks[hash_index]);
DLAddTail(&m_session_bucket[hash_index], &session->elem2);
PthreadRWlockUnlock(owner, &m_session_rw_locks[hash_index]);
@ -640,10 +697,12 @@ void ThreadPoolListener::AddIdleSessionToHead(knl_session_context* session)
m_readySessionList->AddHead(&session->elem);
return;
}
// 如果不启用本地系统缓存,将会话添加到就绪会话列表的头部。
Index hash_index = HASH_INDEX(session->proc_cxt.MyDatabaseId, (uint32)m_session_nbucket);
ResourceOwner owner = LOCAL_SYSDB_RESOWNER;
// 获取本地系统缓存的资源锁,并将会话添加到本地系统缓存和就绪会话列表的头部。
PthreadRWlockWrlock(owner, &m_session_rw_locks[hash_index]);
DLAddHead(&m_session_bucket[hash_index], &session->elem2);
PthreadRWlockUnlock(owner, &m_session_rw_locks[hash_index]);
m_readySessionList->AddHead(&session->elem);
}
}