openGauss-server/src/gausskernel/runtime/executor/execClusterResize.cpp

1239 lines
46 KiB
C++
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

/* -------------------------------------------------------------------------
*
* execClusterResize.cpp
* MPPDB ClusterResizing relevant routines
*
* 部分版权 c 2020 华为技术有限公司
* 部分版权所有 c 1996-2012PostgreSQL 全球开发集团
* 部分版权 c 1994加州大学摄政
*
* IDENTIFICATION
* src/gausskernel/runtime/executor/execClusterResize.cpp
*
* -------------------------------------------------------------------------
*/
#include "postgres.h"
#include "knl/knl_variable.h"
#include "access/tableam.h"
#include "catalog/heap.h"
#include "catalog/index.h"
#include "catalog/namespace.h"
#include "catalog/pg_namespace.h"
#include "catalog/pgxc_group.h"
#include "executor/executor.h"
#include "executor/node/nodeModifyTable.h"
#include "optimizer/clauses.h"
#include "parser/analyze.h"
#include "parser/parsetree.h"
#include "pgxc/pgxc.h"
#include "pgxc/redistrib.h"
#include "tcop/utility.h"
#include "utils/builtins.h"
#include "utils/elog.h"
#include "utils/guc.h"
#include "utils/inval.h"
#include "utils/lsyscache.h"
#include "utils/rel.h"
#include "utils/rel_gs.h"
#include "utils/syscache.h"
#include "utils/snapmgr.h"
#include "storage/lmgr.h"
#include "postgres_ext.h"
/*
* ---------------------------------------------------------------------------------
* 局部函数/变量声明字段*
* ---------------------------------------------------------------------------------
*/
/*删除增量表定义 */
#define Natts_pg_delete_delta 3
#define Anum_pg_delete_delta_xcnodeid_and_dntableoid 1
#define Anum_pg_delete_delta_tablebucketid_and_ctid 2
#define RANGE_SCAN_IN_REDIS "tidge+tid+pg_get_redis_rel_start_ctid+tidle+tid+pg_get_redis_rel_end_ctid+"
static Node* eval_dnstable_func_mutator(
Relation rel, Node* node, StringInfo qual_str, RangeScanInRedis *rangeScanInRedis, bool isRoot);
static inline bool redis_tupleid_retrive_function(const char* funcname, Oid rettype, const Oid* argstype, int nargs);
static inline bool redis_blocknum_retrive_function(const char* funcname, Oid rettype, const Oid* argstype, int nargs);
static inline bool redis_offset_retrive_function(const char* funcname, Oid rettype, const Oid* argstype, int nargs);
#define REDIS_TUPLEID_RETRIVE_FUNCSIG(rettype, argstype, nargs) \
(((nargs) == 1 && (rettype) == TIDOID && (argstype)[0] == TEXTOID) || \
((nargs) == 4 && (rettype) == TIDOID && (argstype)[0] == TEXTOID && (argstype)[1] == NAMEOID && \
(argstype)[2] == INT4OID && (argstype)[3] == INT4OID))
static inline bool redis_tupleid_retrive_function(const char* funcname, Oid rettype, const Oid* argstype, int nargs)
{
if (pg_strcasecmp(funcname, "pg_get_redis_rel_start_ctid") == 0 &&
REDIS_TUPLEID_RETRIVE_FUNCSIG(rettype, argstype, nargs)) {
return true;
}
if (pg_strcasecmp(funcname, "pg_get_redis_rel_end_ctid") == 0 &&
REDIS_TUPLEID_RETRIVE_FUNCSIG(rettype, argstype, nargs)) {
return true;
}
return false;
}
static inline bool redis_offset_retrive_function(const char* funcname, Oid rettype, const Oid* argstype, int nargs)
{
if (pg_strcasecmp(funcname, "pg_tupleid_get_offset") == 0 &&
(nargs == 1 && rettype == INT4OID && argstype[0] == TIDOID)) {
return true;
}
return false;
}
static inline bool redis_blocknum_retrive_function(const char* funcname, Oid rettype, const Oid* argstype, int nargs)
{
if (pg_strcasecmp(funcname, "pg_tupleid_get_blocknum") == 0 &&
(nargs == 1 && rettype == INT8OID && argstype[0] == TIDOID)) {
return true;
}
return false;
}
static inline bool redis_ctid_retrive_function(const char* funcname, Oid rettype, const Oid* argstype, int nargs)
{
if (pg_strcasecmp(funcname, "pg_tupleid_get_ctid_to_bigint") == 0 &&
(nargs == 1 && rettype == INT8OID && argstype[0] == TIDOID)) {
return true;
}
return false;
}
/*
*简介将给定元组的元组记录到pg_delete_delta表中
* -参数:
* @rel更新/删除操作的目标关系
* @tupleid需要记录的元组
*-返回:
* 无返回值
*/
void RecordDeletedTuple(Oid relid, int2 bucketid, const ItemPointer tupleid, const Relation deldelta_rel)
{
Datum values[Natts_pg_delete_delta];
bool nulls[Natts_pg_delete_delta];
HeapTuple tup = NULL;
Assert(deldelta_rel);
/*在重新分发中,表 delete_delta 有 3 列或 2 列。 */
Assert(RelationGetDescr(deldelta_rel)->natts <= 3);
/*循环访问初始化空值和值的属性 */
for (int i = 0; i < Natts_pg_delete_delta; i++) {
nulls[i] = false;
values[i] = (Datum)0;
}
values[Anum_pg_delete_delta_xcnodeid_and_dntableoid - 1] =
UInt64GetDatum(((uint64)u_sess->pgxc_cxt.PGXCNodeIdentifier << 32) | relid);
values[Anum_pg_delete_delta_tablebucketid_and_ctid - 1] =
UInt64GetDatum(((uint64)ItemPointerGetBlockNumber(tupleid) << 16) | ItemPointerGetOffsetNumber(tupleid));
if (BUCKET_NODE_IS_VALID(bucketid)) {
values[Anum_pg_delete_delta_tablebucketid_and_ctid - 1] |= ((uint64)bucketid << 48);
}
/* 记录增量 */
tup = heap_form_tuple(RelationGetDescr(deldelta_rel), values, nulls);
(void)simple_heap_insert(deldelta_rel, tup);
tableam_tops_free_tuple(tup);
}
/*
* - 简介:确定关系是否正在执行群集大小调整操作
* - 参数:
* @rel需要检查的关系
* - 返回:
* @TRUE关系正在调整集群大小
* @FALSE: 关系未调整集群大小
*/
bool RelationInClusterResizing(const Relation rel)
{
Assert(rel != NULL);
/*检查关系的append_mode状态 */
if (!IsInitdb && RelationInRedistribute(rel))
return true;
return false;
}
/*
* - 简要:确定关系是否处于集群调整只读操作下
* - 参数:
* @rel:需要检查的关系
* - 返回:
* @TRUE: 关系处于集群调整大小只读状态
* @FALSE: 关系不处于集群调整大小只读状态
*/
bool RelationInClusterResizingReadOnly(const Relation rel)
{
Assert(rel != NULL);
/*检查关系的append_mode状态 */
if (!IsInitdb && RelationInRedistributeReadOnly(rel))
return true;
return false;
}
/*
* - 简要:确定关系是否处于集群调整只读操作下
* - 参数:
* @rel: 需要检查的关系
* - 返回:
* @TRUE: 关系处于群集调整大小状态endcatchup(写错误)
* @FALSE: 关系不在群集调整大小范围内endcatchup(写错误)
*/
bool RelationInClusterResizingEndCatchup(const Relation rel)
{
Assert(rel != NULL);
/* 检查关系的append_mode状态*/
if (!IsInitdb && RelationInRedistributeEndCatchup(rel))
return true;
return false;
}
/*
* @说明:通过范围变量检查关系是否在重新分配。
* @在range_var:存储关系信息的范围变量。
* @在重新分配中返回:true。
*/
bool CheckRangeVarInRedistribution(const RangeVar* range_var)
{
Relation relation;
Oid relid;
bool in_redis = false;
relid = RangeVarGetRelid(range_var, AccessShareLock, true);
if (OidIsValid(relid)) {
relation = relation_open(relid, NoLock);
/* 如果关系是索引,我们应该检查相关表是否在调整大小。*/
if (RelationIsIndex(relation)) {
Oid heapOid = IndexGetRelation(relid, false);
Relation heapRelation = relation_open(heapOid, AccessShareLock);
in_redis = RelationInClusterResizing(heapRelation);
relation_close(heapRelation, AccessShareLock);
} else {
in_redis = RelationInClusterResizing(relation);
}
relation_close(relation, NoLock);
UnlockRelationOid(relid, AccessShareLock);
}
return in_redis;
}
/*
* - 简要:确定表名是否为delete_delta table。
* - 参数:
* @relname: 目标表名
* - 返回:
* @TRUE: 表为delete_delta表
* @FALSE: 这个表不是delete_delta表
*/
bool RelationIsDeleteDeltaTable(char* delete_delta_name)
{
Oid relid;
uint64 val;
char* endptr = NULL;
HeapTuple tuple;
if (IsInitdb) {
return false;
}
if (strncmp(delete_delta_name, "pg_delete_delta_", 16) != 0) {
return false;
}
val = strtoull(delete_delta_name + 16, &endptr, 0);
if ((errno == ERANGE) || (errno != 0 && val == 0)) {
return false;
}
if (endptr == delete_delta_name + 16 || *endptr != '\0') {
return false;
}
relid = (Oid)val;
if (!OidIsValid(relid)) {
return false;
}
tuple = SearchSysCache1(RELOID, ObjectIdGetDatum(relid));
if (!HeapTupleIsValid(tuple)) {
elog(WARNING, "Table %u related to %s does not exists.", relid, delete_delta_name);
return false;
}
ReleaseSysCache(tuple);
return true;
}
/*
* - 简要:确定进度是否处于集群调整状态
* - 返回:
* @TRUE: 正在调整集群大小
* @FALSE: 进度并不在集群调整中
*/
bool ClusterResizingInProgress()
{
Relation pgxc_group_rel = NULL;
TableScanDesc scan;
HeapTuple tup = NULL;
Datum datum;
bool isNull = false;
bool result = false;
pgxc_group_rel = heap_open(PgxcGroupRelationId, AccessShareLock);
if (!pgxc_group_rel) {
ereport(PANIC, (errcode(ERRCODE_RELATION_OPEN_ERROR), errmsg("can not open pgxc_group")));
}
scan = tableam_scan_begin(pgxc_group_rel, SnapshotNow, 0, NULL);
while ((tup = (HeapTuple) tableam_scan_getnexttuple(scan, ForwardScanDirection)) != NULL) {
datum = heap_getattr(tup, Anum_pgxc_group_in_redistribution, RelationGetDescr(pgxc_group_rel), &isNull);
if ('y' == DatumGetChar(datum)) {
result = true;
break;
}
}
tableam_scan_end(scan);
heap_close(pgxc_group_rel, AccessShareLock);
return result;
}
/*
* -简介:获取delete_delta表的名称
* - 参数:
* @relname: 目标表名
* @delta_delta_name: delete_delta表名的输出值
* @isMultiCatchup: 是不是多追赶delta
* - 返回:
* 无返回值
*/
static inline void RelationGetDeleteDeltaTableName(Relation rel, char* delete_delta_name, bool isMultiCatchup)
{
int rc = 0;
/* 检查输出参数是否没有从调用方palloc()-ed */
if (delete_delta_name == NULL || rel == NULL) {
ereport(ERROR,
(errcode(ERRCODE_INVALID_PARAMETER_VALUE), errmsg("Invalid parameter in function '%s'", __FUNCTION__)));
}
/*
* 查找Relation的关联以获得表的id
* 形成delete_delta表的名称
*/
if (!IsInitdb) {
if (RelationInClusterResizing(rel) && !RelationInClusterResizingReadOnly(rel)) {
if (isMultiCatchup) {
rc = snprintf_s(delete_delta_name,
NAMEDATALEN,
NAMEDATALEN - 1,
REDIS_MULTI_CATCHUP_DELETE_DELTA_TABLE_PREFIX "%u",
RelationGetRelCnOid(rel));
} else {
rc = snprintf_s(delete_delta_name,
NAMEDATALEN,
NAMEDATALEN - 1,
REDIS_DELETE_DELTA_TABLE_PREFIX "%u",
RelationGetRelCnOid(rel));
}
securec_check_ss(rc, "\0", "\0");
} else {
elog(LOG, "rel %s doesn't exist in redistributing", RelationGetRelationName(rel));
rc = snprintf_s(delete_delta_name,
NAMEDATALEN,
NAMEDATALEN - 1,
REDIS_DELETE_DELTA_TABLE_PREFIX "%s",
RelationGetRelationName(rel));
securec_check_ss(rc, "\0", "\0");
}
}
return;
}
/*
* - 简介:获取并打开delete_delta rel
* - 参数:
* @rel: UPDATE/DELETE/TRUNCATE操作的目标关系
* @lockmode: 锁定模式
* @isMultiCatchup: 是不是多追赶delta
* - 返回:
* delete_delta rel
*/
Relation GetAndOpenDeleteDeltaRel(const Relation rel, LOCKMODE lockmode, bool isMultiCatchup)
{
Relation deldelta_rel;
Oid deldelta_relid;
char delete_delta_tablename[NAMEDATALEN];
Oid data_redis_namespace;
errno_t errorno;
errorno = memset_s(delete_delta_tablename, NAMEDATALEN, 0, NAMEDATALEN);
securec_check_c(errorno, "\0", "\0");
RelationGetDeleteDeltaTableName(rel, (char*)delete_delta_tablename, isMultiCatchup);
data_redis_namespace = get_namespace_oid("data_redis", false);
/* 我们将在data_redis模式下获取delete delta关系。 */
deldelta_relid = get_relname_relid(delete_delta_tablename, data_redis_namespace);
if (!OidIsValid(deldelta_relid)) {
/*
* 如果多追赶增量表不存在则返回NULL。否则不是 We should not
* 报告错误因为这是一个有效的案例。多追赶delta表是( Multi catchup delta table is)
* 在每次追赶迭代中被丢弃。
*/
if (isMultiCatchup) {
return NULL;
}
/*
* 为了在扩展期间支持更新或删除我们需要添加2列。
* 更多的列。如果表已经包含了太多的列受maxheapattributennumber的限制
* 我们不再允许更新或删除,但插入语句仍然可以进行。
*/
if (((rel->rd_att->natts > (MaxHeapAttributeNumber - (Natts_pg_delete_delta - 1))) &&
!RELATION_IS_PARTITIONED(rel)) ||
((rel->rd_att->natts > (MaxHeapAttributeNumber - Natts_pg_delete_delta)) && RELATION_IS_PARTITIONED(rel))) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("do not support update or delete on table %s, when do cluster resizing on it.",
RelationGetRelationName(rel)),
errdetail("Can not support online extension, if the table contains too many columns")));
}
/* 错误情况下,不应该出现在这里 */
ereport(ERROR,
(errcode(ERRCODE_UNDEFINED_TABLE),
errmsg("delete delta table %s is not found when do cluster resizing table \"%s\"",
delete_delta_tablename,
RelationGetRelationName(rel))));
}
deldelta_rel = relation_open(deldelta_relid, lockmode);
elog(DEBUG1,
"Delete_delta table %s for relation %s being under cluster resizing is valid.",
delete_delta_tablename,
RelationGetRelationName(rel));
return deldelta_rel;
}
/*
* - 简介:检查在线扩展期间的配置在集群调整中阻止不支持的ddl。
* - 参数:
* @rel: DDL的解析树
* -返回:
* 无返回值
*/
void BlockUnsupportedDDL(const Node* parsetree)
{
if (IsInitdb) {
return;
}
Relation rel = NULL;
Oid relid = InvalidOid;
List* relidlist = NULL;
ListCell* lc = NULL;
LOCKMODE lockmode_getrelid = AccessShareLock;
LOCKMODE lockmode_openrel = AccessShareLock;
/*
* 文件之前,请检查是否存在共享缓存无效消息
* relation. 关系。这是需要覆盖的情况下的名称
* 对象之后已删除并重新创建的rel
* 事务开始:如果我们不刷新旧的syscache条目
* 然后我们将锁定该条目并在稍后遭受错误。
*/
AcceptInvalidationMessages();
switch (nodeTag(parsetree)) {
case T_CreatedbStmt:
case T_AlterDatabaseStmt:
case T_AlterDatabaseSetStmt:
case T_CreateTableSpaceStmt:
case T_DropTableSpaceStmt:
case T_AlterTableSpaceOptionsStmt:
case T_CreateGroupStmt: {
if (ClusterResizingInProgress()) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion", CreateCommandTag((Node*)parsetree))));
}
return;
} break;
case T_DropdbStmt: {
if (ClusterResizingInProgress() && !u_sess->attr.attr_common.xc_maintenance_mode) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion", CreateCommandTag((Node*)parsetree))));
} else if (ClusterResizingInProgress()) {
ereport(WARNING, (errmsg("Drop database during online expansion in maintenance mode!")));
}
return;
} break;
/* 在集群调整大小时阻塞游标 */
case T_PlannedStmt: {
PlannedStmt* stmt = (PlannedStmt*)parsetree;
relidlist = stmt->relationOids;
} break;
/* 当表在集群中调整大小时块RENAME */
case T_RenameStmt: {
RenameStmt* stmt = (RenameStmt*)parsetree;
switch (stmt->renameType) {
case OBJECT_SCHEMA: {
Oid nsOid = get_namespace_oid(stmt->subname, true);
TRANSFER_DISABLE_DDL(nsOid);
break;
}
case OBJECT_TABLE: {
if (stmt->relation != NULL) {
Oid relOid = RangeVarGetRelid(stmt->relation, AccessShareLock, true);
if (OidIsValid(relOid)) {
Oid nsOid = GetNamespaceIdbyRelId(relOid);
UnlockRelationOid(relOid, AccessShareLock);
TRANSFER_DISABLE_DDL(nsOid);
}
}
break;
}
default:
break;
}
if (stmt->relation && CheckRangeVarInRedistribution(stmt->relation))
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
stmt->relation->relname)));
} break;
/* 当表在集群中调整大小时Block ALTER设置模式 */
case T_AlterObjectSchemaStmt: {
AlterObjectSchemaStmt* stmt = (AlterObjectSchemaStmt*)parsetree;
/* 在传输时禁用alter table set schema */
if (stmt->relation != NULL) {
Oid relOid = RangeVarGetRelid(stmt->relation, AccessShareLock, true);
if (OidIsValid(relOid)) {
Oid nsOid = GetNamespaceIdbyRelId(relOid);
UnlockRelationOid(relOid, AccessShareLock);
TRANSFER_DISABLE_DDL(nsOid);
if (stmt->newschema) {
nsOid = get_namespace_oid(stmt->newschema, true);
TRANSFER_DISABLE_DDL(nsOid);
}
}
}
if (stmt->relation && CheckRangeVarInRedistribution(stmt->relation))
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
stmt->relation->relname)));
} break;
/* 当表在集群中调整大小时,阻塞创建索引(仅适用于行表) */
case T_IndexStmt: {
IndexStmt* stmt = (IndexStmt*)parsetree;
if (stmt->relation) {
relid = RangeVarGetRelid(stmt->relation, AccessShareLock, true);
if (OidIsValid(relid)) {
Relation relation = relation_open(relid, NoLock);
bool inRedis = RelationIsRowFormat(relation) && RelationInClusterResizing(relation);
relation_close(relation, NoLock);
UnlockRelationOid(relid, AccessShareLock);
if (inRedis) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
stmt->relation->relname)));
}
}
}
} break;
/* 当表在集群中调整大小时块REINDEX(仅适用于行表) */
case T_ReindexStmt: {
ReindexStmt* stmt = (ReindexStmt*)parsetree;
if (stmt->relation) {
relid = RangeVarGetRelid(stmt->relation, AccessShareLock, true);
if (OidIsValid(relid)) {
/* 在锁表之前释放索引锁以避免死锁 */
UnlockRelationOid(relid, AccessShareLock);
Relation relation = relation_open(relid, NoLock);
bool inRedis = false;
if (RelationIsRelation(relation)) {
inRedis = RelationIsRowFormat(relation) && RelationInClusterResizing(relation);
} else if (RelationIsIndex(relation)) {
Oid heapOid = IndexGetRelation(relid, false);
Relation heapRelation = relation_open(heapOid, AccessShareLock);
inRedis = RelationIsRowFormat(heapRelation) && RelationInClusterResizing(heapRelation);
relation_close(heapRelation, AccessShareLock);
}
relation_close(relation, NoLock);
if (inRedis) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
stmt->relation->relname)));
}
}
}
} break;
/* 当表在集群中调整大小时阻塞ALTER-Table */
case T_AlterTableStmt: {
AlterTableStmt* stmt = (AlterTableStmt*)parsetree;
AlterTableCmd* cmd = NULL;
foreach (lc, stmt->cmds) {
cmd = (AlterTableCmd*)lfirst(lc);
switch (cmd->subtype) {
case AT_TruncatePartition: {
/*
* 当目标处于只读状态时,我们不允许截断分区
*在线扩容时的模式
*/
if (stmt->relation) {
relid = RangeVarGetRelid(stmt->relation, lockmode_getrelid, true);
if (OidIsValid(relid)) {
/* 禁止在传输过程中截断分区 */
if (CheckRangeVarInRedistribution(stmt->relation)) {
Oid nsOid = GetNamespaceIdbyRelId(relid);
TRANSFER_DISABLE_DDL(nsOid);
}
rel = relation_open(relid, NoLock);
if (RelationInClusterResizingReadOnly(rel)) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command with '%s' option during online expansion on "
"'%s' because the object is in read only mode.",
CreateCommandTag((Node*)parsetree),
CreateAlterTableCommandTag(cmd->subtype),
RelationGetRelationName(rel))));
}
relation_close(rel, NoLock);
if (u_sess->attr.attr_sql.enable_parallel_ddl)
UnlockRelationOid(relid, lockmode_getrelid);
}
}
} break;
case AT_AddNodeList:
case AT_DeleteNodeList: {
#ifndef ENABLE_MULTIPLE_NODES
ereport(ERROR,
(errmodule(MOD_FUNCTION),
errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupported feature"),
errdetail("The capability is not supported for openGauss."),
errcause("%s is not supported for openGauss",
cmd->subtype == AT_DeleteNodeList ? "DeleteNode" : "AddNode"),
erraction("NA")));
#endif
} break;
case AT_UpdateSliceLike: {
#ifndef ENABLE_MULTIPLE_NODES
ereport(ERROR,
(errmodule(MOD_FUNCTION),
errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupported feature"),
errdetail("The capability is not supported for openGauss."),
errcause("UpdateSliceLike is not supported for openGauss"),
erraction("NA")));
#endif
if (!ClusterResizingInProgress() && IS_PGXC_COORDINATOR) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Cannot alter table update slice like when it is not "
"under redistribution")));
}
} break;
case AT_ResetRelOptions:
case AT_SetRelOptions:
if (cmd->def) {
List* options = (List*)cmd->def;
ListCell* opt = NULL;
DefElem* def = NULL;
foreach (opt, options) {
def = (DefElem*)lfirst(opt);
if (pg_strcasecmp(def->defname, "append_mode") == 0) {
break;
}
}
/* 如果rel选项包含append_mode则不检查。 */
if (opt != NULL) {
break;
}
}
/* 失败 */
default: {
if (stmt->relation && !u_sess->attr.attr_sql.enable_cluster_resize &&
CheckRangeVarInRedistribution(stmt->relation))
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command with '%s' option during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
CreateAlterTableCommandTag(cmd->subtype),
stmt->relation->relname)));
} break;
}
}
return;
} break;
/* 当集群中的目标表调整大小时阻塞CREATE-RULE语句 */
case T_RuleStmt: {
RuleStmt* stmt = (RuleStmt*)parsetree;
if (stmt->relation) {
relid = RangeVarGetRelid(stmt->relation, lockmode_getrelid, true);
relidlist = list_make1_oid(relid);
}
} break;
/* 当所有者表在集群中调整大小时Block CREATE SEQUENCE设置模式 */
case T_CreateSeqStmt: {
CreateSeqStmt* stmt = (CreateSeqStmt*)parsetree;
List* owned_by = NULL;
DefElem* defel = NULL;
foreach (lc, stmt->options) {
defel = (DefElem*)lfirst(lc);
if (pg_strcasecmp(defel->defname, "owned_by") == 0 && nodeTag(defel->arg) == T_List) {
owned_by = (List*)defel->arg;
}
}
if (owned_by != NULL) {
int owned_len = list_length(owned_by);
if (owned_len != 1) {
List* relname = list_truncate(list_copy(owned_by), owned_len - 1);
RangeVar* r = makeRangeVarFromNameList(relname);
if (r && CheckRangeVarInRedistribution(r))
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
r->relname)));
}
}
} break;
/* 当集群中的所有者表调整大小时阻塞ALTER SEQUENCE */
case T_AlterSeqStmt: {
AlterSeqStmt* stmt = (AlterSeqStmt*)parsetree;
List* owned_by = NIL;
DefElem* defel = NULL;
foreach (lc, stmt->options) {
defel = (DefElem*)lfirst(lc);
if (pg_strcasecmp(defel->defname, "owned_by") == 0 && nodeTag(defel->arg) == T_List) {
owned_by = (List*)defel->arg;
}
}
if (owned_by != NULL) {
int owned_len = list_length(owned_by);
if (owned_len != 1) {
List* relname = list_truncate(list_copy(owned_by), owned_len - 1);
RangeVar* r = makeRangeVarFromNameList(relname);
if (r && CheckRangeVarInRedistribution(r))
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
r->relname)));
}
}
} break;
/* 当表在集群中调整大小时阻塞集群 */
case T_ClusterStmt: {
ClusterStmt* stmt = (ClusterStmt*)parsetree;
if (stmt->relation && CheckRangeVarInRedistribution(stmt->relation))
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
stmt->relation->relname)));
} break;
/* 当表在集群中调整大小时,块真空已满 */
case T_VacuumStmt: {
VacuumStmt* stmt = (VacuumStmt*)parsetree;
if ((stmt->options & VACOPT_VACUUM) || (stmt->options & VACOPT_MERGE)) {
if (stmt->relation) {
relid = RangeVarGetRelid(stmt->relation, lockmode_getrelid, true);
relidlist = list_make1_oid(relid);
}
}
if (stmt->options & VACOPT_FULL) {
if (OidIsValid(relid)) {
rel = relation_open(relid, lockmode_openrel);
if (RelationInClusterResizing(rel)) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport 'VACUUM FULL' command during online expansion on '%s'",
RelationGetRelationName(rel))));
}
relation_close(rel, lockmode_openrel);
}
}
} break;
/* 在集群调整大小时当目标表为只读时块截断DDL */
case T_TruncateStmt: {
ListCell* cell = NULL;
TruncateStmt* stmt = (TruncateStmt*)parsetree;
foreach (cell, stmt->relations) {
RangeVar* rv = (RangeVar*)lfirst(cell);
if (CheckRangeVarInRedistribution(rv)) {
Oid relOid = RangeVarGetRelid(rv, lockmode_getrelid, true);
if (OidIsValid(relOid)) {
Oid nsOid = GetNamespaceIdbyRelId(relOid);
UnlockRelationOid(relOid, lockmode_getrelid);
TRANSFER_DISABLE_DDL(nsOid);
}
}
rel = heap_openrv(rv, lockmode_openrel);
if (RelationInClusterResizingReadOnly(rel)) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s' because the object is in "
"read only mode.",
CreateCommandTag((Node*)parsetree),
RelationGetRelationName(rel))));
}
relation_close(rel, lockmode_openrel);
}
} break;
case T_DropStmt: {
DropStmt* stmt = (DropStmt*)parsetree;
switch (stmt->removeType) {
case OBJECT_TABLE: {
/* 在传输时禁用drop表 */
ListCell* cell = NULL;
foreach (cell, stmt->objects) {
RangeVar* rel = makeRangeVarFromNameList((List*)lfirst(cell));
Oid relOid = RangeVarGetRelid(rel, AccessShareLock, true);
if (OidIsValid(relOid)) {
Oid nsOid = GetNamespaceIdbyRelId(relOid);
UnlockRelationOid(relOid, AccessShareLock);
TRANSFER_DISABLE_DDL(nsOid);
}
}
break;
}
case OBJECT_SCHEMA: {
/* 传输时禁用删除模式 */
ListCell* cell = NULL;
foreach (cell, stmt->objects) {
List* objname = (List*)lfirst(cell);
char* name = NameListToString(objname);
Oid nsOid = get_namespace_oid(name, true);
TRANSFER_DISABLE_DDL(nsOid);
}
break;
}
default:
break;
}
} break;
case T_CreateStmt: {
/* 禁止传输时创建表 */
CreateStmt* stmt = (CreateStmt*)parsetree;
if (stmt->relation != NULL) {
Oid nsOid = RangeVarGetCreationNamespace(stmt->relation);
TRANSFER_DISABLE_DDL(nsOid);
}
} break;
default:
break;
}
foreach (lc, relidlist) {
relid = lfirst_oid(lc);
if (OidIsValid(relid) && get_rel_name(relid) != NULL) {
rel = relation_open(relid, lockmode_openrel);
if (RelationInClusterResizing(rel)) {
ereport(ERROR,
(errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
errmsg("Unsupport '%s' command during online expansion on '%s'",
CreateCommandTag((Node*)parsetree),
RelationGetRelationName(rel))));
}
relation_close(rel, lockmode_openrel);
}
}
}
/*
* - 简介:对于在线扩展,这里评估的是可发布功能模块
* 在优化器中调用FQS评估时我们必须定义函数
* 是稳定的as STABLE
* - 参数:
* @funcid: 在gs_redis范围内创建/删除的用户定义函数的Oid
* - 返回:
* @true: 可交付
* @false: 不可交付
*/
bool redis_func_shippable(Oid funcid)
{
const char* func_name = get_func_name(funcid);
Oid* argstype = NULL;
int nargs;
Oid rettype = InvalidOid;
bool result = false;
if (func_name == NULL) {
ereport(ERROR, (errcode(ERRCODE_UNDEFINED_FUNCTION), errmsg("function with OID %u does not exist", funcid)));
}
/* 获取函数签名 */
rettype = get_func_signature(funcid, &argstype, &nargs);
if (redis_tupleid_retrive_function(func_name, rettype, argstype, nargs)) {
/* Tupleid检索函数可以发布到数据节点 */
result = true;
} else if (redis_offset_retrive_function(func_name, rettype, argstype, nargs)) {
result = true;
} else if (redis_blocknum_retrive_function(func_name, rettype, argstype, nargs)) {
result = true;
} else if (redis_ctid_retrive_function(func_name, rettype, argstype, nargs)) {
result = true;
}
/* pfree */
if (argstype != NULL) {
pfree_ext(argstype);
argstype = NULL;
}
return result;
}
/*
* - 简介:确定给定的函数是否反映了一个非稳定函数
* - 参数:
* @funcid: 要求值的函数oid
* - 返回:
* @result: true:不稳定的 false: 不稳定的函数
*/
bool redis_func_dnstable(Oid funcid)
{
const char* func_name = get_func_name(funcid);
Oid* argstype = NULL;
int nargs;
Oid rettype = InvalidOid;
bool result = false;
if (func_name == NULL) {
ereport(ERROR,
(errcode(ERRCODE_UNDEFINED_FUNCTION),
errmsg("function with OID %u does not exist when checking function dnstable", funcid)));
}
/* 获取函数签名 */
rettype = get_func_signature(funcid, &argstype, &nargs);
if (redis_tupleid_retrive_function(func_name, rettype, argstype, nargs)) {
/* 管状反射函数是不稳定的 */
result = true;
}
return result;
}
/*
* - 简介:将ctid函数求值为const值以避免每次扫描
* 在seqscan中调用元组。
* - 参数:
* @rel: 真正的问题是再分配
* @original_quals: 原始的quals可能包含ctid_funcs
* @isRangeScanInRedis: 这是一个redis范围扫描
* - 返回:
* @new_quals: 函数调用的Quals将被const替换
*/
List* eval_ctid_funcs(Relation rel, List* original_quals, RangeScanInRedis *rangeScanInRedis)
{
StringInfo qual_str = makeStringInfo();
/*
* 由于eval_dnstable_func_mutator的存在我们必须对原始的quals进行复制
* 将修改它。在以后的时间里,将会一次又一次地需要原始的质量
* 要在分区表扫描中重新计算。
*/
List* new_quals = (List*)copyObject((const void*)(original_quals));
rangeScanInRedis->isRangeScanInRedis = false;
rangeScanInRedis->sliceTotal = 0;
rangeScanInRedis->sliceIndex = 0;
(void)eval_dnstable_func_mutator(rel, (Node*)new_quals, qual_str, rangeScanInRedis, true);
pfree_ext(qual_str->data);
pfree_ext(qual_str);
return new_quals;
}
static int32 get_expr_const_val(Node *val){
if (IsA(val, Const) && !((Const*)val)->constisnull && ((Const*)val)->consttype == INT4OID) {
return DatumGetInt32(((Const*)val)->constvalue);
} else {
return 0;
}
}
/*
* - 简介:eval_dnstable_func()的工作库用于将一个稳定函数求值为const
* 值以避免在seqscan中调用每次扫描的元组
* - 参数:
* @rel: 真正的问题是再分配
* @node: 表达式节点
* @qual_str: 谓词模式
* @isRangeScanInRedis: 输出以指示谓词模式是否为redis中的范围扫描
* @isRoot: 我们只想在根级别对谓词模式进行一次比较
* - 返回:
* @result: 表达式树与dn稳定函数const评估
*/
static Node* eval_dnstable_func_mutator(
Relation rel, Node* node, StringInfo qual_str, RangeScanInRedis *rangeScanInRedis, bool isRoot)
{
if (node == NULL)
return NULL;
if (IS_PGXC_COORDINATOR)
return node;
switch (nodeTag(node)) {
case T_FuncExpr: {
FuncExpr* expr = (FuncExpr*)node;
/* 将一个稳定函数扁平化为const值 */
if (redis_func_dnstable(expr->funcid)) {
Node* new_const = NULL;
char* funcname = get_func_name(expr->funcid);
if (funcname == NULL) {
ereport(ERROR,
(errcode(ERRCODE_UNDEFINED_FUNCTION),
errmsg("operation expression function with OID %u does not exist.", expr->funcid)));
}
bool is_func_get_start_ctid = pg_strcasecmp(funcname, "pg_get_redis_rel_start_ctid") == 0;
bool is_func_get_end_ctid = pg_strcasecmp(funcname, "pg_get_redis_rel_end_ctid") == 0;
if (is_func_get_start_ctid || is_func_get_end_ctid){
int32 numSlices = get_expr_const_val((Node*)list_nth(expr->args, 2));
int32 idxSlices = get_expr_const_val((Node*)list_nth(expr->args, 3));
new_const = eval_redis_func_direct(rel, is_func_get_start_ctid, numSlices, idxSlices);
rangeScanInRedis->sliceIndex = idxSlices;
rangeScanInRedis->sliceTotal = numSlices;
} else {
new_const = eval_const_expressions(NULL, node);
}
appendStringInfoString(qual_str, get_func_name(expr->funcid));
appendStringInfoString(qual_str, "+");
return new_const;
}
break;
}
case T_List: {
List* l = (List*)node;
for (int i = 0; i < list_length(l); i++) {
Node* expr = (Node*)list_nth(l, i);
Node* new_expr = eval_dnstable_func_mutator(rel, expr, qual_str, rangeScanInRedis, false);
/*
* 如果将FuncExpr节点求值为T_Const值则为命中
* 将点替换为等号列表。
*/
if (expr && IsA(expr, FuncExpr) && new_expr && IsA(new_expr, Const)) {
l = list_delete_ptr(l, expr);
l = lappend(l, new_expr);
}
}
/*
* 如果在根的谓词类似于“where ctid between pg_get_redis_rel_start_ctid('xx')”
* 和pg_get_redis_rel_end_ctid('xx')"在DN上我们将在扫描节点下推谓词。
*/
if (isRoot && pg_strcasecmp(qual_str->data, RANGE_SCAN_IN_REDIS) == 0) {
rangeScanInRedis->isRangeScanInRedis = true;
}
break;
}
case T_OpExpr: {
OpExpr* opexpr = (OpExpr*)node;
char* funcname = get_func_name(opexpr->opfuncid);
if (funcname == NULL) {
ereport(ERROR,
(errcode(ERRCODE_UNDEFINED_FUNCTION),
errmsg("operation expression function with OID %u does not exist.", opexpr->opfuncid)));
}
appendStringInfoString(qual_str, funcname);
appendStringInfoString(qual_str, "+");
eval_dnstable_func_mutator(rel, (Node*)opexpr->args, qual_str, rangeScanInRedis, false);
break;
}
case T_Var: {
Var* var = (Var*)node;
/* 我们只期望谓词中有tid列 */
if (var->vartype == TIDOID) {
appendStringInfoString(qual_str, "tid");
appendStringInfoString(qual_str, "+");
}
break;
}
default: {
appendStringInfoString(qual_str, nodeTagToString(nodeTag(node)));
appendStringInfoString(qual_str, "+");
break;
}
}
return NULL;
}
/*
* - 简介:获取并打开new_table rel
* - 参数:
* @rel: TRUNCATE操作的目标关系
* - 返回:
* new_table rel
*/
Relation GetAndOpenNewTableRel(const Relation rel, LOCKMODE lockmode)
{
Relation newtable_rel = NULL;
Oid newtable_relid = InvalidOid;
Oid data_redis_namespace;
char new_tablename[NAMEDATALEN];
errno_t errorno = EOK;
errorno = memset_s(new_tablename, NAMEDATALEN, 0, NAMEDATALEN);
securec_check_c(errorno, "\0", "\0");
RelationGetNewTableName(rel, (char*)new_tablename);
data_redis_namespace = get_namespace_oid("data_redis", false);
newtable_relid = get_relname_relid(new_tablename, data_redis_namespace);
if (!OidIsValid(newtable_relid)) {
/* 错误情况下,不应该出现在这里 */
ereport(ERROR,
(errcode(ERRCODE_DATA_EXCEPTION),
errmsg("new table %s is not found when do cluster resizing table \"%s\"",
new_tablename,
RelationGetRelationName(rel))));
}
newtable_rel = relation_open(newtable_relid, lockmode);
elog(LOG,
"New temp table %s for relation %s under cluster resizing is valid.",
new_tablename,
RelationGetRelationName(rel));
return newtable_rel;
}
/*
* - 简介:获得新表的名称
* - 参数:
* @relname: 目标表名
* @newtable_name: 新表名的输出值
* - 返回:
* 无返回值
*/
void RelationGetNewTableName(Relation rel, char* newtable_name)
{
int rc = 0;
/* 检查输出参数是否没有从调用方palloc()-ed */
if (newtable_name == NULL || rel == NULL) {
ereport(ERROR,
(errcode(ERRCODE_INVALID_PARAMETER_VALUE),
errmsg("Invalid parameter in function '%s' when getting the name of new table", __FUNCTION__)));
}
/*
* 查找关系的关联以获得表的关联
* 形成新表的名称
*/
if (!IsInitdb) {
Oid rel_cn_oid = RelationGetRelCnOid(rel);
if (OidIsValid(rel_cn_oid)) {
rc = snprintf_s(newtable_name, NAMEDATALEN, NAMEDATALEN - 1, "data_redis_tmp_%u", rel_cn_oid);
} else {
elog(LOG, "rel %s doesn't exist in redistributing", RelationGetRelationName(rel));
rc = snprintf_s(
newtable_name, NAMEDATALEN, NAMEDATALEN - 1, "data_redis_tmp_%s", RelationGetRelationName(rel));
}
/* 检查安全函数的返回值 */
securec_check_ss(rc, "\0", "\0");
}
return;
}
/*
* - 简介:确定关系是否处于群集调整大小写错误模式
* - 参数:
* @rel: 需要检查的关系
* - 参数:
* @TRUE: 关系处于群集调整大小写错误模式
* @FALSE: 关系不在群集调整大小写错误模式下
*/
bool RelationInClusterResizingWriteErrorMode(const Relation rel)
{
return RelationInClusterResizingReadOnly(rel) ||
(RelationInClusterResizingEndCatchup(rel) && !pg_try_advisory_lock_for_redis(rel));
}