Update streamConsumer.cpp

This commit is contained in:
Nemoo 2023-09-30 21:06:42 +08:00
parent fb37311a96
commit 7f78629009
1 changed files with 106 additions and 98 deletions

View File

@ -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;
}
}
}