Apache Beam vs Spark 2026年版:統一パイプラインと面接質問の完全ガイド

Apache Beam 2.76とSpark 4.2を徹底比較。ポータビリティ、ウィンドウ処理、パフォーマンスの違いを解説し、データエンジニアリング面接で頻出の質問と回答例を紹介します。

Apache Beam vs Spark 2026年版:統一パイプラインと面接質問の完全ガイド

Apache Beam vs Sparkは、現代のデータエンジニアリングにおける最も一般的なアーキテクチャ選択の一つです。Beam 2.76(2026年8月)とSpark 4.2(2026年7月)はどちらもバッチ処理とストリーミングワークロードを処理できますが、設計思想は根本的に異なります。Beamは実行エンジンを抽象化し、Sparkは緊密に統合されたランタイムを提供します。

クイック判断フレームワーク

Beamを選択するのは、ランナー間(Dataflow、Flink、Spark)のポータビリティが重要な場合、またはGoogle Cloud Dataflowを使用する場合です。Sparkを選択するのは、自己管理クラスターを運用する場合、Databricksを使用する場合、またはMLlibとのML統合が必要な場合です。

Beamのポータビリティモデル vs Sparkの統合エンジン

Apache Beamはプログラミングモデルと実行を分離します。単一のパイプライン定義が、コード変更なしでGoogle Cloud Dataflow、Apache Flink、Apache Spark、その他のランナー上で実行できます。この抽象化は、Beam SDKが互換性のあるランナーが解釈するポータブルなパイプライン表現を生成することで実現されます。

Spark 4.2は逆のアプローチを取ります。DataFrame API、Structured Streaming、MLlibは同じCatalystオプティマイザとTungsten実行エンジンを共有します。この緊密な結合により、実際のデータ統計に基づいて実行時にプランを調整するAdaptive Query Executionのような最適化が可能になります。

python
# beam_pipeline.py
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

# Same code runs on Dataflow, Flink, or Spark runner
options = PipelineOptions([
    '--runner=DataflowRunner',  # Switch to FlinkRunner or SparkRunner
    '--project=my-project',
    '--region=us-central1',
    '--temp_location=gs://my-bucket/temp'
])

with beam.Pipeline(options=options) as pipeline:
    (pipeline
     | 'ReadEvents' >> beam.io.ReadFromPubSub(topic='projects/p/topics/events')
     | 'ParseJSON' >> beam.Map(lambda x: json.loads(x))
     | 'FilterValid' >> beam.Filter(lambda e: e.get('status') == 'valid')
     | 'WindowByMinute' >> beam.WindowInto(beam.window.FixedWindows(60))
     | 'CountPerWindow' >> beam.combiners.Count.Globally()
     | 'WriteToBQ' >> beam.io.WriteToBigQuery('project:dataset.table'))

Sparkの同等のコードは、Sparkランタイムに直接結びついています:

python
# spark_streaming.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, window

spark = SparkSession.builder \
    .appName("EventProcessing") \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

# Spark 4.2: ANSI mode enabled by default, stricter type checking
events = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "broker:9092") \
    .option("subscribe", "events") \
    .load()

processed = events \
    .select(from_json(col("value").cast("string"), schema).alias("data")) \
    .filter(col("data.status") == "valid") \
    .groupBy(window(col("data.timestamp"), "1 minute")) \
    .count()

processed.writeStream \
    .format("bigquery") \
    .option("table", "project.dataset.table") \
    .outputMode("append") \
    .start()

Beamのポータビリティはクラウド非依存のアーキテクチャを可能にしますが、変換レイヤーが追加されます。Sparkの直接実行は、同じハードウェア上での同等の操作において、通常より低いレイテンシを示します。

ウィンドウ処理とイベント時間処理の比較

両フレームワークともイベント時間セマンティクスを処理しますが、そのAPIは異なる背景を反映しています。BeamのウィンドウモデルはDataflow Modelペーパー(2015年)に由来し、ウィンドウをパイプラインの第一級要素として扱います。SparkはStructured Streaming向けにウィンドウ処理を適応させ、DataFrame APIと統合しました。

機能Beam 2.76Spark 4.2
固定ウィンドウFixedWindows(duration)window(col, duration)
スライディングウィンドウSlidingWindows(size, period)window(col, size, period)
セッションウィンドウSessions(gap)ネイティブ非対応(flatMapGroupsWithStateを使用)
カスタムウィンドウWindowFnサブクラス限定的
遅延データ処理組み込みトリガーウォーターマーク遅延
許容遅延ウィンドウごとの設定グローバルウォーターマーク

