【问题标题】:Playframework and Twitter Streaming APIPlayframework 和 Twitter 流媒体 API
【发布时间】:2016-08-20 12:14:55
【问题描述】:

如何从 Twitter Streaming API - POST 状态/过滤器读取响应数据? 我已建立连接并收到 200 状态码,但我不知道如何阅读推文。我只想在推文到来时打印它们。

ws.url(url)
.sign(OAuthCalculator(consumerKey, requestToken))
.withMethod("POST")
.stream()
.map { response =>
  if(response.headers.status == 200)
    println(response.body)
} 

编辑:我找到了这个解决方案

ws.url(url)
.sign(OAuthCalculator(consumerKey, requestToken))
.withMethod("POST")
.stream()
.map { response => 
  if(response.headers.status == 200){
    response.body
      .scan("")((acc, curr) => if (acc.contains("\r\n")) curr.utf8String else acc + curr.utf8String)
      .filter(_.contains("\r\n"))
      .map(json => Try(parse(json).extract[Tweet]))
      .runForeach {
        case Success(tweet) =>
          println("-----")
          println(tweet.text)
        case Failure(e) =>
          println("-----")
          println(e.getStackTrace)
      }
  }
}

【问题讨论】:

    标签: scala twitter playframework twitter-streaming-api playframework-2.5


    【解决方案1】:

    流式 WS 请求的响应正文是 Akka Streams Source 字节。由于 Twitter Api 响应是换行符分隔的(通常),您可以使用 Framing.delimiter 将它们拆分为字节块,将块解析为 JSON,然后对它们执行您想要的操作。像这样的东西应该可以工作:

    import akka.stream.scaladsl.Framing
    import scala.util.{Success, Try}
    import akka.util.ByteString
    import play.api.libs.json.{JsSuccess, Json, Reads}
    import play.api.libs.oauth.{ConsumerKey, OAuthCalculator, RequestToken}
    
    case class Tweet(id: Long, text: String)
    object Tweet {
      implicit val reads: Reads[Tweet] = Json.reads[Tweet]
    }
    
    def twitter = Action.async { implicit request =>
      ws.url("https://stream.twitter.com/1.1/statuses/filter.json?track=Rio2016")
          .sign(OAuthCalculator(consumerKey, requestToken))
          .withMethod("POST")
          .stream().flatMap { response =>
        response.body
          // Split up the byte stream into delimited chunks. Note
          // that the chunks are quite big
          .via(Framing.delimiter(ByteString.fromString("\n"), 20000))
          // Parse the chunks into JSON, and then to a Tweet.
          // A better parsing strategy would be to account for all
          // the different possible responses, but here we just
          // collect those that match a Tweet.
          .map(bytes => Try(Json.parse(bytes.toArray).validate[Tweet]))
          .collect {
            case Success(JsSuccess(tweet, _)) => tweet.text
          }
          // Print out each chunk
          .runForeach(println).map { _ =>
            Ok("done")
        }
      }
    }
    

    注意:要实现流,您需要将隐式 Materializer 注入控制器。

    【讨论】:

    • 感谢您的解释
    • 以后可以关闭连接吗?我计划有多个请求来跟踪不同的单词,并且我想在未来的某个时间关闭特定的连接?
    • 查看 Akka 文档中的 Dynamic Stream Handling。一个想法:创建一个共享终止开关,然后使用source.via(killSwitch.flow) 将其添加到流中。在 killswitch 上运行 shutdown() 应该会关闭连接。
    • 所以我已经将我的 WS 流请求添加到 Source(Stream(request)).via(sharedKillSwitch.flow) 但是当我在 killswitch 连接上运行 shutdown() 时仍然打开
    【解决方案2】:

    致电stream() 会返回Future[StreamedResponse]。然后,您必须使用 akka 成语来转换其中的 ByteString 块。类似:

    val stream = ws.url(url)
      .sign(OAuthCalculator(consumerKey, requestToken))
      .withMethod("POST")
      .stream()
    
    stream flatMap { res =>
      res.body.runWith(Sink.foreach[ByteString] { bytes =>
        println(bytes.utf8String)
      })
    }
    

    请注意,我没有测试上面的代码(但它基于 https://www.playframework.com/documentation/2.5.x/ScalaWS 的流响应部分以及来自 http://doc.akka.io/docs/akka/2.4.2/scala/stream/stream-flows-and-basics.html 的接收器描述)

    还请注意,这将在自己的行上打印每个块,我不确定 twitter API 是否会返回每个块的完整 json blob。如果您想在打印之前累积块,您可能需要使用Sink.fold

    【讨论】:

      猜你喜欢
      • 2013-04-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-05-06
      • 1970-01-01
      • 2015-09-15
      • 1970-01-01
      相关资源
      最近更新 更多