将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() } ) }
关键优化点说明
- 协程Mutex:避免传统锁的阻塞问题,适配协程环境,确保同一时间只有一个协程访问共享资源
- 响应式数据流:替换定时轮询为callbackFlow,减少不必要的并发访问和性能消耗
- 空安全检查:用
requireNotNull和check替代断言,提升代码健壮性 - 排序逻辑补全:解决TreeSet依赖元素可比较的要求,避免运行时异常
内容的提问来源于stack exchange,提问作者xephosbot
相关产品推荐
相关产品推荐

