You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

笔记本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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 00:36:08