1. 为什么选择SSE实现Android消息推送

如果你正在开发一个需要实时消息推送的Android应用,比如新闻推送、股票行情或者AI对话流,SSE(Server-Sent Events)可能是个不错的选择。我第一次接触SSE是在做一个智能客服项目时,当时需要一种轻量级的方案来实现服务器到客户端的单向消息推送。

SSE本质上是一种基于HTTP的长连接技术,它最大的特点就是简单。相比WebSocket,SSE不需要额外的协议升级过程,直接使用普通的HTTP连接就能工作。我记得当时用OkHttp-SSE库只花了不到半小时就实现了基本功能,这在快速迭代的项目中特别实用。

在实际项目中,SSE特别适合以下场景:

  • 单向数据流 :比如服务器推送通知、实时行情数据
  • 轻量级实现 :不需要复杂的握手协议
  • 自动重连 :内置的重试机制让连接更稳定
  • 兼容性好 :基于HTTP,能穿透大多数防火墙

2. 基础实现:OkHttp-SSE快速上手

2.1 添加依赖

首先在build.gradle中添加OkHttp和OkHttp-SSE依赖:

dependencies {
    implementation 'com.squareup.okhttp3:okhttp:4.12.0'
    implementation 'com.squareup.okhttp3:okhttp-sse:4.12.0'
}

建议使用最新稳定版,我实测4.12.0版本在连接稳定性上有明显提升。

2.2 创建EventSource监听器

这个监听器是SSE的核心,处理所有连接事件:

class SSEListener : EventSourceListener() {
    override fun onOpen(eventSource: EventSource, response: Response) {
        Log.d("SSE", "连接已建立")
    }
    
    override fun onEvent(eventSource: EventSource, id: String?, 
                        type: String?, data: String) {
        Log.d("SSE", "收到事件: $data")
        // 这里处理业务逻辑
    }
    
    override fun onClosed(eventSource: EventSource) {
        Log.d("SSE", "连接已关闭")
    }
    
    override fun onFailure(eventSource: EventSource, t: Throwable?, 
                          response: Response?) {
        Log.e("SSE", "连接失败", t)
        // 这里实现重连逻辑
    }
}

2.3 建立SSE连接

配置OkHttpClient时,建议设置合理的超时时间:

val client = OkHttpClient.Builder()
    .connectTimeout(30, TimeUnit.SECONDS)
    .readTimeout(0, TimeUnit.SECONDS) // 0表示不超时
    .build()

val request = Request.Builder()
    .url("https://your-server.com/sse-endpoint")
    .addHeader("Accept", "text/event-stream")
    .build()

val eventSource = EventSources.createFactory(client)
    .newEventSource(request, SSEListener())

这里有个坑要注意:readTimeout设为0表示永不超时,这对SSE长连接是必要的。我曾经因为没设置这个参数,导致连接1分钟后自动断开。

3. 后台服务中的稳定性优化

3.1 使用Foreground Service保持连接

Android系统会限制后台网络访问,所以我们需要把SSE连接放在前台服务中:

class SSEService : Service() {
    private var eventSource: EventSource? = null
    
    override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int {
        startForeground(NOTIFICATION_ID, createNotification())
        startSSEConnection()
        return START_STICKY
    }
    
    private fun startSSEConnection() {
        val client = OkHttpClient.Builder()
            .connectTimeout(30, TimeUnit.SECONDS)
            .readTimeout(0, TimeUnit.SECONDS)
            .build()
        
        val request = Request.Builder()
            .url(SSE_URL)
            .build()
        
        eventSource = EventSources.createFactory(client)
            .newEventSource(request, object : EventSourceListener() {
                // 实现监听器方法
            })
    }
    
    override fun onDestroy() {
        eventSource?.cancel()
        super.onDestroy()
    }
}

记得在AndroidManifest.xml中声明服务和使用前台服务权限:

<uses-permission android:name="android.permission.FOREGROUND_SERVICE" />

<service
    android:name=".SSEService"
    android:enabled="true"
    android:exported="false" />

3.2 实现智能重连机制

网络不稳定时,自动重连是关键。我推荐使用指数退避算法:

override fun onFailure(eventSource: EventSource, t: Throwable?, 
                      response: Response?) {
    val retryDelay = calculateRetryDelay(retryCount)
    handler.postDelayed({
        startSSEConnection()
        retryCount++
    }, retryDelay)
}

private fun calculateRetryDelay(retryCount: Int): Long {
    return minOf(
        INITIAL_RETRY_DELAY_MS * (1 shl retryCount),
        MAX_RETRY_DELAY_MS
    ).toLong()
}

建议初始延迟设为1秒,最大不超过30秒。我在项目中实测这个策略能有效应对短时网络波动。

3.3 心跳检测与保活

有些网络环境会主动关闭空闲连接,我们需要实现心跳机制:

private val heartbeatTask = object : Runnable {
    override fun run() {
        if (isConnected) {
            sendHeartbeat()
        }
        handler.postDelayed(this, HEARTBEAT_INTERVAL)
    }
}

private fun sendHeartbeat() {
    // 发送心跳包或空消息
    Log.d("SSE", "发送心跳")
}

在onOpen中启动心跳任务:

override fun onOpen(eventSource: EventSource, response: Response) {
    handler.post(heartbeatTask)
}

4. 性能优化实战技巧

4.1 电量优化策略

长时间运行的SSE连接可能耗电,我们可以:

  1. 根据网络类型调整频率
val connectivityManager = getSystemService(CONNECTIVITY_SERVICE) as ConnectivityManager
val network = connectivityManager.activeNetwork
val caps = connectivityManager.getNetworkCapabilities(network)

val isMetered = !caps?.hasCapability(NET_CAPABILITY_NOT_METERED) ?: true
val interval = if (isMetered) LONG_INTERVAL else SHORT_INTERVAL
  1. 利用WorkManager在Doze模式下工作
val constraints = Constraints.Builder()
    .setRequiredNetworkType(NetworkType.CONNECTED)
    .setRequiresBatteryNotLow(true)
    .build()

val request = OneTimeWorkRequestBuilder<SSEWorker>()
    .setConstraints(constraints)
    .build()

WorkManager.getInstance(context).enqueue(request)

4.2 数据压缩与批处理

如果推送频率高,可以考虑:

  1. 服务端启用gzip压缩
  2. 实现客户端批处理逻辑:
private val messageQueue = mutableListOf<String>()
private val BATCH_SIZE = 10
private val BATCH_TIMEOUT = 5000L // 5秒

fun onNewMessage(message: String) {
    messageQueue.add(message)
    if (messageQueue.size >= BATCH_SIZE) {
        processBatch()
    } else {
        handler.postDelayed(::processBatch, BATCH_TIMEOUT)
    }
}

private fun processBatch() {
    if (messageQueue.isEmpty()) return
    val batch = messageQueue.joinToString("\n")
    // 处理批量消息
    messageQueue.clear()
}

4.3 连接状态管理

实现一个状态机管理连接状态很有帮助:

sealed class ConnectionState {
    object Disconnected : ConnectionState()
    object Connecting : ConnectionState()
    data class Connected(val since: Long) : ConnectionState()
    data class Failed(val error: Throwable) : ConnectionState()
}

private var state: ConnectionState = ConnectionState.Disconnected

fun updateState(newState: ConnectionState) {
    state = newState
    when (newState) {
        is ConnectionState.Failed -> scheduleReconnect()
        // 其他状态处理
    }
}

5. 常见问题解决方案

5.1 后台连接被杀死

这是最常见的问题,我的解决方案是:

  1. 使用Foreground Service并显示持久通知
  2. 在onTaskRemoved中重启服务:
override fun onTaskRemoved(rootIntent: Intent?) {
    val restartIntent = Intent(applicationContext, SSEService::class.java).apply {
        flags = Intent.FLAG_ACTIVITY_NEW_TASK
    }
    startService(restartIntent)
    super.onTaskRemoved(rootIntent)
}
  1. 监听网络变化自动重连:
val networkCallback = object : ConnectivityManager.NetworkCallback() {
    override fun onAvailable(network: Network) {
        if (state is ConnectionState.Disconnected) {
            startSSEConnection()
        }
    }
}

connectivityManager.registerNetworkCallback(
    NetworkRequest.Builder().build(),
    networkCallback
)

5.2 消息顺序错乱