セッションウィンドウは最も明確に違いを示します。Beamはセッションをネイティブのウィンドウ戦略として扱います:

python
# beam_sessions.py
from apache_beam import window

# 30-minute session gap, allow 1 hour late data
windowed = (
    events
    | 'SessionWindow' >> beam.WindowInto(
        window.Sessions(30 * 60),  # 30 min gap closes session
        trigger=beam.trigger.AfterWatermark(
            early=beam.trigger.AfterProcessingTime(60),
            late=beam.trigger.AfterCount(1)
        ),
        allowed_lateness=3600,  # Accept data up to 1 hour late
        accumulation_mode=beam.trigger.AccumulationMode.ACCUMULATING
    )
)

Sparkではセッション処理にステートフル処理が必要です:

python
# spark_sessions.py
from pyspark.sql.streaming import GroupState, GroupStateTimeout

def update_session(key, events, state: GroupState):
    # Manual session management with state API
    session_data = state.getOption() or {"count": 0, "start": None, "end": None}
    
    for event in events:
        ts = event.timestamp
        if session_data["end"] and (ts - session_data["end"]).seconds > 1800:
            # Gap exceeded 30 min, emit previous session
            yield session_data
            session_data = {"count": 0, "start": ts, "end": ts}
        
        session_data["count"] += 1
        session_data["end"] = ts
        if not session_data["start"]:
            session_data["start"] = ts
    
    state.update(session_data)
    state.setTimeoutDuration(30 * 60 * 1000)  # 30 min timeout

# Spark 4.2 Arbitrary State API v2
result = events \
    .groupByKey(lambda e: e.user_id) \
    .flatMapGroupsWithState(
        update_session,
        outputMode="append",
        stateType=session_schema,
        timeoutConf=GroupStateTimeout.ProcessingTimeTimeout
    )

これらのトピックに関する面接対策については、Apache BeamとDataflow面接質問モジュールを参照してください。

パフォーマンス:2026年のベンチマークデータ

BeamはSpark上でランナーオプションの一つとして実行されるため、直接比較には注意が必要です。関連する比較はBeam-on-Dataflow vs ネイティブSparkとなります。

DatabricksとGoogle Cloudからの最近のベンチマークは以下を示しています:

ワークロードSpark 4.2 (Databricks)Beam 2.76 (Dataflow)備考
バッチETL(1TB Parquet)4.2分5.1分SparkのPhotonエンジンの優位性
ストリーミング(100Kイベント/秒)45ms p99レイテンシ120ms p99レイテンシDataflowオートスケーリングのオーバーヘッド
Exactly-onceシンク書き込みネイティブネイティブ2024年以降両方サポート
コスト(継続ワークロード)$0.12/GB処理$0.08/GB処理Dataflow Flex価格

Spark 4.2のパフォーマンス向上は以下の機能から生まれています:

  • ANSIモードがデフォルト:より厳格なSQLセマンティクスでエラーを早期検出
  • VARIANTデータ型:ネイティブの半構造化データ処理
  • Adaptive Query Execution:実行時プラン最適化
  • Java 21サポート:仮想スレッドでオーバーヘッド削減

Beam 2.76は以下で対抗します:

  • Flink 2.0ランナーサポート:本番環境向けステートフルストリーミング
  • CDCオフセット永続化:DebeziumIO FileSystemOffsetRetainer
  • ADK統合:Python SDKでのGoogle Agent Development Kitサポート

Data Engineeringの面接対策はできていますか?

インタラクティブなシミュレーター、flashcards、技術テストで練習しましょう。

面接でよく聞かれる質問:Beam vs Spark

データエンジニアリング職の技術面接では、これらのフレームワークを頻繁に比較します。以下は2026年の面接で出題される質問と、期待される回答の深さです。

Q1: ネイティブSparkではなくBeamを選択するのはどのような場合ですか?

期待される回答ポイント:

  • パイプラインコードを変更せずに実行する必要があるマルチクラウドまたはハイブリッドデプロイメント
  • Dataflowがマネージドインフラストラクチャを提供するGoogle Cloud環境
  • セッションウィンドウやカスタムトリガーを必要とする複雑なイベント時間セマンティクス
  • Dataflowの背景から既存のBeam専門知識を持つチーム

問題のある回答:「Beamはポータブルなので常に優れている。」ポータビリティにはオーバーヘッドコストがあります。

Q2: Beamのランナー抽象化はデバッグにどのように影響しますか?

期待される回答ポイント:

  • スタックトレースはBeam SDKとランナー実装の両方を参照
  • ランナー固有の最適化は異なる動作をする可能性がある(Flinkチェックポイント vs Sparkチェックポイント)
  • Metrics APIは統一監視を提供するが、ランナーダッシュボードは異なる詳細を表示
  • 本番ランナーへのデプロイ前にDirectRunnerでテスト

