Loading src/main/java/com/censoft/flink/StreamingJob.java +7 −2 Original line number Diff line number Diff line Loading @@ -8,6 +8,7 @@ import com.censoft.flink.mqtt.MqttConsumer; import com.censoft.flink.transform.AlgorithmBaseFilterFunction; import com.censoft.flink.transform.AlgorithmTypeFlatMapFunction; import com.censoft.flink.transform.ResultExistFilterFunction; import com.censoft.flink.transform.UpdateLiveFilterFunction; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.api.java.utils.ParameterTool; import org.apache.flink.core.fs.FileSystem; Loading Loading @@ -47,8 +48,12 @@ public class StreamingJob { DataStream<AlgorithmPushDto> mqttStream = env.addSource(new MqttConsumer("/censoft/test/" + sceneId)); //3、进行处理 //3、1 默认 是否存在结果集判断 //3、0 更新场景在线时间 SingleOutputStreamOperator<AlgorithmPushDto> streamOperator = mqttStream .filter(new UpdateLiveFilterFunction(sceneId)); //3、1 默认 是否存在结果集判断 streamOperator = streamOperator .filter(new ResultExistFilterFunction()); //3、2 默认 根据分类,分解多个推送信息 Loading @@ -74,7 +79,7 @@ public class StreamingJob { //3、6 打印输出 outputStreamOperator.print(); outputStreamOperator.writeAsText("D:/word1.txt", FileSystem.WriteMode.OVERWRITE).setParallelism(1); // outputStreamOperator.writeAsText("D:/word1.txt", FileSystem.WriteMode.OVERWRITE).setParallelism(1); //3、7 自动执行 env.execute(); } Loading src/main/java/com/censoft/flink/mapper/AlgorithmSceneDao.java +2 −0 Original line number Diff line number Diff line Loading @@ -11,4 +11,6 @@ public interface AlgorithmSceneDao { public AlgorithmSceneBasePo getAlgorithmSceneBasePo(@Param("sceneId") Long sceneId); public List<AlgorithmScenePiecePo> getAlgorithmScenePieceList(@Param("sceneId") Long sceneId); public int updateLiveBySceneId(Long sceneId); } src/main/java/com/censoft/flink/mapper/AlgorithmSceneDao.xml +7 −0 Original line number Diff line number Diff line Loading @@ -15,6 +15,11 @@ <result property="variableValue" column="variable_value"/> </collection> </resultMap> <update id="updateLiveBySceneId"> UPDATE algorithm_scene_base SET live = NOW() WHERE id = #{0} </update> <select id="getAlgorithmSceneBasePo" resultType="com.censoft.flink.domain.AlgorithmSceneBasePo"> SELECT scb.id AS id, scb.scene_name AS sceneName, Loading @@ -39,4 +44,6 @@ WHERE asp.scene_id = #{sceneId} ORDER BY asp.sort ASC </select> </mapper> No newline at end of file src/main/java/com/censoft/flink/mapper/SqlFactory.java +14 −0 Original line number Diff line number Diff line Loading @@ -42,6 +42,20 @@ public class SqlFactory { return algorithmScenePiecePoList; } public int updateLiveBySceneId(Long sceneId) throws IOException { SqlSession session = initSqlSession(); //指定要执行的SQL语句的id //SQL语句的id=namespace+"."+SQL语句所在标签的id属性的值 String SQLID = "com.censoft.flink.mapper.AlgorithmSceneDao" + "." + "updateLiveBySceneId"; //通过SqlSession对象的方法执行SQL语句 int result = session.update(SQLID,sceneId); //提交事务 session.commit(); //最后我们关闭SqlSession对象 session.close(); return result; } private SqlSession initSqlSession() throws IOException { /* myBatis的核心类是SqlSessionFactory Loading src/main/java/com/censoft/flink/transform/AlgorithmTypeFlatMapFunction.java +0 −1 Original line number Diff line number Diff line Loading @@ -47,7 +47,6 @@ public class AlgorithmTypeFlatMapFunction implements FlatMapFunction<AlgorithmPu }).collect(Collectors.toList()); for (AlgorithmPushDto pushDto : collect) { logger.error("分流:" + pushDto); collector.collect(pushDto); } Loading Loading
src/main/java/com/censoft/flink/StreamingJob.java +7 −2 Original line number Diff line number Diff line Loading @@ -8,6 +8,7 @@ import com.censoft.flink.mqtt.MqttConsumer; import com.censoft.flink.transform.AlgorithmBaseFilterFunction; import com.censoft.flink.transform.AlgorithmTypeFlatMapFunction; import com.censoft.flink.transform.ResultExistFilterFunction; import com.censoft.flink.transform.UpdateLiveFilterFunction; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.api.java.utils.ParameterTool; import org.apache.flink.core.fs.FileSystem; Loading Loading @@ -47,8 +48,12 @@ public class StreamingJob { DataStream<AlgorithmPushDto> mqttStream = env.addSource(new MqttConsumer("/censoft/test/" + sceneId)); //3、进行处理 //3、1 默认 是否存在结果集判断 //3、0 更新场景在线时间 SingleOutputStreamOperator<AlgorithmPushDto> streamOperator = mqttStream .filter(new UpdateLiveFilterFunction(sceneId)); //3、1 默认 是否存在结果集判断 streamOperator = streamOperator .filter(new ResultExistFilterFunction()); //3、2 默认 根据分类,分解多个推送信息 Loading @@ -74,7 +79,7 @@ public class StreamingJob { //3、6 打印输出 outputStreamOperator.print(); outputStreamOperator.writeAsText("D:/word1.txt", FileSystem.WriteMode.OVERWRITE).setParallelism(1); // outputStreamOperator.writeAsText("D:/word1.txt", FileSystem.WriteMode.OVERWRITE).setParallelism(1); //3、7 自动执行 env.execute(); } Loading
src/main/java/com/censoft/flink/mapper/AlgorithmSceneDao.java +2 −0 Original line number Diff line number Diff line Loading @@ -11,4 +11,6 @@ public interface AlgorithmSceneDao { public AlgorithmSceneBasePo getAlgorithmSceneBasePo(@Param("sceneId") Long sceneId); public List<AlgorithmScenePiecePo> getAlgorithmScenePieceList(@Param("sceneId") Long sceneId); public int updateLiveBySceneId(Long sceneId); }
src/main/java/com/censoft/flink/mapper/AlgorithmSceneDao.xml +7 −0 Original line number Diff line number Diff line Loading @@ -15,6 +15,11 @@ <result property="variableValue" column="variable_value"/> </collection> </resultMap> <update id="updateLiveBySceneId"> UPDATE algorithm_scene_base SET live = NOW() WHERE id = #{0} </update> <select id="getAlgorithmSceneBasePo" resultType="com.censoft.flink.domain.AlgorithmSceneBasePo"> SELECT scb.id AS id, scb.scene_name AS sceneName, Loading @@ -39,4 +44,6 @@ WHERE asp.scene_id = #{sceneId} ORDER BY asp.sort ASC </select> </mapper> No newline at end of file
src/main/java/com/censoft/flink/mapper/SqlFactory.java +14 −0 Original line number Diff line number Diff line Loading @@ -42,6 +42,20 @@ public class SqlFactory { return algorithmScenePiecePoList; } public int updateLiveBySceneId(Long sceneId) throws IOException { SqlSession session = initSqlSession(); //指定要执行的SQL语句的id //SQL语句的id=namespace+"."+SQL语句所在标签的id属性的值 String SQLID = "com.censoft.flink.mapper.AlgorithmSceneDao" + "." + "updateLiveBySceneId"; //通过SqlSession对象的方法执行SQL语句 int result = session.update(SQLID,sceneId); //提交事务 session.commit(); //最后我们关闭SqlSession对象 session.close(); return result; } private SqlSession initSqlSession() throws IOException { /* myBatis的核心类是SqlSessionFactory Loading
src/main/java/com/censoft/flink/transform/AlgorithmTypeFlatMapFunction.java +0 −1 Original line number Diff line number Diff line Loading @@ -47,7 +47,6 @@ public class AlgorithmTypeFlatMapFunction implements FlatMapFunction<AlgorithmPu }).collect(Collectors.toList()); for (AlgorithmPushDto pushDto : collect) { logger.error("分流:" + pushDto); collector.collect(pushDto); } Loading