You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Android中Retrofit+RxJava调用Twitter statuses/filter无响应求助

问题分析与解决方案

看起来你遇到的核心问题是:Twitter的statuses/filter是长连接流式API,它会分块持续返回数据,但你的代码没有正确处理流式响应的读取逻辑,导致RxJava的onNext/onError回调无法触发——虽然HTTP状态码返回了200,但响应体的数据流根本没被消费。

下面是针对性的修复建议,一步步来解决问题:

1. 正确处理流式ResponseBody的读取

@Streaming注解只是告诉Retrofit不要一次性把整个响应加载到内存,但你需要手动读取响应流的内容。直接返回Flowable<ResponseBody>不会自动触发onNext,因为ResponseBody本身只是一个容器,需要你主动去读取它的字节流。

修改你的代码,把流式响应拆分成逐行的字符串输出:

private fun getData(resArray: List<String>): Flowable<String> {
    return retrofit().requestFiltered()
        .flatMap { responseBody ->
            // 读取流式响应的每一行(Twitter流式API每行是一个JSON对象)
            val source = responseBody.source()
            Flowable.generate<String, Unit>(
                { Unit },
                { _, emitter ->
                    if (!source.exhausted()) {
                        val line = source.readUtf8Line()
                        if (line != null && line.isNotEmpty()) {
                            emitter.onNext(line)
                        }
                    } else {
                        emitter.onComplete()
                    }
                }
            )
            .doFinally {
                responseBody.close() // 务必关闭流,避免资源泄漏
            }
        }
}

// 同时修改TwitterService接口,移除多余的HeaderMap参数
interface TwitterService {
    // ... 其他方法保持不变

    @Headers("Cache-Control: max-age=640000")
    @POST("statuses/filter.json")
    @Streaming
    fun requestFiltered(
        @Query("track") track: String = "foo",
        @Query("follow") follow: String = "1234"
    ): Flowable<ResponseBody>
}

2. 移除多余的OAuth Header传递

你在requestFiltered里手动添加了oauthToken相关的Header,但实际上SigningInterceptor已经会自动处理OAuth的签名逻辑,手动添加这些Header会导致签名冲突,甚至让请求失效。直接删除getData里的HeaderMap相关代码即可。

3. 调整超时设置适配长连接

流式API是持续保持连接的,你当前设置的30秒读取超时会导致连接被主动断开。建议把读取超时设置为无限制:

private fun getClient(): OkHttpClient {
    val clientBuilder = OkHttpClient.Builder()
        .apply {
            connectTimeout(10, TimeUnit.SECONDS) // 连接超时保留
            writeTimeout(30, TimeUnit.SECONDS) // 写入超时保留
            readTimeout(0, TimeUnit.SECONDS) // 无读取超时,适配长连接
        }
    // ... 其他配置保持不变
}

4. 优化RxJava订阅逻辑

确保线程调度合理,并且添加完整的错误日志(方便排查问题):

private fun calls(): Disposable {
    return retrofit()
        .requestToken(REQUEST_TOKEN)
        .subscribeOn(Schedulers.io())
        .observeOn(Schedulers.io())
        .toFlowable(BackpressureStrategy.BUFFER)
        .flatMap { getData(it.string().split("&")) }
        .observeOn(AndroidSchedulers.mainThread()) // 如果需要在UI线程输出日志/更新UI
        .subscribe(
            { Log.d("TwitterStream", "收到推文: $it") },
            { Log.e("TwitterStream", "流请求出错", it) }, // 打印完整异常栈
            { Log.d("TwitterStream", "流连接结束") }
        )
}

额外检查点

  • 确认你的OAuth凭证(consumer key、secret、access token、secret)完全正确,且账号有权限访问Twitter的流式API;
  • 通过OkHttp日志拦截器查看响应头,确认是否存在Transfer-Encoding: chunked(这是流式响应的标志);
  • 检查Twitter的API限制,确保你的账号没有达到请求频次上限。

内容的提问来源于stack exchange,提问作者davthecoder

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.07 15:52:32