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

新手求教:如何实现支持外部自定义逻辑的嵌套回调机制?

解决方案:完全可行!

首先明确:你的需求不仅可行,而且是封装这类工具库的标准做法——让库聚焦于通用的MQTT订阅/连接逻辑,把具体的回调行为(比如更新UI、处理不同topic的消息)交给上层应用自定义。下面结合你的代码和场景,一步步实现:

一、先解决原代码的核心问题:回调覆盖

你原来的subscribe函数每次调用都会给client设置新的MqttCallback,这会导致后面的订阅覆盖前面的回调,完全无法实现多订阅多逻辑的需求。我们需要改成让库维护每个订阅对应的回调集合,而不是全局覆盖。

1. 定义订阅回调的数据结构

用一个数据类来绑定订阅的topics和对应的回调逻辑:

data class SubscriptionCallback(
    val topics: List<String>,
    val onMsg: ((String, MqttMessage) -> Unit)?,
    val conLost: ((Throwable) -> Unit)?,
    val delComp: ((IMqttDeliveryToken) -> Unit)?
)

2. 维护client与回调列表的关联

用两个集合分别记录每个client的回调列表,以及是否已经给client设置过全局Callback:

// 关联client和它的所有订阅回调
private val clientCallbacks = mutableMapOf<MqttAndroidClient, MutableList<SubscriptionCallback>>()
// 记录已设置全局Callback的client,避免重复设置
private val initializedClients = mutableSetOf<MqttAndroidClient>()

3. 重写subscribe函数

修改后的函数会把每个订阅的回调加入集合,且只给client设置一次全局Callback,触发时遍历匹配的回调:

fun subscribe(
    client: MqttAndroidClient,
    topics: MutableList<String>,
    onMsg: ((String, MqttMessage) -> Unit)? = null,
    conLost: ((Throwable) -> Unit)? = null,
    delComp: ((IMqttDeliveryToken) -> Unit)? = null
) {
    if (!client.isConnected) {
        Log.w("MQTT_LIB", "Client not connected, can't subscribe to topics: $topics")
        conLost?.invoke(IllegalStateException("Client is not connected"))
        return
    }

    // 逐个订阅topic,处理订阅结果
    topics.forEach { topic ->
        client.subscribe(topic, 0, null, object : IMqttActionListener {
            override fun onSuccess(asyncActionToken: IMqttToken?) {
                Log.i("MQTT_LIB", "Successfully subscribed to $topic")
            }

            override fun onFailure(asyncActionToken: IMqttToken?, exception: Throwable?) {
                Log.e("MQTT_LIB", "Failed to subscribe to $topic", exception)
                conLost?.invoke(exception ?: RuntimeException("Subscribe failed without exception"))
            }
        })
    }

    // 将当前订阅的回调加入client的回调列表
    val callbackList = clientCallbacks.getOrPut(client) { mutableListOf() }
    callbackList.add(SubscriptionCallback(topics, onMsg, conLost, delComp))

    // 仅第一次给client设置全局Callback,避免覆盖
    if (!initializedClients.contains(client)) {
        client.setCallback(object : MqttCallback {
            override fun connectionLost(cause: Throwable) {
                Log.i("MQTT_LIB", "Connection lost")
                // 通知所有注册的连接丢失回调
                clientCallbacks[client]?.forEach { it.conLost?.invoke(cause) }
            }

            override fun messageArrived(topic: String, message: MqttMessage) {
                Log.i("MQTT_LIB", "Message arrived: $topic -> ${message.toString()}")
                // 找到订阅了该topic的所有回调并触发
                clientCallbacks[client]?.forEach { callback ->
                    if (topic in callback.topics) {
                        callback.onMsg?.invoke(topic, message)
                    }
                }
            }

            override fun deliveryComplete(token: IMqttDeliveryToken) {
                Log.i("MQTT_LIB", "Delivery complete for token: ${token.messageId}")
                // 通知所有注册的交付完成回调
                clientCallbacks[client]?.forEach { it.delComp?.invoke(token) }
            }
        })
        initializedClients.add(client)
    }
}

