国测 vexdb
-- cn: 50000 -- dn1: 50005 -- dn2: 50010 vsql -d postgres -p 50000 -r SELECT * FROM PG_DIST_NODE; CREATE TABLE t1(c1 INT, c2 TEXT); SELECT create_distributed_table('t1', 'c1'); -- 创建8个分片,dn1和dn2,每个dn有4个分片,即4个表,表名为 t1_$oid SELECT * FROM pg_dist_shard; -- min和max是hash之后的值 SELECT * FROM pg_dist_shard_placement; SELECT * FROM pg_dist_partition; INSERT INTO t1 VALUES(1, 'aa'), (3, 'bb'), (7, 'cc'), (10, 'dd'); vsql -d postgres -p 50005 -r SELECT oid,relname FROM pg_class WHERE relname like '%t1%'; vsql -d postgres -p 50010 -r SELECT oid,relname FROM pg_class WHERE relname like '%t1%'; # cn SELECT create_distributed_table('t1', 'c1'); # --> shard 1 (dn1) SELECT worker_apply_shard_ddl_command (12042, 'CREATE TABLE public.t1 (c1 integer, c2 text) WITH (orientation=row, compression=no, fillfactor=80)'); # dn1 worker_apply_shard_ddl_command ddlnode = ParseTreeNode(sql) RelayEventExtendNames(ddlnode, sharid) # 在语法解析树中,把t1名称,修改为t1_$shardid ProcessUtilityParseTree SELECT worker_apply_shard_ddl_command (12042, 'ALTER TABLE public.t1 OWNER TO shenkun'); # --> shard 2 (dn2) SELECT worker_apply_shard_ddl_command (12043, 'CREATE TABLE public.t1 (c1 integer, c2 text) WITH (orientation=row, compression=no, fillfactor=80)'); SELECT worker_apply_shard_ddl_command (12043, 'ALTER TABLE public.t1 OWNER TO shenkun'); # --> shared 3 (dn1) ... # --> 一共8个分片,每个分片2条SQL! # cn INSERT INTO t1 VALUES(1,'aa'),(3,'bb'),(7,'cc'),(10,'dd'); exec_simple_query pg_parse_query pg_analyze_and_rewrite pg_plan_query [planner hook] distributed_planner CreateDistributedPlan CreateModifyPlan RouterInsertJob (deferredPruning=true) # 返回 CustomScan PortalStart ExecutorStart [DVecExecutorStart] DVecExecutorStart BeginCustomScan → DVecBeginScan DVecBeginModifyScan RegenerateTaskListForInsert BuildRoutesForInsert FindShardInterval ActiveShardPlacementList RebuildQueryStrings DeparseTaskQuery AppendShardIdToName # 生成改名后SQL PortalRun ExecutorRun [ExecutorRun_hook] DVecExecutorRun ExecCustomScan DVecExecScan AdaptiveExecutor RunDistributedExecution ... SendRemoteCommand StartRemoteTransactionBegin BeginTransactiionCommand # 生成:BEGIN ISOLATION LEVEL READ COMMITTED; AssignDistributedTransactionIdCommand # 生成:SELECT assign_distributed_transaction_id(0,20,'...'); SendRemoteCommand # 一次发送2条事务SQL SendRemoteCommand PQsendQuery # INTO t1_$shardid .. CoordinatedRemoteTransactionsPrepare .. SendRemoteCommand # 发送 REPARE TRANSACTION .. SendRemoteCommand # 发送 COMMIT PREPARED .. # --> shard 1 (dn1) BEGIN ISOLATION LEVEL READ COMMITTED; SELECT assign_distributed_transaction_id(0,20,'...'); INSERT INTO public.t1_12042 AS dvec_table_alias (c1, c2) VALUES (1,'aa'::text) REPARE TRANSACTION 'dvecx_...'; COMMIT PREPARED 'dvecx_...'; # 执行分布式事务 INSERT INTO public.t1_12042 AS dvec_table_alias (c1, c2) VALUES (1,'aa'::text) # --> shard 2 (dn2) INSERT INTO public.t1_12043 AS dvec_table_alias (c1, c2) VALUES (10,'dd'::text) # --> shared 3 (d1) ... # --> 一共8个分片,其中4个分片,各收到1条INSERT语句 # 多行在同一分片,INSERT语句合并 # cn SELECT * FROM t1 WHERE c1 > 10; exec_simple_query pg_parse_query pg_analyze_and_rewrite pg_plan_query [planner hook] distributed_planner CreateDistributedPlan GetRouterPlanType (一造 multi-shard fan-out(非 router 计划)) MultiLogicalPlanCreate (logical planner:建逻辑计划) CreatePhysicalDistributedPlan (physical planner) ├─ 为命中的每个分片建一个 Task(本表 8 个分片 → 8 个 Task) UpdateTaskQueryString DeparseTaskQuery deparse_shard_query AppendShardIdToName # t1 → t1_<shardId> # SELECT * + c1,c2 ; WHERE c1>10 → (c1 OPERATOR(pg_catalog.>) 10) # 包成 CustomScan 返回 PortalStart ExecutorStart [DVecExecutorStart] DVecExecutorStart BeginCustomScan → DVecBeginScan PortalRun ExecutorRun [ExecutorRun_hook] DVecExecutorRun ExecCustomScan DVecExecScan AdaptiveExecutor RunDistributedExecution AssignTasksToConnectionsOrWorkerPool # 8 个 Task 分给 DN1/DN2 的连接 ConnectionStateMachine StartRemoteTransactionBegin *发送: BEGIN + assign(每条 DN 连接首次用,只发一次) BeginTransactionCommand # 生成: BEGIN TRANSACTION ISOLATION LEVEL READ COMMITTED AssignDistributedTransactionIdCommand # 生成: SELECT assign_distributed_transaction_id(...) SendRemoteCommand → PQsendQuery *两条合并成一条发出 # 发送: SELECT c1,c2 FROM t1_<shardId> WHERE (c1>10) # DN1 上 4 条 (12042/12044/12046/12048), DN2 上 4 条 (COMMIT) RemoteTransactionCommit SendRemoteCommand → PQsendQuery # 只读 → 普通 COMMIT(无 PREPARE/COMMIT PREPARED,无 2PC) # COMMIT # --> shard 1 (dn1) SELECT c1, c2 FROM public.t1_12042 t1 WHERE (c2 OPERATOR(pg_catalog.>) 'cc'::text) SELECT * FROM t1 WHERE c2 > 'cc'; 密态: cn CREATE TABLE t1 (c1 ..) INSERT INTO t1 VALUES (..) SELECT c1,c2 FROM t1 WHERE .. dn CREATE TABLE t1_$shardid (c1 ..) INSERT INTO t1_$shardid VALUES (..) # values的值不变 SELECT c1,c2 FROM t1_$shardid WHERE .. # WHERE子句不变