新手求教:如何实现支持外部自定义逻辑的嵌套回调机制?
解决方案:完全可行!
首先明确:你的需求不仅可行,而且是封装这类工具库的标准做法——让库聚焦于通用的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
相关产品推荐
相关产品推荐

