Upgrade to Pro — share decks privately, control downloads, hide ads and more …

Dataflow周りで地雷を踏みまくった話

Sponsored · Ship Features Fearlessly Turn features on and off without deploys. Used by thousands of Ruby developers.

 Dataflow周りで地雷を踏みまくった話

Scio + Dataflowでバッチを実装した際に色々な問題に突き当たったときの話

Avatar for mikesorae

mikesorae

July 15, 2020

More Decks by mikesorae

Other Decks in Programming

Transcript

  1. 背景 • 業務の一環でデータ分析レポートを作っている • 1000行〜3000行くらいのSQLをいつも書いている => つらい • 条件文をRubyで組み立てて最終結果もRubyで加工している •

    一ヶ月も経つと何やってるか読めない => つらい • window関数を多用するためデータ抽出が1回8時間〜24時間かかる => つらい • SQLをこね回しているので仕様変更/デバッグ/レビューも時間がかかる => つらい メンテナンス性とデータ処理時間改善のためBQ + Google Cloud Dataflowへ
  2. Apache Beamとは • Batch処理とStreaming処理の統合プログラミングモデル ◦ Batch + Stream = Beam

    • Dataflow用のプログラミングモデルだったが、Googleによって2016年にOSS化された • Dataflow以外にもSpark、Flink、Nemo、Samzaを実行環境として利用できる • JavaとPythonとGo版のSDKを提供している ◦ Go版はDataflowで動くのか未確認
  3. Apache Beamのプログラミングモデル(1) • Pipeline ◦ データの読み込みから出力まで含めた一連のデータ処理タスクを包括する概念 • PCollection ◦ データ処理の対象となる分散データセット

    ◦ テキストファイルのような有限データも、 Pub/Subのような無限データも扱える • PTransform ◦ Pipeline中の各ステップのデータ処理操作 ◦ 1〜複数個のPCollectionを受け取って0〜1個のPCollectionを返す
  4. WordCountの例 https://github.com/apache/beam/blob/master/examples/java/src/main/java/org/apache/beam/examples/MinimalWordCount.java Pipeline p = Pipeline.create(options); p.apply(TextIO.read().from("gs://apache-beam-samples/shakespeare/*")) .apply( FlatMapElements.into(TypeDescriptors.strings()) .via((String

    line) -> Arrays.asList(line.split("[^\\p{L}]+")))) .apply(Filter.by((String word) -> !word.isEmpty())) .apply(Count.perElement()) .apply( MapElements.into(TypeDescriptors.strings()) .via( (KV<String, Long> wordCount) -> wordCount.getKey() + ": " + wordCount.getValue())) .apply(TextIO.write().to("wordcounts"));
  5. 独自の処理関数の定義 // 文字数を数える処理の関数部 static class ComputeWordLengthFn extends DoFn<String, Integer> {

    @ProcessElement public void processElement(@Element String word, OutputReceiver<Integer> out) { // Use OutputReceiver.output to emit the output element. out.output(word.length()); } } Genericsやアノテーションなどで記述量がそこそこ多くなる JavaよりはKotlinの方が楽かと思い選定から除外
  6. JvmWildcardを付けないと実行時エラー typecheck時に required: Iterable<? extends MyValue> found: Iterable<MyValue> で実行時エラーになる。 解決方法:

    MyValueに@JvmWildCardをつけ る。 class MyFn : DoFn< KV<String, Iterable<MyValue>>, // InputT String // OutputT >() { ...
  7. Kotlin + Apache Beam • Tupleが使えない • 自作PairCoderを作ってみたけどJvmSuppresWildcardが必要 • lambdaの型チェックがうまくいかない(コンパイルは通るけど実行時に落ちる)

    • 結局型を自分で書かないといけない箇所が多い • 全体的にリフレクション周りでハマる • 期待していたよりもコード量が減らず、玄人作業になる部分が多かったため、採用を見送 り
  8. Apache BeamとScio • Pipeline => ScioContext • PCollection => SCollection

    • PTransform => .map / .flatMap / .filter / .reduce / .transform • PipelineResult => ScioResult • innerでApache Beamのインスタンスを取り出してApache Beamとして書くこともできる。 Kotlin/Java + Apache Beamよりも圧倒的に書きやすく、 型周りで特にハマることもなかったため採用
  9. Dataflowジョブの実行方式 • Dataflowのジョブの実行方式は大きく分けて2つある(※この後増える) ◦ Dataflow Runner ▪ ビルド/デプロイ/実行をワンセットで行う ▪ 毎回ビルドするので時間がかかる

    ◦ カスタムテンプレート ▪ ビルド/デプロイと実行を分けて行う ▪ 実行のたびにビルドを行わないので何度も実行するならこちらの方が良い ▪ 実行するためにGUIかDataflow APIを使う必要がある
  10. ValueProvider // カスタムテンプレートを使わない場合 public interface MyOptions : DataflowPipelineOptions { @Validation.Required

    fun getMyValue(): String } // カスタムテンプレートを使う場合 public interface MyOptions : DataflowPipelineOptions { @Validation.Required fun getMyValue(): ValueProvider[String] }
  11. Flex Templateの制約 • Dockerイメージのベースイメージは gcr.io/dataflow-templates-base/java8-template-launcher-base (またはjava11) にす る必要がある • 指定されたディレクトリにFat

    Jarを配置しなければならない ◦ 専用のlauncherから呼び出されており、class pathが指定できない ◦ ベースイメージのソースコードも公開されていない
  12. ScioでFat Jarを作る • Apache BeamのmvnではデフォルトでFat Jarだが、Scioでは自作する必要がある • Scioの公式FAQでは、マージ戦略が複雑になるためsbt-assemblyではなくsbt-packか sbt-native-packagerを推奨している •

    調べてみたところ、どちらもFat Jarではなくlibsにjarを配置するタイプだったため使えな かった • よくよくScioの公式のexampleを見るとbuild.sbtでsbt-assemblyを使っていた
  13. sbt-assemblyの設定 https://github.com/spotify/scio/blob/master/build.sbt lazy val assemblySettings = Seq( test in assembly

    := {}, assemblyMergeStrategy in assembly ~= { old => { case s if s.endsWith(".properties") => MergeStrategy.filterDistinctLines case s if s.endsWith("public-suffix-list.txt") => MergeStrategy.filterDistinctLines case s if s.endsWith("pom.xml") => MergeStrategy.last case s if s.endsWith(".class") => MergeStrategy.last case s if s.endsWith(".proto") => MergeStrategy.last case s if s.endsWith("libjansi.jnilib") => MergeStrategy.last ….
  14. Flex Templateのパラメータ渡しの仕組み 1. Dataflow APIのlaunch_flex_templateにパラメータを渡して実行 2. GCPがlauncher EC2インスタンスを作成 3. EC2のカスタムメタデータでパラメータ文字列をlauncherインスタンスに渡す

    a. user-dataでdocker起動用のsystemdのconfが渡される b. template-container-argsでdockerコマンドの引数が渡される c. template-container-argsの引数がsystemdのExecStartに渡される 4. systemdがdockerコンテナを起動する
  15. systemdの設定 #cloud-configs users: - name: cloudservice uid: 2000 write_files: -

    path: /etc/systemd/system/cloudservice.service permissions: 0644 owner: root content: | [Unit] Description=Start a dynamic template launcher docker container Wants=gcr-online.target After=gcr-online.target [Service] Environment="HOME=/home/cloudservice" ExecStartPre=/usr/bin/docker-credential-gcr configure-docker ExecStartPre=/usr/bin/docker run -v /var/log/dataflow/template_launcher:/var/log/dataflow/template_lau ncher -d gcr.io/dataflow-templates-base/template-launcher-logger:stable ExecStart=/usr/bin/docker run --entrypoint /opt/google/dataflow/java_template_launcher -v /var/log/dataflow/template_launcher:/var/log/dataflow/template_lau ncher gcr.io/xxxx/xxxx # ← この後ろにパラメータが追加される runcmd: - systemctl daemon-reload - systemctl start cloudservice.service
  16. ジョブグラフの制約 • ジョブグラフのサイズが 10 MB を超えてはいけません。パイプライン内の条件によって は、ジョブグラフが制限を超えることがあります。こうした条件のうち一般的なものは次の とおりです。 ◦ 大量のメモリ内データを含む

    Create 変換。 ◦ リモート ワーカーへの送信のためにシリアル化される大きな DoFn インスタンス。 ◦ シリアル化する大量のデータを(通常は誤って) pull する匿名内部クラス インスタンスとしての DoFn。 これらの状況を回避するには、パイプラインを再構築することを検討してください。 https://cloud.google.com/dataflow/docs/guides/common-errors
  17. Flex Template用REST APIのリファレンスがない • 公式リファレンスにまだ載っていない ◦ https://cloud.google.com/dataflow/docs/reference/rest?hl=ja ◦ githubでgoogle-api-ruby-clientの実装を調べた •

    サービスアカウントの指定方法が他のAPIと違う ◦ LaunchFlexTemplateRequestのlaunch_parameter.parameters.serviceAccountで指定 ◦ 他の人がGUIから操作していてたまたま発見
  18. 静的型付け / UnitTest • Ruby + SQLの場合一箇所直すたびにドキドキしながら再実行していた • 入出力が型レベルで保証されるようになったので手戻りが減った ◦

    レビューも楽になった • PipelineやFunction単位でUnit Testが書けるようになったため、安心感が大幅に増した • QAでの不具合件数も1桁台で抑えられた