Android SSE消息推送:OkHttp-SSE库在后台服务中的实践与优化
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连接可能耗电,我们可以:
- 根据网络类型调整频率 :
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
- 利用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 数据压缩与批处理
如果推送频率高,可以考虑:
- 服务端启用gzip压缩
- 实现客户端批处理逻辑:
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 后台连接被杀死
这是最常见的问题,我的解决方案是:
- 使用Foreground Service并显示持久通知
- 在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)
}
- 监听网络变化自动重连:
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 内存泄漏预防
长时间运行的服务容易内存泄漏,要注意:
- 及时取消回调:
override fun onDestroy() {
eventSource?.cancel()
handler.removeCallbacksAndMessages(null)
connectivityManager.unregisterNetworkCallback(networkCallback)
super.onDestroy()
}
- 使用WeakReference持有Context:
class SSEService : Service() {
private val contextRef = WeakReference<Context>(this)
// ...
}
- 监控内存使用:
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实现了实时行情推送。最初版本遇到了几个典型问题:
- 后台连接不稳定 :通过结合Foreground Service和WorkManager解决
- 消息丢失 :实现了本地缓存和重传机制
- 电量消耗高 :优化了心跳间隔和网络使用策略
最终实现的架构包含以下组件:
- SSEClient :处理核心连接逻辑
- MessageQueue :缓冲和排序消息
- PersistenceLayer :消息持久化
- StateMonitor :实时监控连接状态
关键优化点:
- 根据网络质量动态调整心跳间隔
- 重要消息确认机制
- 离线消息缓存
- 智能重连策略
在实现过程中,最大的收获是理解了Android后台限制对长连接的影响。通过合理的进程管理和资源调度,最终实现了既稳定又省电的解决方案。
更多推荐




所有评论(0)