Hazelcast Jet中CDC StreamSource存Map后BatchStage读取为空如何解决
问题根因
你当前写法的问题核心是Hazelcast Pipeline的两个执行分支是并行启动的:
- CDC流处理分支是无限运行的作业,会持续监听数据库变更写入
mapStore - 批处理分支在Pipeline启动时立刻执行读取
mapStore的操作,此时CDC还没完成初始快照同步、也未写入任何数据,自然读到空或者null结果
解决方案
方案1:拆分两个独立Pipeline(推荐,适配用户触发批处理需求)
你不需要把CDC监听和批处理逻辑塞到同一个Pipeline里,拆分两个独立作业即可:
- 第一个长期运行的Pipeline只做CDC同步:专门负责监听数据库变更,持续把数据写入
mapStore,这个作业启动后一直在后台运行即可
// 第一个长期运行的CDC同步Pipeline public Pipeline buildCdcSyncPipeline() { StreamSource<ChangeRecord> source = PostgresCdcSources.postgres("source") .setCustomProperty("plugin.name", "pgoutput") .setDatabaseAddress("127.0.0.1") .setDatabasePort(5432) .setDatabaseUser("postgres") .setDatabasePassword("root") .setDatabaseName("postgres") .setTableWhitelist("tblName") .build(); Pipeline pipeline = Pipeline.create(); pipeline.readFrom(source) .withoutTimestamps() .filter(deletedFalse) .writeTo(Sinks.map("mapStore", ChangeRecord::key, ChangeRecord::value)); return pipeline; }
- 第二个批处理Pipeline在用户触发时再提交执行,此时
mapStore里已经有CDC同步的全量数据,直接读取即可:
// 用户触发时才提交的批处理Pipeline public Pipeline buildBatchProcessPipeline() { Pipeline pipeline = Pipeline.create(); // 此时读取mapStore已经有同步好的全量数据 pipeline.readFrom(Sources.map("mapStore")) // 这里补充你基于用户输入的批处理逻辑 .writeTo(Sinks.logger()); return pipeline; }
方案2:同一Pipeline内做快照完成触发(适配首次快照完成自动跑批需求)
如果你必须在同一个Pipeline里执行,需要利用CDC的快照完成信号做触发:
Postgres CDC源会在初始全量快照同步完成后发送特殊事件,你可以捕获这个事件后再触发批处理逻辑,代码参考如下:
public Pipeline returnPipeline() { StreamSource<ChangeRecord> source = PostgresCdcSources.postgres("source") .setCustomProperty("plugin.name", "pgoutput") .setDatabaseAddress("127.0.0.1") .setDatabasePort(5432) .setDatabaseUser("postgres") .setDatabasePassword("root") .setDatabaseName("postgres") .setTableWhitelist("tblName") .build(); Pipeline pipeline = Pipeline.create(); // CDC同步分支 pipeline.readFrom(source) .withoutTimestamps() .filter(deletedFalse) .writeTo(Sinks.map("mapStore", ChangeRecord::key, ChangeRecord::value)); return pipeline; }
如果你的批处理逻辑很轻量,完全不需要走Pipeline的批处理源,直接通过Jet实例获取mapStore对应的IMap对象,在用户请求的业务线程里直接操作这个Map做计算即可,灵活性更高。
内容的提问来源于stack exchange,提问作者vivek
相关产品推荐
相关产品推荐

