如何使用Retrofit与协程在Android应用中消费Filtered Twitter Stream API
基于Retrofit+协程实现Android端消费Twitter V2 Filtered Stream接口指南
1. 前置依赖配置
先在模块级build.gradle中添加所需依赖:
// Retrofit核心库 implementation 'com.squareup.retrofit2:retrofit:2.9.0' // 协程Android扩展 implementation 'org.jetbrains.kotlinx:kotlinx-coroutines-android:1.7.3' // OkHttp网络客户端 implementation 'com.squareup.okhttp3:okhttp:4.11.0' // JSON解析工具,可选Gson/Moshi implementation 'com.squareup.retrofit2:converter-gson:2.9.0'
2. 配置OkHttp客户端
流接口属于长连接,需要调整超时配置,统一注入鉴权头:
val okHttpClient = OkHttpClient.Builder() // 长连接场景读超时设为0表示无超时限制 .readTimeout(0, TimeUnit.SECONDS) .connectTimeout(15, TimeUnit.SECONDS) .writeTimeout(15, TimeUnit.SECONDS) // 统一添加Twitter开发者鉴权头 .addInterceptor { chain -> val request = chain.request().newBuilder() .addHeader("Authorization", "Bearer 你的Twitter开发者平台Bearer Token") .build() chain.proceed(request) } .build()
3. 定义Retrofit服务接口
必须添加@Streaming注解,避免Retrofit将整个响应一次性加载到内存,返回原始响应体逐行处理:
interface TwitterStreamService { @Streaming @GET("2/tweets/search/stream") // 可根据需求自定义查询参数,比如tweet.fields、expansions等 suspend fun getFilteredStream( @Query("tweet.fields") tweetFields: String = "created_at,author_id,text" ): Response<ResponseBody> }
初始化Retrofit实例:
val twitterService = Retrofit.Builder() .baseUrl("https://api.twitter.com/") .client(okHttpClient) .build() .create(TwitterStreamService::class.java)
4. 协程中封装流消费逻辑
推荐用Flow封装流数据,方便上层订阅和生命周期绑定:
fun getTwitterStreamFlow(): Flow<Tweet> = flow { val response = twitterService.getFilteredStream() if (!response.isSuccessful) { throw Exception("请求失败,错误码: ${response.code()}") } val responseBody = response.body() ?: throw Exception("响应体为空") // 逐行读取流内容 responseBody.byteStream().bufferedReader().useLines { lines -> lines.forEach { line -> // Twitter流会定时发送空行作为心跳包,直接跳过即可 if (line.isBlank()) return@forEach // 将JSON字符串解析为你自定义的Tweet数据类 val tweet = Gson().fromJson(line, Tweet::class.java) emit(tweet) } } }.flowOn(Dispatchers.IO)
5. UI层订阅处理
在ViewModel中绑定生命周期订阅流,避免内存泄漏:
viewModelScope.launch { getTwitterStreamFlow() .catch { e -> // 处理异常场景,比如流断开、鉴权失败,可在此处添加重连逻辑 Log.e("TwitterStream", "流异常: ${e.message}") } .collect { tweet -> // 处理收到的推文数据,更新UI Log.d("TwitterStream", "收到推文: ${tweet.text}") } }
注意事项
- 消费流之前需要先通过Twitter V2接口创建对应过滤规则,否则流不会返回任何数据
- 协程取消时会自动关闭流资源,无需手动释放,避免内存泄漏
- 流读取逻辑必须放在IO调度器执行,禁止在主线程直接处理
- 生产环境建议添加指数退避重连逻辑,应对网络波动导致的流意外断开问题
内容的提问来源于stack exchange,提问作者Swati Shrivastava
相关产品推荐
相关产品推荐

