笔记本MQTT broker对接安卓IoT客户端 Java取dangerLevel及topic存MySQL问题
问题修复方案
问题1:dangerLevel存库为null的原因及修复
根因
- 执行顺序错误:你在订阅回调中先调用
parseMqttPayload执行数据库插入,之后才给timeEvent赋值,插入时timeEvent为null,如果数据库对应字段设置了非空约束,会直接触发SQL异常导致整条插入失败,表现为dangerLevel为空 calculateDangerLevel本身逻辑不会返回null,哪怕传感器解析异常返回-1,也会命中默认分支返回No_Risk- 静态变量并发风险:多个消息同时到达时,静态变量
timeEvent会被互相覆盖,也会导致数据异常
修复思路
调整执行顺序,先计算timeEvent再执行插入操作,同时把需要的参数通过方法传递,避免依赖静态变量
问题2:获取topic存入数据库的修复
根因
- 订阅回调的第一个入参就是当前消息的topic值,你没有把这个值传递到
parseMqttPayload方法中 - 数据库写入时写死了字符串常量
"topicCheck",实际要传入topic变量值
核心修改代码
1. 修改parseMqttPayload方法签名,增加入参
// 新增topic和timeEvent入参,移除对全局静态变量的依赖 private static void parseMqttPayload(String payload, String topic, String timeEvent) { String[] payloadTokens = payload.split(","); // Parse the location if (Objects.equals(payloadTokens[0], "null") || Objects.equals(payloadTokens[1], "null")) return; // Parse the gps reading double y = Double.valueOf(payloadTokens[0]); double x = Double.valueOf(payloadTokens[1]); // Parse the battery reading double battery = Double.valueOf(payloadTokens[2]); // Parse the sensor readings double smokeSensorReading = -1; try { smokeSensorReading = Double.valueOf(payloadTokens[3]); } catch (NumberFormatException e) {LOGGER.warn("smoke解析失败",e);} double gasSensorReading = -1; try { gasSensorReading = Double.valueOf(payloadTokens[4]); } catch (NumberFormatException e) {LOGGER.warn("gas解析失败",e);} double tempSensorReading = -1; try { tempSensorReading = Double.valueOf(payloadTokens[5]); } catch (NumberFormatException e) {LOGGER.warn("temp解析失败",e);} double uvSensorReading = -1; try { uvSensorReading = Double.valueOf(payloadTokens[6]); } catch (NumberFormatException e) {LOGGER.warn("uv解析失败",e);} LOGGER.warn("y: {}", y); LOGGER.warn("x: {}", x); LOGGER.warn("battery: {}", battery); LOGGER.warn("smoke: {}", smokeSensorReading); LOGGER.warn("gas: {}", gasSensorReading); LOGGER.warn("temp: {}", tempSensorReading); LOGGER.warn("uv: {}", uvSensorReading); String dangerLevel = calculateDangerLevel(smokeSensorReading,gasSensorReading,tempSensorReading,uvSensorReading); LOGGER.warn("danger: {}",dangerLevel); //connect,insert,update data in sql via ResultSet try (Connection conn = DriverManager.getConnection( "jdbc:mysql:// localhost:3306/mqttdemo?useSSL=false&serverTimezone=UTC&useLegacyDatetimeCode=false", "xxxx","xxxx")) { Statement stmt2 = conn.createStatement(ResultSet.TYPE_SCROLL_INSENSITIVE, ResultSet.CONCUR_UPDATABLE); ResultSet result = stmt2.executeQuery("SELECT * FROM sensorsdata"); result.moveToInsertRow(); result.updateInt("id", 0); // 把固定字符串改为传入的topic变量 result.updateString(2, topic); result.updateDouble("cordY", y); result.updateDouble("cordX", x); result.updateDouble("battery", battery); result.updateDouble("sensor1", smokeSensorReading); result.updateDouble("sensor2", gasSensorReading); result.updateDouble("sensor3", tempSensorReading); result.updateDouble("sensor4", uvSensorReading); result.updateString(10, dangerLevel); result.updateString(11,timeEvent); result.insertRow(); result.last(); System.out.println("id = " + result.getInt("id")); result.close(); stmt2.close(); } catch (SQLException e) { LOGGER.error("数据库插入失败",e); } }
2. 调整订阅回调执行顺序
subscriber.subscribe(topicProperty, (topic, msg) -> { byte[] payload = msg.getPayload(); LOGGER.debug("Message received: topic={}, payload={}", topic, new String(payload)); // 先计算事件时间 LocalDateTime localDate = LocalDateTime.now(); DateTimeFormatter dtf = DateTimeFormatter.ofPattern("dd/MM/yy h:mm:ss"); String timeEvent =dtf.format(localDate); // 再调用解析方法,传入topic和timeEvent parseMqttPayload(new String(payload), topic, timeEvent); });
优化建议
- 数据库连接建议使用连接池,避免每次插入新建连接导致的性能损耗和连接泄露
- 可以给
calculateDangerLevel补充传感器解析失败(值为-1)的异常分支判断,避免异常数据导致风险等级误判 - 移除不再使用的全局静态变量
dangerLevel、topicCheck、timeEvent,避免并发场景下的数据覆盖问题
内容的提问来源于stack exchange,提问作者Mixalis Pikoulas
相关产品推荐
相关产品推荐

