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

Goでデータパイプラインを作ろう

 Goでデータパイプラインを作ろう

■ イベント
golang.tokyo #44
https://golangtokyo.connpass.com/event/399472/

■登壇概要
タイトル:Goでデータパイプラインを作ろう
登壇者:技術本部 Data Intelligence Engineering Unit 上島 剛

■ 技術本部 採用情報
https://media.sansan-engineering.com

Avatar for SansanTech

SansanTech PRO

August 05, 2026

More Decks by SansanTech

Other Decks in Technology

Transcript

  1. 上島 剛 (Tsuyoshi Ueshima) Sansan株式会社 技術本部 Data Intelligence Engineering Unit

    サーバーサイドエンジニア - Sansan Data Intelligence - 企業・事業所マスターデータ基盤 開発 趣味 - コーヒー、キャンプ、ガジェット 2
  2. Apache Beam と バッチ処理 とストリーミング処理 両方に対応した並列データ処理パイプラインを 「定義するため 統合モデル 」オープンソース -

    Go に native 対応(SDK - もともと Google Java・Python、そして Go) Dataflow プログラミングモデルがベース ◦ Google を中心としたオープンソースプロジェクト 7
  3. パイプラインを実行するランナー 同じパイプライン定義を、様々な分散処理バックエンドで実行できる (以下 Google Cloud Dataflow Direct Runner - 各バックエンド互換

    Apache Spark Runner 一部 ) Apache Flink Runner APIに変換されて実行 → ポータビリティに優れる - ただし Spark・Flink 固有 機能 Beam から 使えない制限もある 8
  4. Sansan Data Intelligence マスターデータ構築で採用 - Sansan Data Intelligence 企業・事業所データ ◦

    企業 取引先データを最新・正確に整える データ品質管理サービス 約1,000万件 - 国税庁・各省庁 オープンデータや企業ウェブサイ トなど複数 ソースを統合し、企業・事業所 以上を扱うバッチ処理 データパイプライン マス ターデータ (DB)を構築 - Apache Beam Go SDK で開発 - Google Cloud Dataflow 上で運用 9
  5. データパイプライン 全体像とコンセプト 例) Google Cloud Storage から読み込み、変換して BigQuery へ書き出すまで 流れ

    GCS TextIO PCollection ParDo PCollection BigQueryIO BigQuery Source I/O Connector 生データ (string) PTransform 変換後データ (struct) I/O Connector Sink これ全体が 1 つ Pipeline I/O Connector PCollection PTransform データ 入出力(Source / Sink)を パイプラインを流れるデータ 集合 データを変換する処理 担う ParDoを使うと、Go 関数を実装 できる 11
  6. 実行してみる ① GCS 上 入力ファイル gs://my-bucket/input/records.jsonl {"name":"alice","score":100} {"name":"bob","score":200} {"name":"charlie","score":300} ③

    BigQuery に書き込まれたデータ my-project:my_dataset.records name score alice 100 bob 200 charlie 300 ② 実行コマンド go run . \ --runner direct \ --project my-project \ 3件 --input \ 書き込まれる gs://my-bucket/input/records.jsonl \ レコードがそ まま BigQuery テーブルに --table \ my-project:my_dataset.records 14
  7. PTransform と I/O Connector 主な PTransform - ParDo:基本 - 変換。

    Go I/O Connector - TextIO(GCS / S3 / local) GroupByKey:キーごとに集約 - BigQuery / Spanner Flatten:複数 - Database(MySQL / PostgreSQL 互換) - Parquet … 関数を実装 PCollection を統合 15
  8. ドメインコードを、サービスとパイプラインで共有できる 前提 - 並列度調整・シャーディング・順序/状態管理・リトライ → Goで ビジネスロジック 実装に集中できる サービスとパイプラインがドメインモデル・ロジック・型を共通化できる 二重管理がなく、修正漏れによる振る舞い

    - Beam / Dataflow(Runner) に任せられる cgo で C ライブラリ依存 独自 ズレがなくなる(データ構造を変えても壊れる箇所をコンパイラが検出) コードも共通資産として扱える 複雑な住所正規化を行う共通ライブラリを、サービスとパイプライン 両方で利用できている https://buildersbox.corp-sansan.com/entry/2025/06/20/120000 17
  9. ローカルでテストできる - 処理 通常 Go 関数として書ける で go test で検証できる

    - パイプライン全体もローカル( Direct Runner)で検証できる go test 上で使えるテストヘルパーが用意されている 18
  10. 注意点:他 SDKと 比較 Go SDK 基本機能 例:BigQueryIO - 個人的に 揃うが、一部機能

    未サポート Storage Read API、SpannerIO 実務上 MutationGroup など 支障 感じていない - 必要なら Cross-language transforms で Java SDK等 を内部的に呼び出せる - Go で完結したけれ 自作もできる 19
  11. まとめ - Apache Beam なら、Go でデータパイプラインを書ける Go native SDK。バッチもストリーミングも扱える -

    サービスとコード資産を共有でき、ローカルでテストもできる 同じ Go でドメインコードを共通化/go test でそ まま検証 - 1,000万件規模でも、実際に運用できている Sansan - 次 企業マスターデータ基盤で実運用中 パイプライン、 Go で書いてみませんか 20
  12. Appendix 参考文献 Apache Beam - Apache Beam Overview Google Cloud

    Dataflow - https://beam.apache.org/get-started/beam-overview/ - Apache Beam Go SDK https://beam.apache.org/documentation/sdks/go/ - Dataflow 概要 https://cloud.google.com/dataflow/docs/overview - Go パイプラインを作成する https://cloud.google.com/dataflow/docs/guides/create-pipeline-go Beam Programming Guide https://beam.apache.org/documentation/programming-guide/ - Runner Capability Matrix https://beam.apache.org/documentation/runners/capability-matrix/ 21
  13. 22