Skip to content

Instantly share code, notes, and snippets.

@torao
Created March 27, 2014 04:11
Show Gist options
  • Select an option

  • Save torao/9800055 to your computer and use it in GitHub Desktop.

Select an option

Save torao/9800055 to your computer and use it in GitHub Desktop.
finagle-stream を使用して HTTP コンテンツ部分を非同期で送信するためのサンプル。
import com.twitter.concurrent.{Broker, Offer}
import com.twitter.finagle.Service
import com.twitter.finagle.builder.ServerBuilder
import com.twitter.finagle.stream.{EOF, StreamResponse, Stream}
import com.twitter.util.Future
import java.net.InetSocketAddress
import org.jboss.netty.buffer.{ChannelBuffers, ChannelBuffer}
import org.jboss.netty.handler.codec.http._
import scala.concurrent.ExecutionContext.Implicits.global
import scala.util.{Failure, Success}
/**
* Finagle を使用して HTTP で接続を受け付け 300 までのフィボナッチ数列を非同期で計算 &
* 送信するストリームサーバ。
*
* 画像ファイルのように 1 度にメモリに乗せるのがしんどい Blob データを送信するのに
* [[com.twitter.finagle.http.Http]] や [[com.twitter.finagle.http.RichHttp]]
* の振り替え策として使用しなきゃならないんだけど、本家のサンプルでは Offer で小細工して
* たりストリームの終了方法が実装されていなかったりで分かりづらいので、ちょっとダサいかも
* しれないがもう少しフォーカスを絞って実装しています。
*
* 本家のサンプルはこちら:
* https://github.com/twitter/finagle/blob/master/finagle-example/src/main/scala/com/twitter/finagle/example/stream/StreamServer.scala
*/
object Server {
def main(args:Array[String]):Unit = {
// サーバの起動
ServerBuilder()
.codec(Stream())
.bindTo(new InetSocketAddress(8088))
.name("http-streaming")
.build(service)
}
val service = new Service[HttpRequest, StreamResponse] {
def apply(request:HttpRequest):Future[StreamResponse] = {
// 非同期処理から送信データを渡す Broker (まぁ Queue だな)
val broker = new Broker[ChannelBuffer]()
// 送信を完了するときに EOF を流す Broker
val eof = new Broker[Throwable]()
// フィボナッチ数列を算出する処理
def fact(n:Int):BigInt = if(n == 0) 1 else n * fact(n - 1)
// 非同期でフィボナッチ数列の計算を開始
concurrent.future {
(0 to 300).foreach{ i =>
val value = fact(i) // フィボナッチ数列を計算
val binary = s"$value\n".getBytes("UTF-8") // 送信用データ
broker ! ChannelBuffers.wrappedBuffer(binary) // バッファを共有する場合は copiedBuffer() 使え
}
}.onComplete {
case Success(_) => eof ! EOF // 完了したらふつーにストリーミングを終了
case Failure(ex) => eof ! ex // なんかエラーが発生したらエラーで終わり (Finagleのログにスタックトレースが出てクライアントと切断する)
}
// 割と単純な処理なので StreamResponse のインスタンス自体は apply() の実行スレッド内で作成しているが、
// レスポンスヘッダがすぐに確定しない場合などは Promise を使用しても良いんじゃないすかね
Future(new StreamResponse {
/**
* クライアント側から切断された場合に呼ばれる。
* 非同期処理から Broker へ渡されるデータをそのまま破棄。
*/
override def release():Unit = {
// 本当はここで実行中の処理を中断できると良い
broker.recv.foreach{ _ => None }
}
/**
* Finagle が次のデータを要求したときに呼ばれる。
* Broker 経由で Offer を取り出しているが、返値の Offer はまだ値が確定していない
* かもしれない。確定していない場合は次の broker ! obj で値が確定し Finagle 側で
* 参照することが出来る (そういう意味では Offer は Future に似ている)。
*/
override def messages:Offer[ChannelBuffer] = broker.recv
/**
* レスポンスのステータスやヘッダとして使用するオブジェクト。スーパークラスはメソッ
* ドだが構築時に確定しているなら val でオーバーライドしてやれば良い。と言うかタイ
* ミングによって可変のレスポンスを返せるようにしてもどのタイミングで呼ばれるかよく
* 分からない。
*/
override val httpResponse:HttpResponse = locally {
val res = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK)
res.headers().set("Content-Type", "text/plain; charset=UTF-8")
res.headers().set("Connection", "close")
res.setChunked(true)
res
}
/**
* Finagle から非同期処理でエラーが発生していないか確認するための Offer を得る
* ために呼ばれる。この Offer が有効な値を返すと TCP 接続がクローズされる。
* サーバ側からストリーミングを終了したい場合はここに [[EOF]] を流してやるのが
* 標準らしい (なんかびみょいが)。
*/
override def error:Offer[Throwable] = eof.recv
})
}
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment