求助:使用WSO2 Streaming Integrator连接PostgreSQL并开启CDC监听模式构建POC
WSO2 Streaming Integrator 连接PostgreSQL CDC监听POC实现步骤
前置准备
- 确认PostgreSQL版本≥9.4(支持逻辑复制),WSO2 Streaming Integrator(SI)已安装完成
- 在PostgreSQL中创建测试库和示例表:
CREATE DATABASE cdc_test; CREATE TABLE cdc_test.public.users ( id SERIAL PRIMARY KEY, name VARCHAR(50), email VARCHAR(100) UNIQUE );
配置PostgreSQL开启CDC
修改
postgresql.conf文件,设置以下参数:wal_level = logical max_replication_slots = 10 max_wal_senders = 10 wal_sender_timeout = 60000保存后重启PostgreSQL服务。
创建具备复制权限的CDC用户:
CREATE USER cdc_user WITH REPLICATION LOGIN PASSWORD 'cdc_pass123';修改
pg_hba.conf,添加WSO2 SI节点的访问权限:host replication cdc_user <WSO2_SI_IP>/32 md5重启PostgreSQL使配置生效。
为目标表创建发布者:
CREATE PUBLICATION cdc_pub FOR TABLE cdc_test.public.users;
配置WSO2 SI的PostgreSQL驱动
- 下载对应版本的PostgreSQL JDBC驱动(如
postgresql-42.6.0.jar),将其放入<SI_HOME>/lib目录 - 重启WSO2 SI服务,确保驱动加载成功
编写Siddhi CDC流处理脚本
创建postgresql_cdc_poc.siddhi文件,内容如下:
@App:name("PostgreSQLCDCPOC") @App:description("Capture PostgreSQL CDC events for users table") define stream UserChangeStream (id int, name string, email string, operation string); @source(type='cdc', url='jdbc:postgresql://<PG_HOST>:5432/cdc_test', username='cdc_user', password='cdc_pass123', table.name='public.users', operation='insert,update,delete', @map(type='keyvalue')) define stream CDCSourceStream (id int, name string, email string, operation string); @sink(type='log', prefix='[CDC EVENT]') define stream LogOutputStream (id int, name string, email string, operation string); from CDCSourceStream select id, name, email, operation insert into LogOutputStream;
部署与验证
- 将编写好的Siddhi脚本上传至WSO2 SI Dashboard的Siddhi Apps页面,或直接放入
<SI_HOME>/deployment/siddhi-files目录 - 启动该Siddhi应用
- 在PostgreSQL中对
users表执行增、删、改操作:INSERT INTO users(name, email) VALUES('Alice', 'alice@test.com'); UPDATE users SET email='alice_new@test.com' WHERE id=1; DELETE FROM users WHERE id=1; - 查看WSO2 SI的日志文件(
<SI_HOME>/logs/wso2carbon.log),若出现类似以下输出,说明CDC事件已被成功捕获:[CDC EVENT] Event{timestamp=1699999999999, data=[1, Alice, alice@test.com, insert], isExpired=false} [CDC EVENT] Event{timestamp=1700000000000, data=[1, Alice, alice_new@test.com, update], isExpired=false} [CDC EVENT] Event{timestamp=1700000000001, data=[1, Alice, alice_new@test.com, delete], isExpired=false}
注意事项
- 确保PostgreSQL与WSO2 SI之间网络连通,5432端口开放
- WSO2 SI会自动创建PostgreSQL复制槽,若应用停止后需清理,可执行:
SELECT pg_drop_replication_slot('<slot_name>');(槽名通常为siddhi_cdc_<table_name>) - 若需监听多个表,可修改
table.name为逗号分隔的表名,或调整CREATE PUBLICATION包含多个表
内容的提问来源于stack exchange,提问作者aditi21
相关产品推荐
相关产品推荐

