This is an automated email from the ASF dual-hosted git repository. oppenheimer01 pushed a commit to branch cbdb-postgres-merge in repository https://gitbox.apache.org/repos/asf/cloudberry.git
commit ca3424694a84053af1581e32cfcac7933b88075f Author: liushengsong <[email protected]> AuthorDate: Wed Apr 29 17:29:23 2026 +0800 Fix resgroup: restore pgstat_report_resgroup, fix is_session_in_group, fix duplicate totalExecuted Three issues fixed: 1. pgstat_report_resgroup was commented out during PG16 merge (MERGE16_FIXME). Restore all call sites so pg_stat_activity correctly shows rsgid/rsgname. 2. In check_and_unassign_from_resgroup, the Assign-to-Bypass transition calls UnassignResGroup which clears st_rsgid, but pgstat_report_resgroup was not called afterward to restore it. This caused pg_stat_activity to show rsgname='unknown' for queries that switch from normal assign to bypass mode. 3. is_session_in_group plpython function used `ps -ef | grep con{session_id}` to find process PIDs, but PG16 removed con{session_id} from process titles. Empty grep result meant empty set, and empty.issubset(any) = True, so the function always returned true. Fixed to use gp_stat_activity JOIN gp_segment_configuration instead. 4. Remove duplicate totalExecuted++ in check_and_unassign_from_resgroup that caused double-counting when a query transitions from Assign to Bypass state. --- src/backend/utils/activity/backend_status.c | 21 +++++++++ src/backend/utils/resgroup/resgroup.c | 54 +++++++++------------- src/include/utils/backend_status.h | 1 + .../resgroup/resgroup_auxiliary_tools_v1.out | 10 ++-- .../resgroup/resgroup_auxiliary_tools_v2.out | 10 ++-- .../sql/resgroup/resgroup_auxiliary_tools_v1.sql | 27 ++++++----- .../sql/resgroup/resgroup_auxiliary_tools_v2.sql | 25 +++++----- 7 files changed, 82 insertions(+), 66 deletions(-) diff --git a/src/backend/utils/activity/backend_status.c b/src/backend/utils/activity/backend_status.c index 36ae361b250..647c4482f09 100644 --- a/src/backend/utils/activity/backend_status.c +++ b/src/backend/utils/activity/backend_status.c @@ -729,6 +729,27 @@ pgstat_report_xact_timestamp(TimestampTz tstamp) PGSTAT_END_WRITE_ACTIVITY(beentry); } +/* ---------- + * pgstat_report_resgroup() - + * + * Called to update the resource group id in MyBEEntry. + * ---------- + */ +void +pgstat_report_resgroup(Oid groupId) +{ + volatile PgBackendStatus *beentry = MyBEEntry; + + if (!beentry) + return; + + PGSTAT_BEGIN_WRITE_ACTIVITY(beentry); + + beentry->st_rsgid = groupId; + + PGSTAT_END_WRITE_ACTIVITY(beentry); +} + /* ---------- * pgstat_read_current_status() - * diff --git a/src/backend/utils/resgroup/resgroup.c b/src/backend/utils/resgroup/resgroup.c index 987bae3c7ad..98284d27a09 100644 --- a/src/backend/utils/resgroup/resgroup.c +++ b/src/backend/utils/resgroup/resgroup.c @@ -53,6 +53,7 @@ #include "utils/memutils.h" #include "utils/ps_status.h" #include "utils/cgroup.h" +#include "utils/backend_status.h" #include "utils/resgroup.h" #include "utils/resource_manager.h" #include "utils/session_state.h" @@ -1290,8 +1291,7 @@ groupAcquireSlot(ResGroupInfo *pGroupInfo, bool isMoveQuery) /* got one, lucky */ group->totalExecuted++; LWLockRelease(ResGroupLock); - /* MERGE16_FIXME report data to pastat */ - //pgstat_report_resgroup(group->groupId); + pgstat_report_resgroup(group->groupId); return slot; } } @@ -1328,7 +1328,7 @@ groupAcquireSlot(ResGroupInfo *pGroupInfo, bool isMoveQuery) group->totalExecuted++; LWLockRelease(ResGroupLock); -// pgstat_report_resgroup(group->groupId); + pgstat_report_resgroup(group->groupId); return slot; } @@ -1567,7 +1567,7 @@ AssignResGroupOnMaster(void) /* Update pg_stat_activity statistics */ bypassedGroup->totalExecuted++; -// pgstat_report_resgroup(bypassedGroup->groupId); + pgstat_report_resgroup(bypassedGroup->groupId); /* Initialize the fake slot */ bypassedSlot.group = groupInfo.group; @@ -1633,7 +1633,7 @@ UnassignResGroup(void) bypassedGroup = NULL; /* Update pg_stat_activity statistics */ -// pgstat_report_resgroup(InvalidOid); + pgstat_report_resgroup(InvalidOid); return; } @@ -1668,7 +1668,7 @@ UnassignResGroup(void) if (Gp_role == GP_ROLE_DISPATCH) SIMPLE_FAULT_INJECTOR("unassign_resgroup_end_qd"); -// pgstat_report_resgroup(InvalidOid); + pgstat_report_resgroup(InvalidOid); } /* @@ -1812,7 +1812,7 @@ waitOnGroup(ResGroupData *group, bool isMoveQuery) * not enough to store a full Oid, so we set groupId out-of-band, * via the backend entry. */ -// pgstat_report_resgroup(group->groupId); + pgstat_report_resgroup(group->groupId); /* * Mark that we are waiting on resource group @@ -2448,7 +2448,6 @@ static void groupWaitQueuePush(ResGroupData *group, PGPROC *proc) { dclist_head *waitQueue; - PGPROC *headProc; Assert(LWLockHeldByMeInMode(ResGroupLock, LW_EXCLUSIVE)); Assert(!procIsWaiting(proc)); @@ -2457,13 +2456,10 @@ groupWaitQueuePush(ResGroupData *group, PGPROC *proc) groupWaitQueueValidate(group); waitQueue = &group->waitProcs; - headProc = (PGPROC *) &waitQueue->dlist.head; - dclist_insert_before(waitQueue, &headProc->links, &proc->links); + dclist_push_tail(waitQueue, &proc->links); groupWaitProcValidate(proc, waitQueue); - waitQueue->count++; - Assert(groupWaitQueueFind(group, proc)); } @@ -2483,14 +2479,12 @@ groupWaitQueuePop(ResGroupData *group) waitQueue = &group->waitProcs; - proc = (PGPROC *) waitQueue->dlist.head.next; + proc = dclist_head_element(PGPROC, links, waitQueue); groupWaitProcValidate(proc, waitQueue); Assert(groupWaitQueueFind(group, proc)); Assert(proc->resSlot == NULL); - dclist_delete_from(waitQueue, &proc->links); - - waitQueue->count--; + dclist_delete_from_thoroughly(waitQueue, &proc->links); return proc; } @@ -2513,9 +2507,7 @@ groupWaitQueueErase(ResGroupData *group, PGPROC *proc) waitQueue = &group->waitProcs; groupWaitProcValidate(proc, waitQueue); - dclist_delete_from(waitQueue, &proc->links); - - waitQueue->count--; + dclist_delete_from_thoroughly(waitQueue, &proc->links); } /* @@ -2772,31 +2764,27 @@ resgroupDumpGroup(StringInfo str, ResGroupData *group) static void resgroupDumpWaitQueue(StringInfo str, dclist_head *queue) { - PGPROC *proc; + dlist_iter iter; + bool first = true; appendStringInfo(str, "\"wait_queue\":{"); appendStringInfo(str, "\"wait_queue_size\":%d,", queue->count); appendStringInfo(str, "\"wait_queue_content\":["); - proc = (PGPROC *)dclist_next_node(queue, &queue->dlist.head); - - if (!ShmemAddrIsValid(&proc->links)) + dclist_foreach(iter, queue) { - appendStringInfo(str, "]},"); - return; - } + PGPROC *proc = dlist_container(PGPROC, links, iter.cur); + + if (!first) + appendStringInfo(str, ","); + first = false; - while (proc) - { appendStringInfo(str, "{"); appendStringInfo(str, "\"pid\":%d,", proc->pid); appendStringInfo(str, "\"resWaiting\":%s,", procIsWaiting(proc) ? "true" : "false"); appendStringInfo(str, "\"resSlot\":%d", slotGetId(proc->resSlot)); appendStringInfo(str, "}"); - proc = (PGPROC *)dclist_next_node(queue, &queue->dlist.head); - if (proc) - appendStringInfo(str, ","); } appendStringInfo(str, "]},"); } @@ -3329,7 +3317,7 @@ HandleMoveResourceGroup(void) cgroupOpsRoutine->attachcgroup(self->groupId, MyProcPid, self->caps.cpuMaxPercent == CPU_MAX_PERCENT_DISABLED); -// pgstat_report_resgroup(self->groupId); + pgstat_report_resgroup(self->groupId); } /* @@ -3698,7 +3686,7 @@ check_and_unassign_from_resgroup(PlannedStmt* stmt) } while (!groupIncBypassedRef(&groupInfo)); bypassedGroup = groupInfo.group; - bypassedGroup->totalExecuted++; + pgstat_report_resgroup(bypassedGroup->groupId); bypassedSlot.group = groupInfo.group; bypassedSlot.groupId = groupInfo.groupId; diff --git a/src/include/utils/backend_status.h b/src/include/utils/backend_status.h index f01acf0b01d..d8fa3854ca1 100644 --- a/src/include/utils/backend_status.h +++ b/src/include/utils/backend_status.h @@ -325,6 +325,7 @@ extern void pgstat_report_query_id(uint64 query_id, bool force); extern void pgstat_report_tempfile(size_t filesize); extern void pgstat_report_appname(const char *appname); extern void pgstat_report_xact_timestamp(TimestampTz tstamp); +extern void pgstat_report_resgroup(Oid groupId); extern const char *pgstat_get_backend_current_activity(int pid, bool checkUser); extern const char *pgstat_get_crashed_backend_activity(int pid, char *buffer, int buflen); diff --git a/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v1.out b/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v1.out index 17bfbb23fee..6968a020ec3 100644 --- a/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v1.out +++ b/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v1.out @@ -117,10 +117,10 @@ CREATE 0: CREATE OR REPLACE FUNCTION is_session_in_group(pid integer, groupname text) RETURNS BOOL AS $$ import subprocess sql = "select sess_id from pg_stat_activity where pid = '%d'" % pid result = plpy.execute(sql) session_id = result[0]['sess_id'] sql = "select groupid from gp_toolkit.gp_resgroup_config where groupname='%s'" % groupname result = plpy.execute(sql) groupid = result[0]['groupid'] -sql = "select hostname from gp_segment_configuration group by hostname" result = plpy.execute(sql) hosts = [_['hostname'] for _ in result] -def get_result(host): stdout = subprocess.run(["ssh", "{}".format(host), "ps -ef | grep postgres | grep con{} | grep -v grep | awk '{{print $2}}'".format(session_id)], check=True, stdout=subprocess.PIPE).stdout session_pids = stdout.splitlines() -path = "/sys/fs/cgroup/cpu/gpdb/{}/cgroup.procs".format(groupid) stdout = subprocess.run(["ssh", "{}".format(host), "cat {}".format(path)], check=True, stdout=subprocess.PIPE).stdout cgroups_pids = stdout.splitlines() -return set(session_pids).issubset(set(cgroups_pids)) -for host in hosts: if not get_result(host): return False return True +sql = """select sc.hostname, array_agg(pa.pid::text) as pids from gp_stat_activity pa join gp_segment_configuration sc on pa.gp_segment_id = sc.content and sc.role = 'p' where pa.sess_id = %d group by sc.hostname""" % session_id result = plpy.execute(sql) +for row in result: host = row['hostname'] session_pids = set(row['pids']) +if not session_pids: continue +path = "/sys/fs/cgroup/cpu/gpdb/{}/cgroup.procs".format(groupid) stdout = subprocess.run(["ssh", "{}".format(host), "cat {}".format(path)], stdout=subprocess.PIPE, check=True).stdout cgroups_pids = set(stdout.decode().splitlines()) +if not session_pids.issubset(cgroups_pids): return False return True $$ LANGUAGE plpython3u; CREATE diff --git a/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v2.out b/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v2.out index 05a8ea81b73..92c2e1fa33f 100644 --- a/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v2.out +++ b/src/test/isolation2/expected/resgroup/resgroup_auxiliary_tools_v2.out @@ -133,11 +133,11 @@ CREATE 0: CREATE OR REPLACE FUNCTION is_session_in_group(pid integer, groupname text) RETURNS BOOL AS $$ import subprocess sql = "select sess_id from pg_stat_activity where pid = '%d'" % pid result = plpy.execute(sql) session_id = result[0]['sess_id'] sql = "select groupid from gp_toolkit.gp_resgroup_config where groupname='%s'" % groupname result = plpy.execute(sql) groupid = result[0]['groupid'] -sql = "select hostname from gp_segment_configuration group by hostname" result = plpy.execute(sql) hosts = [_['hostname'] for _ in result] -def get_result(host): stdout = subprocess.run(["ssh", "{}".format(host), "ps -ef | grep postgres | grep con{} | grep -v grep | awk '{{print $2}}'".format(session_id)], stdout=subprocess.PIPE, check=True).stdout session_pids = stdout.splitlines() -path = "/sys/fs/cgroup/gpdb/{}/queries/cgroup.procs".format(groupid) stdout = subprocess.run(["ssh", "{}".format(host), "cat {}".format(path)], stdout=subprocess.PIPE, check=True).stdout cgroups_pids = stdout.splitlines() -return set(session_pids).issubset(set(cgroups_pids)) -for host in hosts: if not get_result(host): return False return True +sql = """select sc.hostname, array_agg(pa.pid::text) as pids from gp_stat_activity pa join gp_segment_configuration sc on pa.gp_segment_id = sc.content and sc.role = 'p' where pa.sess_id = %d group by sc.hostname""" % session_id result = plpy.execute(sql) +for row in result: host = row['hostname'] session_pids = set(row['pids']) +if not session_pids: continue +path = "/sys/fs/cgroup/gpdb/{}/queries/cgroup.procs".format(groupid) stdout = subprocess.run(["ssh", "{}".format(host), "cat {}".format(path)], stdout=subprocess.PIPE, check=True).stdout cgroups_pids = set(stdout.decode().splitlines()) +if not session_pids.issubset(cgroups_pids): return False return True $$ LANGUAGE plpython3u; CREATE diff --git a/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v1.sql b/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v1.sql index 1e09cd70809..cd71081bae6 100644 --- a/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v1.sql +++ b/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v1.sql @@ -255,23 +255,26 @@ $$ LANGUAGE plpython3u; result = plpy.execute(sql) groupid = result[0]['groupid'] - sql = "select hostname from gp_segment_configuration group by hostname" + sql = """select sc.hostname, array_agg(pa.pid::text) as pids + from gp_stat_activity pa + join gp_segment_configuration sc + on pa.gp_segment_id = sc.content and sc.role = 'p' + where pa.sess_id = %d + group by sc.hostname""" % session_id result = plpy.execute(sql) - hosts = [_['hostname'] for _ in result] - def get_result(host): - stdout = subprocess.run(["ssh", "{}".format(host), "ps -ef | grep postgres | grep con{} | grep -v grep | awk '{{print $2}}'".format(session_id)], - check=True, stdout=subprocess.PIPE).stdout - session_pids = stdout.splitlines() + for row in result: + host = row['hostname'] + session_pids = set(row['pids']) - path = "/sys/fs/cgroup/cpu/gpdb/{}/cgroup.procs".format(groupid) - stdout = subprocess.run(["ssh", "{}".format(host), "cat {}".format(path)], check=True, stdout=subprocess.PIPE).stdout - cgroups_pids = stdout.splitlines() + if not session_pids: + continue - return set(session_pids).issubset(set(cgroups_pids)) + path = "/sys/fs/cgroup/cpu/gpdb/{}/cgroup.procs".format(groupid) + stdout = subprocess.run(["ssh", "{}".format(host), "cat {}".format(path)], stdout=subprocess.PIPE, check=True).stdout + cgroups_pids = set(stdout.decode().splitlines()) - for host in hosts: - if not get_result(host): + if not session_pids.issubset(cgroups_pids): return False return True diff --git a/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v2.sql b/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v2.sql index 38e0a24d3a3..9f123b0e5c1 100644 --- a/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v2.sql +++ b/src/test/isolation2/sql/resgroup/resgroup_auxiliary_tools_v2.sql @@ -261,23 +261,26 @@ $$ LANGUAGE plpython3u; result = plpy.execute(sql) groupid = result[0]['groupid'] - sql = "select hostname from gp_segment_configuration group by hostname" + sql = """select sc.hostname, array_agg(pa.pid::text) as pids + from gp_stat_activity pa + join gp_segment_configuration sc + on pa.gp_segment_id = sc.content and sc.role = 'p' + where pa.sess_id = %d + group by sc.hostname""" % session_id result = plpy.execute(sql) - hosts = [_['hostname'] for _ in result] - def get_result(host): - stdout = subprocess.run(["ssh", "{}".format(host), "ps -ef | grep postgres | grep con{} | grep -v grep | awk '{{print $2}}'".format(session_id)], - stdout=subprocess.PIPE, check=True).stdout - session_pids = stdout.splitlines() + for row in result: + host = row['hostname'] + session_pids = set(row['pids']) + + if not session_pids: + continue path = "/sys/fs/cgroup/gpdb/{}/queries/cgroup.procs".format(groupid) stdout = subprocess.run(["ssh", "{}".format(host), "cat {}".format(path)], stdout=subprocess.PIPE, check=True).stdout - cgroups_pids = stdout.splitlines() - - return set(session_pids).issubset(set(cgroups_pids)) + cgroups_pids = set(stdout.decode().splitlines()) - for host in hosts: - if not get_result(host): + if not session_pids.issubset(cgroups_pids): return False return True --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
