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

如何基于Akka-HTTP、Twitter4J实现返回推文的API接口?

解决思路与代码重构方案

首先得戳破核心问题:你遇到的类型不匹配和响应无法返回,根源是Twitter4J的Streaming API是基于回调的,不是同步返回结果的——你之前的filterTwitterStreamByUserID调用的twitterStream.filter()是无返回值的(本质是void方法),它只会通过Counter这个Listener推送推文,而不是直接返回推文集合。另外你的Counter把推文直接打到控制台,没有对外暴露可被API订阅的数据流,所以接口拿不到数据。

下面分步骤解决:


1. 重构Counter:让推文流可被API订阅

把原来直接绑定到控制台Sink的逻辑改成对外暴露推文Source,这样API可以订阅这个流来获取数据:

class Counter extends StatusAdapter {
  implicit val system = ActorSystem("EmojiTrends")
  implicit val materializer = ActorMaterializer()
  implicit val executionContext = system.dispatcher
  private val logger = Logging(system, classOf[Counter])

  // 用队列转Publisher,对外暴露推文Source
  private val statusQueue = Source.queue[Status](1000, OverflowStrategy.backpressure)
    .toMat(Sink.asPublisher(false))(Keep.right)
    .run()

  // 对外提供可订阅的推文Source
  def tweetSource: Source[Status, NotUsed] = Source.fromPublisher(statusQueue)

  override def onStatus(status: Status): Unit = {
    // 非阻塞入队,避免阻塞Listener线程
    statusQueue.offer(status).onComplete {
      case Success(_) => // 入队成功
      case Failure(ex) => logger.error(ex, "Failed to add tweet to queue")
    }
  }
}

2. 修正过滤方法:明确无返回值

既然filter()是启动流的操作,不需要返回值,直接改成Unit类型:

def filterTwitterStreamByUserID(twitterStream: TwitterStream, userId: Array[Long]): Unit = {
  twitterStream.filter(new FilterQuery().follow(userId))
}

3. 修改API路由:订阅推文流并返回结果

Akka HTTP支持直接处理Future和Source,我们可以收集指定数量的推文(或设置超时)后返回给客户端。同时要注意避免重复启动同一个用户的过滤:

trait ApiRoute extends Databases {
  // 初始化资源:只做一次,不要每次请求都创建
  private val twitterStream = TwitterStreamFilters.configureTwitterStream()
  private val counter = new Counter
  twitterStream.addListener(counter)
  // 记录已启动过滤的用户ID,避免重复请求
  private val activeFilters = scala.collection.mutable.Set[Long]()

  // 引入JSON序列化:用spray-json为例,需要添加依赖
  import spray.json._
  import DefaultJsonProtocol._

  // 自定义Status的JSON格式(按需调整字段)
  implicit val tweetFormat: RootJsonFormat[(Long, String, String)] = jsonFormat3((s: Status) => 
    (s.getId, s.getText, s.getCreatedAt.toString)
  )

  val routes = pathPrefix("tweets") {
    pathSingleSlash {
      complete("Twitter Stream API is running!")
    }
  } ~ pathPrefix("tweets" / LongNumber) { userId =>
    get {
      // 仅第一次请求时启动过滤
      if (!activeFilters.contains(userId)) {
        TwitterStreamFilters.filterTwitterStreamByUserID(twitterStream, Array(userId))
        activeFilters.add(userId)
      }

      // 收集10条推文,或等待5秒超时(按需调整)
      val tweetsFuture = counter.tweetSource
        .take(10)
        .map(status => (status.getId, status.getText, status.getCreatedAt.toString))
        .runWith(Sink.seq)
        .withTimeout(5.seconds, Seq.empty)

      // Akka HTTP自动处理Future,将序列转为JSON响应
      complete(tweetsFuture)
    }
  }
}

关于数据库的建议

要不要用数据库取决于你的需求:

  • 如果只需要实时推文:不需要数据库,直接返回流中收集的实时数据即可。
  • 如果需要历史推文查询:必须用数据库——Twitter4J的Streaming API只能获取实时推送的内容,无法回溯历史推文,这时你需要把Counter中收到的推文持久化到数据库(比如PostgreSQL、MongoDB),然后API从数据库查询历史数据。
  • 如果需要缓存近期推文:可以用Redis做缓存,把最近的推文存起来,避免用户重复请求时重新收集。

额外注意事项

  1. JSON序列化依赖:需要添加spray-json或circe的依赖到你的build.sbt,否则Akka HTTP无法将推文序列化为JSON响应。
  2. 资源清理:在ActorSystem关闭时记得关闭twitterStream,避免资源泄漏。
  3. 限流与重试:Twitter API有速率限制,建议添加重试机制和限流逻辑,避免被封禁。
  4. Server-Sent Events(可选):如果想要实时推送推文给客户端(而不是一次性返回集合),可以用Akka HTTP的SSE支持,让客户端保持长连接接收实时推文。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:19:29