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

将TDLib Java示例转Kotlin协程时遭遇并发异常求助

问题解决:TDLib Kotlin协程版并发异常修复与代码重构

核心问题分析

  • mainChatList 使用的 TreeSet 本身非线程安全,多个协程同时执行读写操作(setChatPositions 修改集合、getData 遍历集合)会触发并发修改异常
  • setChatPositions 中对 TdApi.Chat 对象的属性修改(chat.positions)未做同步,和 getData 的读取操作存在竞态条件
  • getData 采用定时轮询的方式获取数据,不仅效率低,还会加剧并发冲突的概率

重构方案与代码实现

1. 引入协程友好的同步机制

使用Kotlin协程的 Mutex 替代传统锁,确保对共享资源的访问互斥,避免并发冲突:

import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import java.util.concurrent.ConcurrentHashMap

private val chats: ConcurrentHashMap<Long, TdApi.Chat> = ConcurrentHashMap()
private val mainChatList: NavigableSet<OrderedChat> = TreeSet()
private val mutex = Mutex() // 协程互斥锁,保护mainChatList和chat对象的修改

2. 同步修改集合与Chat对象

将 setChatPositions 改为挂起函数,用 mutex.withLock 包裹所有对共享资源的操作:

private suspend fun setChatPositions(chat: TdApi.Chat, positions: Array<TdApi.ChatPosition?>) {
    mutex.withLock {
        // 移除旧的主列表位置
        for (position in chat.positions) {
            if (position.list.constructor == TdApi.ChatListMain.CONSTRUCTOR) {
                val isRemoved = mainChatList.remove(OrderedChat(chat.id, position))
                check(isRemoved) { "Failed to remove chat from main list" }
            }
        }
        // 更新chat的positions
        chat.positions = positions
        // 添加新的主列表位置
        for (position in chat.positions) {
            if (position?.list?.constructor == TdApi.ChatListMain.CONSTRUCTOR) {
                val isAdded = mainChatList.add(OrderedChat(chat.id, position))
                check(isAdded) { "Failed to add chat to main list" }
            }
        }
    }
}

3. 优化数据获取逻辑(替换定时轮询)

用 callbackFlow 实现响应式数据发射,仅当共享资源变化时才发送更新,同时保证读取操作在互斥锁保护下进行:

import kotlinx.coroutines.channels.awaitClose
import kotlinx.coroutines.flow.callbackFlow
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.launch

fun getData(): Flow<List<Chat>> = callbackFlow {
    // 定义更新数据的函数
    suspend fun sendUpdatedData() {
        val chatList = mutex.withLock {
            mainChatList.map { orderedChat ->
                val chat = requireNotNull(chats[orderedChat.chatId]) { "Chat not found: ${orderedChat.chatId}" }
                Chat(
                    id = chat.id,
                    title = chat.title.ifEmpty { "Deleted account" },
                    photo = chat.photo?.let {
                        ProfilePhoto(
                            thumbnail = it.minithumbnail!!.data,
                            file = it.small
                        )
                    },
                    isPinned = orderedChat.position.isPinned,
                    unreadCount = chat.unreadCount,
                    lastMessage = chat.lastMessage
                )
            }.toList()
        }
        trySend(chatList)
    }

    // 首次发送初始数据
    sendUpdatedData()

    // 启动协程监听更新,这里保留原轮询逻辑,实际可替换为TDLib的更新事件流
    val updateJob = launch {
        while (true) {
            sendUpdatedData()
            delay(1000L)
        }
    }

    awaitClose { updateJob.cancel() }
}

4. 补全OrderedChat的排序逻辑

TreeSet 要求元素实现 Comparable 接口,否则会抛出类型转换异常,补充排序逻辑:

data class OrderedChat(
    val chatId: Long,
    val position: TdApi.ChatPosition
) : Comparable<OrderedChat> {
    override fun compareTo(other: OrderedChat): Int {
        // 按照ChatPosition的order排序,值越小排名越靠前(匹配置顶逻辑)
        return position.order.compareTo(other.position.order)
    }
}

5. 优化load函数中的协程调用

setChatPositions 已改为挂起函数,直接在协程上下文调用即可:

suspend fun load(limit: Int) = withContext(Dispatchers.IO) {
    awaitAll(
        async { telegramRepository.loadChats(limit) },
        async {
            telegramRepository.newChatFlow
                .onEach { chat ->
                    chats[chat.id] = chat
                    val positions = chat.positions.clone()
                    chat.positions = arrayOfNulls(0)
                    setChatPositions(chat, positions)
                }.collect()
        },
        async {
            telegramRepository.chatLastMessageFlow
                .onEach { updateChat ->
                    chats[updateChat.chatId]?.let { chat ->
                        chat.lastMessage = updateChat.lastMessage
                        setChatPositions(chat, updateChat.positions)
                    }
                }.collect()
        },
        async {
            telegramRepository.chatPositionFlow
                .onEach { updateChat ->
                    chats[updateChat.chatId]?.let { chat ->
                        // 原有的newPositions计算逻辑保持不变
                        var i = 0
                        for (k in chat.positions.indices) {
                            if (chat.positions[k].list.constructor == TdApi.ChatListMain.CONSTRUCTOR) {
                                break
                            }
                            i++
                        }
                        val newPositions = arrayOfNulls<TdApi.ChatPosition>(
                            size = chat.positions.size + (if (updateChat.position.order == 0L) 0 else 1)
                                    - (if (i < chat.positions.size) 1 else 0)
                        )
                        var pos = 0
                        if (updateChat.position.order != 0L) {
                            newPositions[pos++] = updateChat.position
                        }
                        for (j in chat.positions.indices) {
                            if (i != j) {
                                newPositions[pos++] = chat.positions[j]
                            }
                        }
                        check(pos == newPositions.size)
                        setChatPositions(chat, newPositions)
                    }
                }.collect()
        }
    )
}

关键优化点说明

  1. 协程Mutex:避免传统锁的阻塞问题,适配协程环境,确保同一时间只有一个协程访问共享资源
  2. 响应式数据流:替换定时轮询为callbackFlow,减少不必要的并发访问和性能消耗
  3. 空安全检查:用 requireNotNull 和 check 替代断言,提升代码健壮性
  4. 排序逻辑补全:解决TreeSet依赖元素可比较的要求,避免运行时异常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 07:52:57