SSE本身保证消息顺序,但在多线程环境下需要注意:

private val messageLock = Any()

override fun onEvent(eventSource: EventSource, id: String?, 
                    type: String?, data: String) {
    synchronized(messageLock) {
        // 处理消息
    }
}

对于需要严格顺序的业务,可以使用序列号:

data class SSEEvent(
    val seq: Long,
    val type: String,
    val data: String
)

private var lastSeq = -1L

fun processEvent(event: SSEEvent) {
    if (event.seq <= lastSeq) return // 丢弃旧消息
    lastSeq = event.seq
    // 处理消息
}

5.3 内存泄漏预防

长时间运行的服务容易内存泄漏,要注意:

  1. 及时取消回调:
override fun onDestroy() {
    eventSource?.cancel()
    handler.removeCallbacksAndMessages(null)
    connectivityManager.unregisterNetworkCallback(networkCallback)
    super.onDestroy()
}
  1. 使用WeakReference持有Context:
class SSEService : Service() {
    private val contextRef = WeakReference<Context>(this)
    
    // ...
}
  1. 监控内存使用:
val isLowMemory = (activityManager.memoryInfo.availMem < 
                   activityManager.memoryInfo.threshold)
if (isLowMemory) {
    releaseNonCriticalResources()
}

6. 高级应用场景

6.1 结合WorkManager实现可靠推送

对于关键消息,可以结合WorkManager确保送达:

class SSEWorker(context: Context, params: WorkerParameters) 
    : CoroutineWorker(context, params) {

    override suspend fun doWork(): Result {
        val sse = SSEClient()
        try {
            sse.connect()
            return Result.success()
        } catch (e: Exception) {
            return if (runAttemptCount < MAX_RETRY) {
                Result.retry()
            } else {
                Result.failure()
            }
        }
    }
}

6.2 多通道消息合并

当有多个SSE流时,可以合并处理:

class MultiSSEClient {
    private val sources = mutableMapOf<String, EventSource>()
    private val eventChannel = Channel<SSEEvent>(capacity = 100)
    
    fun addSource(id: String, url: String) {
        val listener = object : EventSourceListener() {
            override fun onEvent(source: EventSource, 
                               id: String?, type: String?, 
                               data: String) {
                eventChannel.trySend(SSEEvent(id, type, data))
            }
        }
        val source = createEventSource(url, listener)
        sources[id] = source
    }
    
    suspend fun collectEvents(block: (SSEEvent) -> Unit) {
        for (event in eventChannel) {
            block(event)
        }
    }
}

6.3 与Room数据库集成

持久化推送消息到本地数据库:

@Dao
interface MessageDao {
    @Insert
    suspend fun insert(message: Message)
    
    @Query("SELECT * FROM messages ORDER BY timestamp DESC")
    fun getAll(): Flow<List<Message>>
}

class MessageRepository(private val dao: MessageDao) {
    private val sseListener = object : EventSourceListener() {
        override fun onEvent(source: EventSource, 
                           id: String?, type: String?, 
                           data: String) {
            viewModelScope.launch {
                dao.insert(Message(content = data))
            }
        }
    }
}

7. 监控与调试技巧

7.1 连接状态日志

记录详细的连接日志有助于排查问题:

private fun logConnectionState(state: String, extra: String? = null) {
    val logEntry = """
        |时间: ${SimpleDateFormat("HH:mm:ss").format(Date())}
        |状态: $state
        |网络: ${getNetworkType()}
        |电量: ${getBatteryLevel()}%
        |${extra ?: ""}
        """.trimMargin()
    
    writeToLogFile(logEntry)
}

7.2 使用Stetho调试

集成Facebook的Stetho库可以实时查看SSE连接:

debugImplementation 'com.facebook.stetho:stetho-okhttp3:1.6.0'

然后在调试构建中启用:

if (BuildConfig.DEBUG) {
    client.addNetworkInterceptor(StethoInterceptor())
}

7.3 性能指标监控

跟踪关键性能指标:

class PerformanceMonitor {
    private val connectTimes = mutableListOf<Long>()
    private val messageDelays = mutableListOf<Long>()
    
    fun recordConnectTime(duration: Long) {
        connectTimes.add(duration)
        if (connectTimes.size > 100) connectTimes.removeAt(0)
    }
    
