Created
March 27, 2014 04:11
-
-
Save torao/9800055 to your computer and use it in GitHub Desktop.
finagle-stream を使用して HTTP コンテンツ部分を非同期で送信するためのサンプル。
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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