二、上层应用自定义逻辑的调用示例

现在你可以在不同应用/场景中,调用subscribe时传入完全不同的逻辑:

场景1:Android应用中更新UI

// 订阅天气topic,收到消息后切换到主线程更新TextView
subscribe(
    client = mqttClient,
    topics = mutableListOf("weather/beijing"),
    onMsg = { topic, msg ->
        runOnUiThread {
            tvWeather.text = "北京天气:${msg.toString()}"
        }
    },
    conLost = { cause ->
        runOnUiThread {
            tvStatus.text = "连接丢失:${cause.message}"
        }
    }
)

场景2:后台服务中处理设备状态

// 订阅设备状态topic,收到异常消息触发告警
subscribe(
    client = mqttClient,
    topics = mutableListOf("device/001/status"),
    onMsg = { topic, msg ->
        val status = msg.toString()
        Log.d("DEVICE_SERVICE", "设备001状态:$status")
        if (status == "ERROR") {
            sendAlarmNotification("设备001异常!")
        }
    }
)

三、补充:配套的取消订阅逻辑

为了避免内存泄漏,需要添加取消订阅的函数,移除对应的回调:

fun unsubscribe(client: MqttAndroidClient, topics: List<String>) {
    // 取消MQTT订阅
    topics.forEach { client.unsubscribe(it) }
    
    // 移除对应topic的回调
    clientCallbacks[client]?.removeAll { callback ->
        callback.topics.intersect(topics.toSet()).isNotEmpty()
    }
    
    // 如果client没有剩余回调,清理资源
    if (clientCallbacks[client].isNullOrEmpty()) {
        client.setCallback(null)
        clientCallbacks.remove(client)
        initializedClients.remove(client)
    }
}

四、对应你最初的MyClass场景

其实这个思路和你提到的MyClass/YourClass场景完全一致,只需要把回调逻辑作为参数传递给方法即可:

class MyClass {
    // 用Kotlin函数类型简化回调定义
    fun connectToService(onSuccess: () -> Unit, onFailure: () -> Unit) {
        // 模拟计算过程
        val isSuccess = performCalculation()
        if (isSuccess) onSuccess() else onFailure()
    }

    private fun performCalculation(): Boolean {
        // 这里写你的实际计算逻辑
        return true
    }
}

class YourClass {
    fun myFunc() {
        val myClass = MyClass()
        
        // 自定义逻辑1:成功时reinventTheWheel,失败时squareIt
        myClass.connectToService(
            onSuccess = { reinventTheWheel() },
            onFailure = { squareIt() }
        )
        
        // 自定义逻辑2:成功时combAWatermelon,失败不处理
        myClass.connectToService(
            onSuccess = { combAWatermelon() },
            onFailure = {}
        )
        
        // 自定义逻辑3:成功时sharpenMyHammer,失败不处理
        myClass.connectToService(
            onSuccess = { sharpenMyHammer() },
            onFailure = {}
        )
    }

    private fun reinventTheWheel() {}
    private fun squareIt() {}
    private fun combAWatermelon() {}
    private fun sharpenMyHammer() {}
}

关键注意事项

  • 线程切换:MQTT的回调默认在后台线程触发,Android中更新UI必须用runOnUiThread、Coroutine或Handler切换到主线程。
  • 内存泄漏:如果在Activity/Fragment中调用订阅,务必在生命周期结束时(比如onDestroy)调用unsubscribe,避免持有Activity引用导致泄漏。
  • 回调去重:如果同一个topic被多次订阅,会触发所有对应的回调,根据你的需求可以调整逻辑(比如只保留最新的回调)。

内容的提问来源于stack exchange,提问作者Luis Fernandez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 18:12:35