在马丁清洁架构中,Worker监听UDP后如何通知ViewModel?
问题
我的应用里有一个UDP Socket服务器,用来监听广播消息,然后通过接口通知UI更新RecyclerView添加消息。之前我用线程监听端口,现在想把项目改成遵循Robert Martin的清洁架构,所以打算把线程换成WorkManager Worker。但现在需要找到一种在Worker接收到广播消息时通知ViewModel的方法,请问怎么用符合清洁架构的方式实现?
之前的实现代码如下:
通知器接口
interface PartyDiscovery { fun onPartyDiscovered(party: Party) fun onPartyEnded(party: Party) }
线程实现
class PartyDiscoveryThread(val discovery: PartyDiscovery) : Thread("PartyDiscovery") { private val TAG = "PartyDiscovery" override fun run() { super.run() val buffer = ByteArray(2048) var socket: DatagramSocket? = null try { Log.e(TAG, "run: Just Started") socket = DatagramSocket(Party.defaultPort) socket.broadcast = true val packet = DatagramPacket(buffer, buffer.size) socket.soTimeout = 100 while (!interrupted()) { try { socket.receive(packet) val jsonData = JsonParser.parseString(buffer.decodeToString(0, packet.length)).asJsonObject val party = Party( jsonData["partyId"].asLong, jsonData["partyRoom"].asInt, jsonData["owner"].asString ) if (jsonData["action"].asString == "started") discovery.onPartyDiscovered(party) else if(jsonData["action"].asString=="ended") discovery.onPartyEnded(party) } catch (e: SocketTimeoutException) { continue } } } catch (e: java.lang.Exception) { Log.e(TAG, "run: Party Discovery Crashed", e.cause) } socket!!.close() Log.e(TAG, "run: Dying") } }
Activity实现
class MainActivity : ComponentActivity(), PartyDiscovery { lateinit var discovery: PartyDiscoveryThread var isDiscovering = false private val TAG = "MainActivity" override fun onResume() { super.onResume() startDiscovery() } fun startDiscovery() { if (!isDiscovering) { discovery = PartyDiscoveryThread(this) discovery!!.start() isDiscovering = true } } fun stopDiscovering() { if (isDiscovering) { discovery!!.interrupt() isDiscovering = false } } override fun onPause() { stopDiscovering() super.onPause() } override fun onPartyDiscovered(party: Party) { // 添加到recyclerview } override fun onPartyEnded(party: Party) { // 从recyclerview移除 } }
解决方案
遵循清洁架构的核心是依赖反转——高层模块(表现层、领域层)不依赖低层模块(数据层),两者都依赖抽象;同时分层隔离,领域层作为核心不依赖任何外部框架。下面是具体实现步骤:
1. 重构领域层(核心抽象)
先定义业务相关的抽象,完全脱离Android框架:
事件封装密封类
用密封类统一封装Party的两种事件:
sealed class PartyEvent { data class Discovered(val party: Party) : PartyEvent() data class Ended(val party: Party) : PartyEvent() }
领域Repository接口
定义数据操作的抽象接口,规定业务需要的能力:
interface PartyRepository { // 观察Party事件流 fun observePartyEvents(): Flow<PartyEvent> // 启动广播监听 fun startDiscovery() // 停止广播监听 fun stopDiscovery() }
2. 实现数据层(依赖领域层,处理具体逻辑)
数据层负责实现领域层的抽象,处理UDP监听、数据存储等具体技术细节,依赖WorkManager、Room等框架,但被抽象隔离。
用Room存储Party状态
通过Room持久化Party的活跃状态,Worker接收到消息后更新数据库,ViewModel通过观察数据库变化获取事件:
// Room Entity @Entity(tableName = "parties") data class PartyEntity( @PrimaryKey val partyId: Long, val partyRoom: Int, val owner: String, val isActive: Boolean // true表示派对活跃,false表示已结束 ) // Room Dao @Dao interface PartyDao { // 观察活跃的派对列表 @Query("SELECT * FROM parties WHERE isActive = 1") fun observeActiveParties(): Flow<List<PartyEntity>> // 插入或更新派对状态 @Insert(onConflict = OnConflictStrategy.REPLACE) suspend fun insertOrUpdateParty(party: PartyEntity) }
实现Repository
实现领域层的PartyRepository,整合WorkManager和Room:
class PartyRepositoryImpl( private val partyDao: PartyDao, private val workManager: WorkManager ) : PartyRepository { override fun observePartyEvents(): Flow<PartyEvent> { // 将Room的活跃派对流转换为业务事件流 return partyDao.observeActiveParties() .map { entities -> // 对比当前列表与上一次列表的差异,生成对应的Discovered/Ended事件 val currentParties = entities.map { it.toParty() }.toSet() val previousParties = _previousParties.getAndSet(currentParties) val discovered = currentParties - previousParties val ended = previousParties - currentParties buildList { discovered.forEach { add(PartyEvent.Discovered(it)) } ended.forEach { add(PartyEvent.Ended(it)) } } } .flatMapConcat { events -> flow { events.forEach { emit(it) } } } } private val _previousParties = mutableStateOf(emptySet<Party>()) override fun startDiscovery() { // 启动唯一的WorkManager任务,避免重复监听 val workRequest = OneTimeWorkRequestBuilder<PartyDiscoveryWorker>() .setConstraints(Constraints.Builder() .setRequiredNetworkType(NetworkType.CONNECTED) .build()) .build() workManager.enqueueUniqueWork( "PartyDiscovery", ExistingWorkPolicy.REPLACE, workRequest ) } override fun stopDiscovery() { // 取消监听任务 workManager.cancelUniqueWork("PartyDiscovery") } // 转换Entity到业务模型 private fun PartyEntity.toParty(): Party { return Party(partyId, partyRoom, owner) } }
实现WorkManager Worker
替换原有的线程,用Worker处理UDP监听逻辑,收到消息后更新Room数据库:
class PartyDiscoveryWorker( context: Context, params: WorkerParameters, private val partyDao: PartyDao // 通过依赖注入(如Hilt)获取 ) : CoroutineWorker(context, params) { private val TAG = "PartyDiscoveryWorker" override suspend fun doWork(): Result { val buffer = ByteArray(2048) var socket: DatagramSocket? = null try { Log.e(TAG, "doWork: 开始监听广播") socket = DatagramSocket(Party.defaultPort) socket.broadcast = true val packet = DatagramPacket(buffer, buffer.size) socket.soTimeout = 100 // 用isStopped判断WorkManager是否要求停止任务 while (!isStopped) { try { socket.receive(packet) val jsonData = JsonParser.parseString( buffer.decodeToString(0, packet.length) ).asJsonObject val party = Party( jsonData["partyId"].asLong, jsonData["partyRoom"].asInt, jsonData["owner"].asString ) // 根据消息类型更新数据库 when (jsonData["action"].asString) { "started" -> { partyDao.insertOrUpdateParty( PartyEntity(party.partyId, party.partyRoom, party.owner, isActive = true) ) } "ended" -> { partyDao.insertOrUpdateParty( PartyEntity(party.partyId, party.partyRoom, party.owner, isActive = false) ) } } } catch (e: SocketTimeoutException) { continue } } } catch (e: Exception) { Log.e(TAG, "doWork: 监听崩溃", e.cause) return Result.failure() } finally { socket?.close() Log.e(TAG, "doWork: 停止监听") } return Result.success() } }
3. 实现表现层(依赖领域层,驱动UI)
ViewModel依赖领域层的PartyRepository,处理事件并暴露UI状态;Activity观察ViewModel的状态更新RecyclerView。
ViewModel实现
class MainViewModel(private val partyRepository: PartyRepository) : ViewModel() { // 暴露给UI的活跃派对列表状态流 private val _activeParties = MutableStateFlow<List<Party>>(emptyList()) val activeParties: StateFlow<List<Party>> = _activeParties init { // 收集领域层的事件流,更新UI状态 viewModelScope.launch { partyRepository.observePartyEvents().collect { event -> val currentList = _activeParties.value.toMutableList() when (event) { is PartyEvent.Discovered -> { if (!currentList.contains(event.party)) { currentList.add(event.party) _activeParties.value = currentList } } is PartyEvent.Ended -> { currentList.remove(event.party) _activeParties.value = currentList } } } } } // 转发控制命令到Repository fun startDiscovery() = partyRepository.startDiscovery() fun stopDiscovery() = partyRepository.stopDiscovery() }
Activity实现
class MainActivity : ComponentActivity() { private val viewModel: MainViewModel by viewModels() // 通过ViewModelProvider或Hilt获取 private lateinit var adapter: PartiesAdapter override fun onCreate(savedInstanceState: Bundle?) { super.onCreate(savedInstanceState) setContentView(R.layout.activity_main) // 初始化RecyclerView和Adapter adapter = PartiesAdapter() binding.recyclerView.adapter = adapter // 观察ViewModel的状态,更新UI lifecycleScope.launch { repeatOnLifecycle(Lifecycle.State.STARTED) { viewModel.activeParties.collect { parties -> adapter.submitList(parties) } } } } override fun onResume() { super.onResume() viewModel.startDiscovery() } override fun onPause() { viewModel.stopDiscovery() super.onPause() } }
为什么符合清洁架构?
- 依赖反转:所有模块都依赖领域层的抽象(
PartyRepository),表现层和数据层不直接依赖彼此,解耦性强。 - 分层隔离:领域层是核心,不依赖任何Android框架;数据层处理技术细节,被抽象隔离;表现层只负责UI逻辑,不关心数据来源。
- 可测试性:可以轻松用Mock实现
PartyRepository来测试ViewModel,无需依赖WorkManager或Room。
内容的提问来源于stack exchange,提问作者Ebrahim Karimi
相关产品推荐
相关产品推荐

