Hazelcast Jet加载数据库数据到Map后无法作为Pipeline源读取问题
问题原因
- 核心原因是
Sources.map("t1")属于批数据源,仅会在Jet作业启动的瞬间读取一次Map当前的全量数据,读取完成后不再响应后续写入的新数据。如果你的执行顺序是先提交Pipeline作业,再调用sourceStage方法写入数据,作业启动时Map为空,自然无法读取到任何内容。 - 少部分情况为Jet实例不一致:写入数据使用的Jet实例和运行Pipeline作业的Jet实例不属于同一个集群,或者Map的命名空间不匹配,导致作业读取的Map和你写入数据的Map不是同一个实例。
- 极端情况为序列化配置缺失:如果存入Map的键值类型未完成Jet序列化注册,部分版本会出现静默读取失败的问题,你当前用的是String类型,该问题概率极低。
解决方法
- 调整执行顺序:先完成Map数据写入,再提交Pipeline作业,确保作业启动时Map中已经存在全量数据,参考执行流程:
// 第一步:写入数据到Map sourceStage(jet, pipeline); // 第二步:确认数据写入完成 System.out.println("Map已写入数据量:" + jet.getMap("t1").size()); // 第三步:提交Pipeline作业 jet.newJob(pipeline).join();
- 若需要Pipeline能持续读取后续写入Map的新数据,将批源替换为流式Map日志源即可,修改后的读取逻辑如下:
pipeLine.readFrom(Sources.mapJournal("t1", JournalInitialPosition.START_FROM_OLDEST)) .map(map -> { System.err.println(map.getKey() + "---" + map.getValue()); return map; }).writeTo(Sinks.logger());
- 校验集群一致性:确认写入数据、提交作业使用的是同一个Jet集群实例,不存在跨集群访问问题。
内容的提问来源于stack exchange,提问作者vivek
相关产品推荐
相关产品推荐

