You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求助:使用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

  1. 修改postgresql.conf文件,设置以下参数:

    wal_level = logical
    max_replication_slots = 10
    max_wal_senders = 10
    wal_sender_timeout = 60000
    

    保存后重启PostgreSQL服务。

  2. 创建具备复制权限的CDC用户:

    CREATE USER cdc_user WITH REPLICATION LOGIN PASSWORD 'cdc_pass123';
    
  3. 修改pg_hba.conf,添加WSO2 SI节点的访问权限:

    host    replication     cdc_user     <WSO2_SI_IP>/32     md5
    

    重启PostgreSQL使配置生效。

  4. 为目标表创建发布者:

    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;

部署与验证

  1. 将编写好的Siddhi脚本上传至WSO2 SI Dashboard的Siddhi Apps页面,或直接放入<SI_HOME>/deployment/siddhi-files目录
  2. 启动该Siddhi应用
  3. 在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;
    
  4. 查看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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 23:36:20