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
相关产品推荐
相关产品推荐

