ScalaでカスタムETLを始める方法
Scalaは、Sparkネイティブな開発、型安全性、そして単一のバッチジョブからより大きなデータ基盤へ成長できる実装を求めるときに強力な選択肢です。
ETLでScalaが向いている場面
ETLロジックが十分に複雑で、強い抽象化、再利用可能なライブラリ、コンパイル時チェックが保守コストを本当に下げる場合、Scalaは特に有効です。Apache Sparkとの相性も非常に高いです。
パイプラインを単なるスクリプト集ではなく、ソフトウェアとして扱いたいならScalaが適しています。
What Spark, Hadoop, and HDFS Actually Are
These technologies are related, but they are not the same layer of the stack.
Spark
Spark is a distributed compute engine. It parallelizes reads, joins, aggregations, and writes across many workers. For ETL, think of Spark as the execution layer rather than the long-term storage layer.
Hadoop
Hadoop is a broader ecosystem that historically bundled distributed storage, cluster resource management, and batch processing. Modern teams often borrow concepts from Hadoop without running a classic Hadoop cluster end to end.
HDFS
HDFS is Hadoop Distributed File System. It stores files across many machines and was the standard storage layer for self-hosted big data clusters before object storage became the default in many cloud environments.
Catalog and table metadata
As soon as multiple jobs and analysts rely on the same curated datasets, you also need a catalog or metastore. It keeps schemas, partitions, and table definitions consistent across your ETL and query tools.
Should You Use HDFS?
Usually not as a default. HDFS still has a place, but many modern ETL stacks use object storage instead.
- - Choose HDFS mainly when you operate your own cluster, want data locality, and already accept the operational cost of managing distributed storage nodes.
- - Prefer S3, ADLS, or GCS when your compute can be ephemeral, your storage should scale independently, or you want simpler disaster recovery and cross-service integration.
- - Do not adopt Hadoop just because you use Spark. Spark can run on Kubernetes, YARN, Databricks, EMR, Synapse, or other managed runtimes without HDFS being the center of the design.
正しい構造で始める
主要レイヤー
- - 設定: パス、シークレット、環境、出力先テーブル
- - Reader: JDBC、ファイル、キュー、APIアダプタ
- - Transform: 正規化、付加、重複排除、業務ルール
- - Validation: スキーマ、件数、ドメインルールの検証
- - Writer: lake、warehouse、またはサービス向け出力
初期設計の目標
- - 再実行で重複しないidempotent load
- - batchやpartitionを含む構造化ログ
- - extract・transform・load間の明確な失敗境界
- - 巨大なジョブではなく小さくテスト可能な変換
- - 初日から実行環境に合ったデプロイモデル
Do You Have to Use Spark?
No. Spark is powerful, but it is not the mandatory starting point for every ETL pipeline.
Use Spark when
You need distributed joins, large historical scans, partitioned batch outputs, or multi-hundred-gigabyte to multi-terabyte processing windows that no longer fit comfortably on one machine.
Skip Spark when
Your ETL is mostly SQL against Postgres, the data volume is still modest, one server can finish the workload inside the SLA, and operational simplicity matters more than horizontal scale.
Common alternatives
Start with Postgres SQL, dbt, Airflow plus Python, DuckDB, or Polars when the pipeline is small or medium. Move to Spark when the workload proves it needs a distributed compute engine.
プロジェクトの足場
最小構成のSpark ETLプロジェクトは、build.sbt、1つのmainクラス、一般的なsrc/main/scala構成から始められます。サンプルのバージョンは対象クラスターに合わせて調整してください。
name := "custom-etl"
version := "0.1.0"
scalaVersion := "2.13.17"
val sparkVersion = "4.1.2" // Replace to match your target cluster
libraryDependencies ++= Seq(
"org.apache.spark" %% "spark-sql" % sparkVersion % "provided",
"com.typesafe" % "config" % "1.4.3",
"org.scalatest" %% "scalatest" % "3.2.19" % Test
)最小ジョブの骨格
小さなスタータージョブでも、入力・出力パスを引数で受け取り、SparkSessionを作成し、絞った変換チェーンを実行し、Parquetのような決定的な形式へ書き出すべきです。
package dev.ryware.etl
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{col, to_timestamp}
object CustomEtlJob {
def main(args: Array[String]): Unit = {
val inputPath = args(0)
val outputPath = args(1)
val spark = SparkSession.builder()
.appName("custom-etl-job")
.getOrCreate()
val curated = spark.read
.option("header", "true")
.csv(inputPath)
.filter(col("transaction_id").isNotNull)
.withColumn("amount", col("amount").cast("decimal(18,2)"))
.withColumn("updated_at", to_timestamp(col("updated_at")))
curated.write
.mode("overwrite")
.parquet(outputPath)
spark.stop()
}
}When Postgres Is Enough and When to Move to a Warehouse
Raw database size matters, but query concurrency, data shape, retention, and the cost of long analytical scans matter more.
| Workload stage | When Postgres is still reasonable | When you should plan the move | Typical target |
|---|---|---|---|
| Early analytics | Up to roughly 50-100 GB of curated analytical data, a few internal dashboards, light concurrency, and daily batch updates. | Stay put unless queries already compete with transactional traffic or refresh windows are missing the SLA. | Postgres plus SQL, dbt, and simple orchestration. |
| Growing reporting stack | Roughly 100-300 GB, moderate joins, fewer than 10-20 active analysts, and mostly structured tables. | Plan a warehouse when refresh jobs grow brittle, vacuum and index tuning turns into constant work, or semi-structured data starts piling up. | Postgres can still work, but start evaluating BigQuery, Snowflake, Redshift, Synapse, or a lakehouse pattern. |
| Heavy analytics | Possible only with careful tuning, but the tradeoff gets worse once scans span hundreds of gigabytes, concurrency climbs, and retention grows fast. | Move when BI workloads, historical backfills, and ML feature extraction all want the same data at the same time. | Columnar warehouse or lakehouse with separated storage and compute. |
| Platform scale | Rarely the right long-term shape once you are in multi-terabyte territory, ingesting many domains, or supporting self-service analytics. | At this point the question is usually not whether to move, but which warehouse, lakehouse, governance model, and cost controls fit the business. | Warehouse or lakehouse platform with catalog, partitioning, and workload isolation. |
What If the Source Files Are XML and You Need XPath?
Spark can process XML, but the best approach depends on how irregular the documents are and how much XPath logic you need.
Spark can read XML
On open-source Spark, teams commonly add the spark-xml package and define an explicit row tag and schema. That works well when the XML documents are repetitive enough to map into tabular records.
XPath is possible, but expensive
XPath-style extraction is reasonable for a few known fields, but deeply nested or dynamic XPath rules can become CPU-heavy and awkward to test at scale. Flatten the documents as early as you can.
Sometimes pre-processing is better
If the XML is highly irregular, namespace-heavy, or validated against complex XSD rules, it can be cleaner to parse it first with a dedicated XML library in Scala or Java and then land normalized JSON or Parquet for downstream Spark work.
Operational advice
Keep the raw XML, version the mapping rules, capture malformed files separately, and build test fixtures from real samples. Special formats fail at the edges, so sample coverage matters more than happy-path demos.
次に本番で追加すべきこと
検証と品質
- - curated outputを書き出す前にスキーマ検証を追加する
- - 各実行ごとの行数やnull率を保存する
- - 壊れたpartitionを黙って補正せず、拒否または隔離する
テストとパッケージング
- - 小さなインメモリDataFrameで変換関数を単体テストする
- - ソース・シンクアダプタの統合テストを追加する
- - sbtでパッケージし、spark-submitやクラスタランチャーで実行する
最初の現実的なマイルストーン
最初の目標はデータ基盤全体の構築ではありません。1つのデータソースを確実に読み、1本のきれいな変換経路を作り、1つのcurated outputを書き出し、その実行を観測可能にすることです。
マイルストーン1
決定的な出力を持つ単一のバッチジョブ。
マイルストーン2
設定駆動の環境と検証ルール。
マイルストーン3
オーケストレーション、アラート、lineage、品質ベースライン。
関連ETLガイド
プラットフォーム選定、異常検知、AWS Glue固有の実装オプションもご覧ください。
AWS vs Azure vs GCP vs On-Premises for ETL
Compare managed ETL stacks, hybrid patterns, and the tools teams commonly use on each platform.
Anomaly Detection in ETL Pipelines
See which data and operational signals matter, how to baseline them, and how to react before bad data spreads.
AWS Glue for Anomaly Detection, Data Quality, and Debugging
Use Glue Data Quality, historical row-count checks, and run-time logging to catch ETL issues quickly.