    fun recordMessageDelay(delay: Long) {
        messageDelays.add(delay)
        if (messageDelays.size > 100) messageDelays.removeAt(0)
    }
    
    fun getStats(): String {
        return """
            平均连接时间: ${connectTimes.average()}ms
            平均消息延迟: ${messageDelays.average()}ms
            连接成功率: ${(connectTimes.size.toDouble() / (connectTimes.size + retryCount)) * 100}%
            """
    }
}

8. 测试策略

8.1 模拟服务器实现

使用MockWebServer进行本地测试:

val server = MockWebServer()
server.enqueue(
    MockResponse()
        .setBody("data: test message\n\n")
        .setHeader("Content-Type", "text/event-stream")
)

server.start()
val url = server.url("/sse")

// 使用这个url测试SSE客户端

8.2 自动化测试用例

编写Espresso测试验证SSE功能:

@RunWith(AndroidJUnit4::class)
class SSETest {
    @get:Rule
    val activityRule = ActivityScenarioRule(MainActivity::class.java)

    @Test
    fun testSSEReception() {
        val mockResponse = "data: test message\n\n"
        val server = MockWebServer().apply {
            enqueue(MockResponse().setBody(mockResponse))
            start()
        }
        
        onView(withId(R.id.url_input)).perform(
            replaceText(server.url("/sse").toString())
        )
        
        onView(withId(R.id.connect_button)).perform(click())
        
        onView(withId(R.id.message_view))
            .check(matches(withText(containsString("test message"))))
    }
}

8.3 压力测试

模拟高并发场景:

@Test
fun testHighLoad() {
    val client = OkHttpClient()
    val request = Request.Builder()
        .url(SSE_URL)
        .build()
    
    val latch = CountDownLatch(100)
    
    repeat(100) {
        thread {
            val source = EventSources.createFactory(client)
                .newEventSource(request, object : EventSourceListener() {
                    override fun onEvent(source: EventSource, 
                                       id: String?, type: String?, 
                                       data: String) {
                        latch.countDown()
                    }
                })
        }
    }
    
    assertTrue(latch.await(30, TimeUnit.SECONDS))
}

9. 安全最佳实践

9.1 认证与加密

确保SSE连接安全:

val client = OkHttpClient.Builder()
    .addInterceptor { chain ->
        val request = chain.request().newBuilder()
            .addHeader("Authorization", "Bearer $token")
            .build()
        chain.proceed(request)
    }
    .sslSocketFactory(sslContext.socketFactory, trustManager)
    .build()

9.2 数据验证

验证接收到的消息:

override fun onEvent(source: EventSource, id: String?, 
                    type: String?, data: String) {
    if (!isValidData(data)) {
        source.cancel()
        return
    }
    // 处理有效数据
}

private fun isValidData(data: String): Boolean {
    return try {
        JsonParser.parseString(data)
        true
    } catch (e: Exception) {
        false
    }
}

9.3 防止DDoS攻击

实现客户端限流:

private val rateLimiter = RateLimiter.create(10.0) // 每秒10条

override fun onEvent(source: EventSource, id: String?, 
                    type: String?, data: String) {
    if (!rateLimiter.tryAcquire()) {
        Log.w("SSE", "消息频率过高,丢弃")
        return
    }
    // 处理消息
}

10. 实际项目经验分享

在最近的一个金融项目中,我们使用SSE实现了实时行情推送。最初版本遇到了几个典型问题:

  1. 后台连接不稳定 :通过结合Foreground Service和WorkManager解决
  2. 消息丢失 :实现了本地缓存和重传机制
  3. 电量消耗高 :优化了心跳间隔和网络使用策略

最终实现的架构包含以下组件:

  • SSEClient :处理核心连接逻辑
  • MessageQueue :缓冲和排序消息
  • PersistenceLayer :消息持久化
  • StateMonitor :实时监控连接状态

关键优化点:

  • 根据网络质量动态调整心跳间隔
  • 重要消息确认机制
  • 离线消息缓存
  • 智能重连策略

在实现过程中,最大的收获是理解了Android后台限制对长连接的影响。通过合理的进程管理和资源调度,最终实现了既稳定又省电的解决方案。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