PostgreSQL 9.5双Schema同步需求:OLTP事务库操作历史同步至OLAP数仓
实现PostgreSQL 9.5中OLTP Schema到OLAP Schema的操作同步
嘿,针对你在PostgreSQL 9.5环境下,需要把OLTP事务Schema的所有操作历史同步到OLAP数据仓库Schema的需求,我整理了几个适配这个版本的靠谱方案,你可以根据自己的场景选择:
方案一:触发器+自定义函数实现实时同步
这是最直接的实时同步方式,适合要求操作一发生就同步到OLAP的场景:
- 第一步,给OLTP里的目标表(比如
oltp_schema.employee_address)写个触发器函数,用来捕获INSERT/UPDATE/DELETE操作的数据,然后同步到OLAP的对应表:
CREATE OR REPLACE FUNCTION sync_employee_address() RETURNS TRIGGER AS $$ BEGIN IF (TG_OP = 'INSERT') THEN INSERT INTO olap_schema.employee_address SELECT NEW.*; RETURN NEW; ELSIF (TG_OP = 'UPDATE') THEN UPDATE olap_schema.employee_address SET address = NEW.address, updated_at = NEW.updated_at WHERE employee_id = NEW.employee_id; RETURN NEW; ELSIF (TG_OP = 'DELETE') THEN DELETE FROM olap_schema.employee_address WHERE employee_id = OLD.employee_id; RETURN OLD; END IF; END; $$ LANGUAGE plpgsql;
- 第二步,把这个函数绑定到目标表的触发器上:
CREATE TRIGGER trigger_sync_employee_address AFTER INSERT OR UPDATE OR DELETE ON oltp_schema.employee_address FOR EACH ROW EXECUTE FUNCTION sync_employee_address();
- 小提醒:要确保OLAP Schema的表结构和OLTP的完全匹配;另外触发器会给OLTP的操作增加一点延迟,高并发场景下最好先做性能测试。
方案二:逻辑复制(Logical Replication)
PostgreSQL 9.5已经原生支持逻辑复制了,这种方式适合需要全量+增量同步,且不想给OLTP带来太多性能负担的场景:
- 第一步,先把数据库的
wal_level设置为logical,修改postgresql.conf后重启数据库生效:
wal_level = logical
- 第二步,在OLTP的目标表上创建发布(Publication):
CREATE PUBLICATION oltp_publication FOR TABLE oltp_schema.employee_address;
- 第三步,在OLAP这边创建订阅(Subscription),指向刚才的发布:
CREATE SUBSCRIPTION olap_subscription CONNECTION 'dbname=your_database user=your_username host=localhost' PUBLICATION oltp_publication WITH (copy_data = true); -- copy_data=true会先同步全量数据,之后自动同步增量操作
- 小提醒:逻辑复制要求同步的表必须有主键;你可以灵活选择要同步的表,还能过滤特定操作;相比触发器,它对OLTP的性能影响小很多,适合大规模数据同步。
方案三:定时任务实现准实时同步
如果对同步延迟要求不高(比如允许几分钟甚至几小时的延迟),可以用定时任务的方式:
- 第一步,写个同步SQL脚本,用
INSERT ... ON CONFLICT来处理新增和更新,同时处理删除操作:
-- 同步员工地址表:不存在则插入,存在则更新 INSERT INTO olap_schema.employee_address (employee_id, address, created_at, updated_at) SELECT employee_id, address, created_at, updated_at FROM oltp_schema.employee_address ON CONFLICT (employee_id) DO UPDATE SET address = EXCLUDED.address, updated_at = EXCLUDED.updated_at; -- 删除OLAP中已在OLTP里移除的记录 DELETE FROM olap_schema.employee_address oa WHERE NOT EXISTS ( SELECT 1 FROM oltp_schema.employee_address oe WHERE oe.employee_id = oa.employee_id );
- 第二步,用Cron或者PostgreSQL的
pg_cron扩展定时执行这个脚本。比如用pg_cron设置每分钟同步一次:
SELECT cron.schedule('sync_employee_address_job', '* * * * *', '上述完整的同步SQL语句');
- 小提醒:这种方式适合数据量不大、对实时性要求不高的场景;要注意处理数据冲突,避免重复同步或者数据丢失。
内容的提问来源于stack exchange,提问作者Harinath Arasu
相关产品推荐
相关产品推荐

