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

Siddhi CDC PostgreSQL应用未输出数据日志问题咨询

问题描述

我创建了用于捕获PostgreSQL数据库表数据变更的Siddhi CDC应用,配置如下:

@App:name('post')
@source(type = 'cdc' ,url = 'jdbc:postgresql://postgres:5432/shipment_db',
username = 'postgresuser', password = 'postgrespw',
table.name = 'public.shipments', operation = 'insert',  plugin.name='pgoutput',slot.name='postslot',
@map(type='keyvalue', @attributes(shipment_id = 'shipment_id', order_id = 'order_id',date_created='date_created',status='status')))
define stream inputStream (shipment_id long, order_id long,date_created string, status string);

@sink(type = 'log')
define stream OutputStream (shipment_id long, date_created string);

@info(name = 'query1')
from inputStream
select shipment_id, date_created
insert into OutputStream;

我将siddhi-io-cdc-2.0.12.jar、siddhi-core-5.1.21.jar放置于./files/bundles目录,org.wso2.carbon.si.metrics.core-3.0.57.jar和postgresql-42.3.3.jar放置于./files/jars目录,基于Siddhi Docker微服务文档构建了名为siddhiimgpostgres的Docker镜像。

运行命令为:

docker run -it --net postgres-docker_default --rm -p 8006:8006 -v /home/me/siddhi-apps:/apps siddhiimgpostgres:tag1 -Dapps=/apps/post.siddhi 

当前仅能看到Debezium的快照日志(显示表中有11条记录等信息),但无法看到Siddhi输出的具体数据日志,请问这是什么原因?

可能的原因及解决方法
  • 快照操作类型未被CDC源接收
    你的CDC源配置了operation = 'insert',但Debezium快照阶段推送的现有数据操作类型为snapshot,而非insert。Siddhi CDC源仅会处理配置中指定的操作类型,因此快照数据无法进入inputStream。
    解决:将operation参数修改为'insert,snapshot',或直接移除该参数(默认接收所有操作类型)。

  • 日志Sink的输出级别未匹配
    Siddhi的log sink默认日志级别可能较高,若Docker容器的日志捕获级别未对应,会导致Sink输出被过滤。
    解决:在sink配置中显式指定日志级别,示例如下:

    @sink(type = 'log', level='INFO')
    define stream OutputStream (shipment_id long, date_created string);
    

    同时确认Docker容器的日志输出设置,确保能捕获INFO及以上级别的日志。

  • 依赖Jar包缺失或版本不兼容
    现有Jar包可能缺少Debezium PostgreSQL连接器相关依赖,或siddhi-io-cdc与连接器版本不匹配,导致快照数据无法转换为Siddhi流事件。
    解决:确保镜像中包含与siddhi-io-cdc-2.0.12.jar兼容的debezium-connector-postgres系列Jar包。

  • 映射配置与快照数据结构不匹配
    Debezium快照阶段推送的数据可能包含嵌套结构(如字段在after对象下),当前keyvalue映射无法正确提取字段。
    解决:改用json映射并指定正确的字段路径,示例如下:

    @map(type='json', @attributes(shipment_id = 'after.shipment_id', order_id = 'after.order_id', date_created='after.date_created', status='after.status'))
    

内容的提问来源于stack exchange,提问作者Vinoja

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 18:39:19