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

WSO2 ESB通过定时任务每小时跨库同步增量数据的序列编写问题

WSO2 ESB 每小时增量同步任务seq.test实现

整个序列按「计算增量时间窗口→查询源库新增数据→写入目标库→空结果兜底」的流程实现,直接用ESB自带的中介器、DB操作组件即可,不需要额外开发插件。

前置准备

先在WSO2 ESB管理控制台配置好两个数据源:

  • 源库数据源JNDI名设为jdbc/sourceDB
  • 目标库数据源JNDI名设为jdbc/targetDB
  • 提前确认源表存在记录创建时间字段(如create_time,datetime/timestamp类型),作为增量过滤的依据

seq.test主序列代码

<sequence name="seq.test" xmlns="http://ws.apache.org/ns/synapse">
    <!-- 计算当前时间、1小时前的时间,格式和数据库存储的时间格式对齐 -->
    <script language="js">
        <![CDATA[
            var now = new Date();
            var oneHourAgo = new Date(now.getTime() - 3600 * 1000);
            function formatTime(date) {
                var y = date.getFullYear();
                var m = String(date.getMonth()+1).padStart(2,'0');
                var d = String(date.getDate()).padStart(2,'0');
                var h = String(date.getHours()).padStart(2,'0');
                var i = String(date.getMinutes()).padStart(2,'0');
                var s = String(date.getSeconds()).padStart(2,'0');
                return y+'-'+m+'-'+d+' '+h+':'+i+':'+s;
            }
            mc.setProperty("syncStartTime", formatTime(oneHourAgo));
            mc.setProperty("syncEndTime", formatTime(now));
        ]]>
    </script>

    <!-- 查询源库1小时内的新增数据,替换成实际业务表名、字段名 -->
    <dblookup>
        <connection><pool><dsName>jdbc/sourceDB</dsName></pool></connection>
        <statement>
            <sql><![CDATA[SELECT id, sms_content, create_time, sender FROM t_sms_record WHERE create_time >= ? AND create_time < ?]]></sql>
            <parameter expression="$ctx:syncStartTime" type="VARCHAR"/>
            <parameter expression="$ctx:syncEndTime" type="VARCHAR"/>
            <result column="id" name="id"/>
            <result column="sms_content" name="sms_content"/>
            <result column="create_time" name="create_time"/>
            <result column="sender" name="sender"/>
        </statement>
    </dblookup>

    <!-- 无新增数据直接结束流程 -->
    <filter xpath="$ctx:id != ''">
        <then>
            <!-- 遍历查询结果,逐条写入目标库,数据量大可替换为批量写入逻辑 -->
            <iterate expression="//row" id="syncIterate">
                <target sequence="seq.write.target"/>
            </iterate>
        </then>
        <else>
            <log level="custom">
                <property name="syncTask" value="当前时间窗口无新增数据,同步结束"/>
                <property name="timeRange" expression="concat($ctx:syncStartTime,' ~ ',$ctx:syncEndTime)"/>
            </log>
            <drop/>
        </else>
    </filter>
</sequence>

写入目标库的子序列seq.write.target代码

<sequence name="seq.write.target" xmlns="http://ws.apache.org/ns/synapse">
    <!-- 提取单条记录的字段值 -->
    <property name="target_id" expression="//id/text()"/>
    <property name="target_content" expression="//sms_content/text()"/>
    <property name="target_createtime" expression="//create_time/text()"/>
    <property name="target_sender" expression="//sender/text()"/>

    <!-- 执行目标库插入,替换成实际目标表、字段 -->
    <dbreport>
        <connection><pool><dsName>jdbc/targetDB</dsName></pool></connection>
        <statement>
            <sql><![CDATA[INSERT INTO t_sms_record_target (id,sms_content,create_time,sender) VALUES (?,?,?,?)]]></sql>
            <parameter expression="$ctx:target_id" type="VARCHAR"/>
            <parameter expression="$ctx:target_content" type="VARCHAR"/>
            <parameter expression="$ctx:target_createtime" type="VARCHAR"/>
            <parameter expression="$ctx:target_sender" type="VARCHAR"/>
        </statement>
    </dbreport>

    <log level="custom">
        <property name="syncTask" value="单条数据同步完成"/>
        <property name="dataId" expression="$ctx:target_id"/>
    </log>
</sequence>

注意事项

  • 你当前任务配置用的interval="3600"是从服务启动时间开始算间隔1小时,如果要固定每小时整点执行,避免时间偏移,把task里的trigger替换为<trigger cron="0 0 * * * ?"/>即可。
  • 如果数据库时区和ESB部署节点时区不一致,需要在JS时间计算逻辑里加时区偏移修正,避免查询的时间窗口错位漏数据。
  • 单小时增量数据超过1000条时,不要用逐条写入的逻辑,把查询到的结果集拼接成批量插入SQL执行,降低数据库IO消耗。
  • 多节点部署ESB时,要在序列开头加分布式锁逻辑,避免多个节点同时触发任务导致重复插入数据。
  • 源表的时间字段必须加索引,否则每小时查询全表扫描会严重影响数据库性能。
  • 仅靠时间窗口做增量过滤可能因为数据库时间回拨、任务执行延迟导致漏数/重数,建议给源表加sync_status同步状态字段,查询后标记已同步状态,下次查询过滤已同步数据,可靠性更高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.31 21:12:12