From 7f7862900913a2e5c4ff5d071cae72f5a72eeaa4 Mon Sep 17 00:00:00 2001 From: Nemoo <18210033860@163.com> Date: Sat, 30 Sep 2023 21:06:42 +0800 Subject: [PATCH] Update streamConsumer.cpp --- .../process/stream/streamConsumer.cpp | 204 +++++++++--------- 1 file changed, 106 insertions(+), 98 deletions(-) diff --git a/src/gausskernel/process/stream/streamConsumer.cpp b/src/gausskernel/process/stream/streamConsumer.cpp index 85d801d56..1471b6093 100755 --- a/src/gausskernel/process/stream/streamConsumer.cpp +++ b/src/gausskernel/process/stream/streamConsumer.cpp @@ -14,7 +14,7 @@ * ------------------------------------------------------------------------- * * streamConsumer.cpp - * Support methods for class StreamConsumer. + * 支持StreamConsumer的方法 * * IDENTIFICATION * src/gausskernel/process/stream/streamConsumer.cpp @@ -36,20 +36,22 @@ extern GlobalNodeDefinition* global_node_definition; +// 构造函数,初始化StreamConsumer对象 StreamConsumer::StreamConsumer(MemoryContext context) : StreamObj(context, STREAM_CONSUMER) { - m_sharedContext = NULL; - m_originProducerNodeList = NULL; - m_ready = false; - m_expectProducer = NULL; - m_currentProducerNum = 0; + m_sharedContext = NULL; // 初始化共享内存上下文为NULL + m_originProducerNodeList = NULL; // 初始化原始生产者节点列表为NULL + m_ready = false; // 初始化就绪状态为false + m_expectProducer = NULL; // 初始化预期生产者为NULL + m_currentProducerNum = 0; // 初始化当前生产者数量为0 } +// 析构函数,释放StreamConsumer对象的资源 StreamConsumer::~StreamConsumer() { - m_originProducerNodeList = NULL; - m_expectProducer = NULL; - m_sharedContext = NULL; + m_originProducerNodeList = NULL; // 将原始生产者节点列表设为NULL + m_expectProducer = NULL; // 将预期生产者设为NULL + m_sharedContext = NULL; // 将共享内存上下文设为NULL } /* @@ -61,12 +63,16 @@ StreamConsumer::~StreamConsumer() * @param[IN] transType: transport type * @return: void */ + /* + 函数通过参数设置了StreamConsumer对象的各个成员变量,包括传输类型、共享上下文、传输对象等,还根据传入的执行生产者节点信息, + 初始化了预期的生产者信息,并将其加入哈希表中 + */ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc desc, StreamTransType transType, StreamSharedContext* sharedContext) { int i = 0; bool found = false; - + int producerNum = 0; AutoMutexLock streamLock(&m_streamInfoLock); AutoMutexLock copyLock(&nodeDefCopyLock); @@ -79,18 +85,19 @@ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc d Assert(producerNum > 0); - m_parallel_desc = desc; - m_currentProducerNum = 0; - m_connNum = producerNum; - m_key = key; - m_ready = false; - m_transtype = transType; - m_sharedContext = sharedContext; - m_transport = (StreamTransport**)MemoryContextAllocZero(m_memoryCxt, producerNum * sizeof(StreamTransport*)); - m_expectProducer = (StreamConnInfo*)MemoryContextAllocZero(m_memoryCxt, producerNum * sizeof(StreamConnInfo)); + m_parallel_desc = desc; // 设置并行描述符 + m_currentProducerNum = 0; // 当前生产者数量置零 + m_connNum = producerNum; // 连接数量等于生产者数量 + m_key = key; // 设置StreamKey + m_ready = false; // 初始状态为未就绪 + m_transtype = transType; // 设置传输类型 + m_sharedContext = sharedContext; // 共享上下文 + m_transport = (StreamTransport**)MemoryContextAllocZero(m_memoryCxt, producerNum * sizeof(StreamTransport*)); // 分配传输对象数组内存 + m_expectProducer = (StreamConnInfo*)MemoryContextAllocZero(m_memoryCxt, producerNum * sizeof(StreamConnInfo)); // 分配预期生产者信息内存 - /* Initialize the origin nodelist */ + /* 初始化原始生产者节点列表 */ m_originProducerNodeList = NIL; + #ifdef ENABLE_MULTIPLE_NODES copyLock.lock(); ListCell* nodelistCell = NULL; @@ -112,6 +119,8 @@ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc d } copyLock.unLock(); #endif + + // 初始化传输对象数组 for (i = 0; i < producerNum; i++) { int nodeNameLen = 0; libcommaddrinfo* libcommaddr = NULL; @@ -131,7 +140,7 @@ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc d scomm->m_addr->streamKey.consumerSmpId = m_key.smpIdentifier; } - HOLD_INTERRUPTS(); /* Add this macro for double safety. */ + HOLD_INTERRUPTS(); // 加入此宏以确保安全性 streamLock.lock(); AutoContextSwitch streamInfoCxtGuard(StreamInfoContext); @@ -139,8 +148,7 @@ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc d StreamElement* element = (StreamElement*)hash_search(m_streamInfoTbl, &m_key, HASH_ENTER, &found); if (element == NULL) { streamLock.unLock(); - ereport( - ERROR, (errcode(ERRCODE_OUT_OF_MEMORY), errmsg("Failed to create stream element due to out of memory"))); + ereport(ERROR, (errcode(ERRCODE_OUT_OF_MEMORY), errmsg("Failed to create stream element due to out of memory"))); } if (found == false) { @@ -154,16 +162,15 @@ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc d if (NULL == element->value) { hash_search(m_streamInfoTbl, &key, HASH_REMOVE, NULL); streamLock.unLock(); - ereport(ERROR, - (errcode(ERRCODE_OUT_OF_MEMORY), errmsg("Failed to generate stream element due to out of memory"))); + ereport(ERROR, (errcode(ERRCODE_OUT_OF_MEMORY), errmsg("Failed to generate stream element due to out of memory"))); } element->value->connNum = 0; element->value->connInfoSize = 0; - /* No need to allocate the array size. */ + /* 不需要分配数组大小。 */ element->value->connInfo = NULL; } else { - /* WTF? find a duplicate element in the hash table. */ + /* 在哈希表中找到重复的元素。 */ if (element->value == NULL || element->value->consumer) { streamLock.unLock(); ereport(ERROR, @@ -173,7 +180,7 @@ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc d } } - /* Found true means some information has not been updated, protected by streamLock. */ + /* 找到为真表示某些信息尚未更新,由streamLock保护。 */ if (found == true) { Assert(element->value->connInfo != NULL); updateTransportInfo(element->value); @@ -215,65 +222,64 @@ void StreamConsumer::init(StreamKey key, List* execProducerNodes, ParallelDesc d */ void StreamConsumer::deInit() { - AutoMutexLock streamHashLock(&m_streamInfoLock); + AutoMutexLock streamHashLock(&m_streamInfoLock); // 加锁,保护哈希表操作 StreamElement* delinfo = NULL; - HOLD_INTERRUPTS(); - streamHashLock.lock(); + HOLD_INTERRUPTS(); // 持有中断,禁止中断发生 + streamHashLock.lock(); // 加锁,保护哈希表操作 - /* Do not need de init. */ - if (m_init == false) { - if (m_threadSyncObjInit == true) { - pthread_mutex_destroy(&m_mutex); - pthread_cond_destroy(&m_cond); - m_threadSyncObjInit = false; + /* 不需要去初始化。 */ + if (m_init == false) { // 如果未初始化 + if (m_threadSyncObjInit == true) { // 如果线程同步对象已经初始化 + pthread_mutex_destroy(&m_mutex); // 销毁互斥锁 + pthread_cond_destroy(&m_cond); // 销毁条件变量 + m_threadSyncObjInit = false; // 标记线程同步对象未初始化 } - streamHashLock.unLock(); - RESUME_INTERRUPTS(); + streamHashLock.unLock(); // 解锁,释放锁 + RESUME_INTERRUPTS(); // 恢复中断 return; } - releaseCommStream(); + releaseCommStream(); // 释放通信流 - /* Close local stream connection. */ + /* 关闭本地流连接。 */ if (NULL != m_sharedContext) { - gs_memory_close_conn(m_sharedContext, m_connNum, u_sess->stream_cxt.smp_id); + gs_memory_close_conn(m_sharedContext, m_connNum, u_sess->stream_cxt.smp_id); // 关闭共享内存中的连接 } - if (m_ready || u_sess->stream_cxt.dummy_thread == true) { - delinfo = (StreamElement*)hash_search(m_streamInfoTbl, &m_key, HASH_REMOVE, NULL); - if (delinfo != NULL) { + if (m_ready || u_sess->stream_cxt.dummy_thread == true) { // 如果已经就绪或者是虚拟线程 + delinfo = (StreamElement*)hash_search(m_streamInfoTbl, &m_key, HASH_REMOVE, NULL); // 从哈希表中移除该消费者信息 + if (delinfo != NULL) { // 如果找到了消费者信息 if (delinfo->value->connInfo != NULL) { - pfree_ext(delinfo->value->connInfo); + pfree_ext(delinfo->value->connInfo); // 释放连接信息内存 delinfo->value->connInfo = NULL; } - pfree_ext(delinfo->value); + pfree_ext(delinfo->value); // 释放消费者信息内存 delinfo->value = NULL; } } else { /* - * set flag to release net port in case some producer not ready and element - * in stream info hash table can not get removed in wakeUpConsumer, - * or consumer get canceled when waiting producer ready. + * 设置标志以释放网络端口,在某些生产者未就绪的情况下, + * 或者在等待生产者就绪时,如果消费者被取消。 */ - u_sess->stream_cxt.global_obj->setNeedClean(true); + u_sess->stream_cxt.global_obj->setNeedClean(true); // 设置需要清理标志 } - if (m_threadSyncObjInit == true) { - pthread_mutex_destroy(&m_mutex); - pthread_cond_destroy(&m_cond); - m_threadSyncObjInit = false; + if (m_threadSyncObjInit == true) { // 如果线程同步对象已经初始化 + pthread_mutex_destroy(&m_mutex); // 销毁互斥锁 + pthread_cond_destroy(&m_cond); // 销毁条件变量 + m_threadSyncObjInit = false; // 标记线程同步对象未初始化 } - m_init = false; - streamHashLock.unLock(); + m_init = false; // 标记消费者未初始化 + streamHashLock.unLock(); // 解锁,释放锁 - /* Cleanup the original producer list */ - m_originProducerNodeList = NIL; + /* 清理原始生产者列表 */ + m_originProducerNodeList = NIL; // 清空原始生产者节点列表 - RESUME_INTERRUPTS(); + RESUME_INTERRUPTS(); // 恢复中断 } /* @@ -283,12 +289,13 @@ void StreamConsumer::deInit() */ void StreamConsumer::releaseCommStream() { - if (m_transport != NULL) { - for (int i = 0; i < m_connNum; i++) { - Assert(t_thrd.int_cxt.ImmediateInterruptOK == false); - StreamCOMM* scomm = (StreamCOMM*)m_transport[i]; - if (scomm != NULL) - scomm->release(); + if (m_transport != NULL) { // 如果传输对象数组不为空 + for (int i = 0; i < m_connNum; i++) { // 遍历所有连接 + Assert(t_thrd.int_cxt.ImmediateInterruptOK == false); // 断言当前不可中断 + + StreamCOMM* scomm = (StreamCOMM*)m_transport[i]; // 获取当前连接的通信对象 + if (scomm != NULL) // 如果通信对象不为空 + scomm->release(); // 释放通信资源 } } } @@ -301,11 +308,11 @@ void StreamConsumer::releaseCommStream() */ int StreamConsumer::getNodeIdx(const char* nodename) { - for (int i = 0; i < m_connNum; i++) { - if (pg_strncasecmp(m_expectProducer[i].nodeName, nodename, strlen(nodename)) == 0) - return m_expectProducer[i].nodeIdx; + for (int i = 0; i < m_connNum; i++) { // 遍历所有生产者节点 + if (pg_strncasecmp(m_expectProducer[i].nodeName, nodename, strlen(nodename)) == 0) // 如果节点名匹配 + return m_expectProducer[i].nodeIdx; // 返回节点索引 } - return -1; + return -1; // 如果未找到匹配节点,返回-1 } /* @@ -318,17 +325,17 @@ void StreamConsumer::findUnconnectProducer(StringInfo str) { bool found = false; - for (int j = 0; j < m_connNum; j++) { + for (int j = 0; j < m_connNum; j++) { // 遍历所有预期的生产者节点 found = false; - for (int i = 0; i < m_currentProducerNum; i++) { - if (strcmp(m_expectProducer[j].nodeName, m_transport[i]->m_nodeName) == 0) { - found = true; + for (int i = 0; i < m_currentProducerNum; i++) { // 遍历当前已连接的生产者节点 + if (strcmp(m_expectProducer[j].nodeName, m_transport[i]->m_nodeName) == 0) { // 如果找到匹配节点 + found = true; // 标记为已连接 break; } } - if (!found) - appendStringInfo(str, " %s", m_expectProducer[j].nodeName); + if (!found) // 如果未找到匹配节点(即未连接) + appendStringInfo(str, " %s", m_expectProducer[j].nodeName); // 将节点名追加到字符串中 } } @@ -342,21 +349,21 @@ int StreamConsumer::getFirstUnconnectedProducerNodeIdx() bool found = false; int nodeIdx = -1; - for (int j = 0; j < m_connNum; j++) { + for (int j = 0; j < m_connNum; j++) { // 遍历所有预期的生产者节点 found = false; - for (int i = 0; i < m_currentProducerNum; i++) { - if (strcmp(m_expectProducer[j].nodeName, m_transport[i]->m_nodeName) == 0) { - found = true; + for (int i = 0; i < m_currentProducerNum; i++) { // 遍历当前已连接的生产者节点 + if (strcmp(m_expectProducer[j].nodeName, m_transport[i]->m_nodeName) == 0) { // 如果找到匹配节点 + found = true; // 标记为已连接 break; } } - if (!found) { - nodeIdx = m_expectProducer[j].nodeIdx; + if (!found) { // 如果未找到匹配节点(即未连接) + nodeIdx = m_expectProducer[j].nodeIdx; // 设置未连接节点的索引 break; } } - return nodeIdx; + return nodeIdx; // 返回第一个未连接节点的索引 } /* @@ -366,7 +373,7 @@ int StreamConsumer::getFirstUnconnectedProducerNodeIdx() */ void StreamConsumer::waitProducerReady() { - /* For local stream, producer will never connect to consumer. Consumer is ready, just return. */ + /* 对于本地流,生产者永远不会连接到消费者。消费者准备好了,直接返回。 */ if (STREAM_IS_LOCAL_NODE(m_parallel_desc.distriType)) { m_ready = true; return; @@ -384,7 +391,7 @@ void StreamConsumer::waitProducerReady() Assert(t_thrd.int_cxt.ImmediateInterruptOK == false); streamLock.lock(); - /* 900s timeout. */ + /* 900秒的超时时间。 */ struct timespec timer; int ret; int ntimes = 1; @@ -419,11 +426,11 @@ void StreamConsumer::waitProducerReady() ereport(ERROR, (errmodule(MOD_STREAM), errcode(ERRCODE_CONNECTION_TIMED_OUT), - errmsg("Distribute query initializing network connection timeout. un-connected nodes: %s", + errmsg("Distribute query initializing network connection timeout. un-connected nodes:%s", str.data))); } - /* Check for interrupts.(cancel signal?). */ + /* 检查中断(取消信号?)。 */ CHECK_FOR_INTERRUPTS(); streamLock.lock(); } @@ -442,6 +449,7 @@ void StreamConsumer::waitProducerReady() return; } + /* * @Description: Wake up consumer and let it work * @@ -465,7 +473,7 @@ bool StreamConsumer::wakeUpConsumerCallBack(CommStreamKey commKey, StreamConnInf streamLock.lock(); - /* Check if can accept connection now */ + /* 检查是否可以接受连接。 */ if (StreamNodeGroup::checkStreamConnectPermission(key.queryId) == false) { streamLock.unLock(); return false; @@ -474,9 +482,9 @@ bool StreamConsumer::wakeUpConsumerCallBack(CommStreamKey commKey, StreamConnInf AutoContextSwitch streamInfoCxtGuard(StreamInfoContext); /* - * Register me in the global consumer table if the consumer thread has not been started. - * Libcomm r_flow_ctrl thread call this, must not use palloc and elog. - * so use isLibcommThread flag, return false where malloc failed. + * 如果消费者线程尚未启动,则在全局消费者表中注册我。 + * Libcomm r_flow_ctrl线程调用这个函数,不能使用palloc和elog。 + * 所以使用isLibcommThread标志,在malloc失败时返回false。 */ StreamElement* element = (StreamElement*)hash_search(m_streamInfoTbl, &key, HASH_ENTER, &found); if (element == NULL) { @@ -485,7 +493,7 @@ bool StreamConsumer::wakeUpConsumerCallBack(CommStreamKey commKey, StreamConnInf return false; } - /* If StreamElement not register by the consumer thread, we save the information. */ + /* 如果StreamElement尚未由消费者线程注册,则保存信息。 */ if (found == false) { element->value = (StreamValue*)palloc0_noexcept(sizeof(StreamValue)); if (NULL == element->value) { @@ -516,7 +524,7 @@ bool StreamConsumer::wakeUpConsumerCallBack(CommStreamKey commKey, StreamConnInf element->value->consumer = NULL; } else { if (element->value) { - /* Consumer thread has not been registered yet. */ + /* 消费者线程尚未注册。 */ if (element->value->consumer == NULL) { element->value->connInfo[element->value->connNum].port.libcomm_layer.gsock = connInfo.port.libcomm_layer.gsock; @@ -528,7 +536,7 @@ bool StreamConsumer::wakeUpConsumerCallBack(CommStreamKey commKey, StreamConnInf element->value->connInfo[element->value->connNum].producerSmpId = connInfo.producerSmpId; element->value->connNum++; - /* Check if need realloc to remember more connection info. */ + /* 检查是否需要重新分配内存以记住更多的连接信息。 */ if (element->value->connNum == element->value->connInfoSize) { StreamConnInfo* new_connInfo = NULL; int new_connInfoSize = 2 * element->value->connInfoSize; @@ -557,10 +565,10 @@ bool StreamConsumer::wakeUpConsumerCallBack(CommStreamKey commKey, StreamConnInf element->value->connInfoSize = new_connInfoSize; } } else { - /* Have registered, so we update the stream info. */ + /* 已经注册,所以我们更新流信息。 */ res = element->value->consumer->updateStreamCommInfo(&connInfo); - /* We can free the memory, no longer need it. */ + /* 我们可以释放内存了,不再需要它。 */ if (element->value->connInfo != NULL) { pfree_ext(element->value->connInfo); element->value->connInfo = NULL; @@ -588,7 +596,7 @@ bool StreamConsumer::updateStreamCommInfo(StreamConnInfo* connInfo) streamLock.lock(); - /* WTF, duplicate plan id encounter or plan error. */ + /* 出现重复的计划ID或计划错误。 */ if (m_currentProducerNum >= m_connNum) { streamLock.unLock(); return false; @@ -632,4 +640,4 @@ void StreamConsumer::updateTransportInfo(StreamValue* val) if (m_currentProducerNum == m_connNum) m_ready = true; } -} +} \ No newline at end of file