Q3: 両フレームワークでのExactly-onceセマンティクスを説明してください。

期待される回答:

python
# Beam: exactly-once via runner guarantees
# Dataflow provides exactly-once for both sources and sinks
# The SDK handles deduplication and checkpoint coordination

with beam.Pipeline() as p:
    (p
     | beam.io.ReadFromPubSub(subscription='...')  # Exactly-once read
     | beam.Map(process)
     | beam.io.WriteToBigQuery(...)  # Exactly-once write with retries
    )

# Spark: exactly-once via checkpointing and idempotent sinks
spark.readStream \
    .format("kafka") \
    .load() \
    .writeStream \
    .option("checkpointLocation", "/checkpoint")  # State recovery
    .foreachBatch(idempotent_write)  # Application-level deduplication
    .start()

その他のストリーミング面接トピックについては、PySparkモジュールをご覧ください。

Q4: SparkバッチジョブをBeamに移行するにはどうしますか?

この質問は両APIの理解をテストします。主要なポイント:

  1. DataFrame操作をPCollectionとTransformにマッピング
  2. spark.readを適切なBeam I/Oコネクタに置き換え
  3. UDFをbeam.Mapまたはbeam.ParDo関数に変換
  4. パーティショニングを異なる方法で処理(BeamのReshuffle vs Sparkのrepartition
  5. 本番ランナーへのデプロイ前にDirectRunnerでテスト

適切なツールの選択:決定マトリックス

選択は技術的な能力よりも組織的なコンテキストに依存します。両フレームワークともほとんどのデータエンジニアリングワークロードを適切に処理します。

要因Beamが有利Sparkが有利
クラウドプロバイダーGoogle CloudAWS EMR、Databricks、オンプレミス
チームの専門知識既存のDataflow経験既存のSpark/PySparkスキル
ML統合限定的(別ツール)MLlib、Spark ML
インタラクティブ分析設計されていないSpark SQL、ノートブック
セッションウィンドウネイティブサポート手動ステート管理
コストモデル従量課金(Dataflow)クラスタープロビジョニング
ベンダーロックイン低い(複数ランナー)高い(Spark固有コード)

ETL/ELTパターンの決定については、フレームワークを選択する前にデータ量とレイテンシ要件を検討してください。

実世界のアーキテクチャ:ハイブリッドアプローチ

多くの組織は両フレームワークを使用しています。一般的なパターン:

text
[ストリーミング取り込み]     [バッチ処理]     [ML訓練]
        |                        |                    |
   Beam/Dataflow           Spark on Databricks    Spark MLlib
        |                        |                    |
        v                        v                    v
    BigQuery  <----- dbt -----> Delta Lake -----> Model Registry

BeamはDataflowのオートスケーリングがトラフィックパターンに合うストリーミング取り込みを処理します。Sparkはクラスター経済が持続的コンピューティングに有利なバッチワークロードを処理します。両方が統一データウェアハウスにフィードします。

このアーキテクチャはApache Sparkチュートリアルで実装の詳細とともに登場します。

今すぐ練習を始めましょう!

面接シミュレーターと技術テストで知識をテストしましょう。

Beam vs Spark選択の重要ポイント

  • Beam 2.76はDataflow、Flink 2.0、Spark間のランナーポータビリティを提供し、デプロイの柔軟性と引き換えにパフォーマンスをトレードオフ
  • Spark 4.2はVARIANT型やAdaptive Query Executionなどの機能で、SQL、ストリーミング、MLワークロード間のより緊密な統合を実現
  • セッションウィンドウと複雑なトリガーはBeamのネイティブウィンドウモデルが有利
  • 同等のハードウェアでのバッチパフォーマンスは、Catalyst/Tungsten最適化によりSparkが通常優位
  • 面接質問はどちらかのフレームワークを優れていると宣言するよりもトレードオフに焦点
  • 両フレームワークを使用するハイブリッドアーキテクチャは本番環境で一般的
  • コスト比較はワークロードパターンに依存:DataflowのGB単位課金 vs Sparkクラスタープロビジョニング
今日のチャレンジ

Data Engineering のバグを見つけられますか

実際のコード、隠れたバグ、1日1回。アカウントなしで試せます。

Anthony Fillion-Maillet

執筆

Anthony Fillion-Maillet

SharpSkill 創業者

10 年以上フルスタック開発に携わっています。SharpSkill を運営し、ここで公開される内容に責任を負っています。

2026年9月10日 更新

共有

関連記事