如何在Postgres的AGE顶点表与普通行表间同步数据?
解决方案:通过PostgreSQL触发器同步关系表与Apache AGE顶点
前提假设
- 关系表名为
site_rel,包含字段:id,name(唯一约束),type,full_cost,part_cost - AGE图名为
my_graph,顶点标签为Site
1. 编写同步触发器函数
插入同步函数
当关系表新增记录时,自动在AGE中创建对应顶点:
CREATE OR REPLACE FUNCTION sync_site_insert() RETURNS TRIGGER AS $$ BEGIN PERFORM ag_catalog.cypher( 'my_graph', $$ CREATE (s:Site {name: $name, type: $type, full_cost: $full_cost, part_cost: $part_cost}) $$, ag_catalog.map( array['name', 'type', 'full_cost', 'part_cost'], array[NEW.name::text, NEW.type::text, NEW.full_cost::text, NEW.part_cost::text] ) ); RETURN NEW; END; $$ LANGUAGE plpgsql;
更新同步函数
当关系表记录更新时,匹配name(唯一键)更新对应AGE顶点的属性:
CREATE OR REPLACE FUNCTION sync_site_update() RETURNS TRIGGER AS $$ BEGIN PERFORM ag_catalog.cypher( 'my_graph', $$ MATCH (s:Site {name: $old_name}) SET s.name = $new_name, s.type = $type, s.full_cost = $full_cost, s.part_cost = $part_cost $$, ag_catalog.map( array['old_name', 'new_name', 'type', 'full_cost', 'part_cost'], array[OLD.name::text, NEW.name::text, NEW.type::text, NEW.full_cost::text, NEW.part_cost::text] ) ); RETURN NEW; END; $$ LANGUAGE plpgsql;
删除同步函数
当关系表记录删除时,删除对应AGE顶点:
CREATE OR REPLACE FUNCTION sync_site_delete() RETURNS TRIGGER AS $$ BEGIN PERFORM ag_catalog.cypher( 'my_graph', $$ MATCH (s:Site {name: $name}) DELETE s $$, ag_catalog.map( array['name'], array[OLD.name::text] ) ); RETURN OLD; END; $$ LANGUAGE plpgsql;
2. 创建触发器绑定到关系表
将上述函数关联到site_rel表的增删改操作:
-- 插入触发器 CREATE TRIGGER trigger_site_insert AFTER INSERT ON site_rel FOR EACH ROW EXECUTE FUNCTION sync_site_insert(); -- 更新触发器 CREATE TRIGGER trigger_site_update AFTER UPDATE ON site_rel FOR EACH ROW EXECUTE FUNCTION sync_site_update(); -- 删除触发器 CREATE TRIGGER trigger_site_delete AFTER DELETE ON site_rel FOR EACH ROW EXECUTE FUNCTION sync_site_delete();
3. 关键注意事项
- 唯一键依赖:必须保证
site_rel表的name字段有唯一约束,否则同步时可能出现匹配错误或重复顶点。 - 事务一致性:触发器与原表操作在同一事务中,任何一步失败都会导致整个操作回滚,确保两份数据的一致性。
- 数据类型适配:示例中把数值型字段转为文本存入AGE属性,若需要保留数值计算能力,可去掉
::text转换,直接传入数值类型。 - 初始数据批量同步:如果是首次配置,需要先把关系表现有数据导入AGE,执行以下语句:
PERFORM ag_catalog.cypher( 'my_graph', $$ UNWIND $sites AS site CREATE (s:Site {name: site.name, type: site.type, full_cost: site.full_cost, part_cost: site.part_cost}) $$, ag_catalog.map( array['sites'], array[(SELECT json_agg(row_to_json(s)) FROM site_rel s)] ) );
内容的提问来源于stack exchange,提问作者PurpleHaze_Mustang
相关产品推荐
相关产品推荐

