Apache Flink 2026完全ガイド:ストリーム処理、イベント時間、面接対策

Apache Flink 2.3のストリーム処理アーキテクチャ、イベント時間セマンティクス、ウィンドウ処理、状態管理を解説。データエンジニアリング面接で頻出する質問と回答も網羅。

Apache Flink ストリーム処理アーキテクチャ図

Apache Flink 2.3は、Exactly-Once セマンティクスとサブ秒レイテンシでストリーム処理を大規模に実行する。バッチ指向のシステムとは異なり、Flinkはデータが到着した時点で継続的に処理を行うため、リアルタイム分析、不正検知、イベント駆動アーキテクチャにおいて最適なフレームワークとなっている。

面接の重要ポイント

FlinkがSpark Streamingと異なる点は真のストリーム処理にある。Flinkはイベント時間セマンティクスでイベントを1つずつ処理するが、Spark Streamingはマイクロバッチを処理時間をデフォルトとして処理する。

Flinkは分散アーキテクチャ上で動作し、JobManagerが複数のTaskManager間の作業を調整する。各TaskManagerはジョブの並列オペレーターの一部を実行するタスクスロットを持つ。この分離により、分散チェックポイントを通じて耐障害性を維持しながら、水平スケールが可能となる。

Flinkのデータフローモデルは、計算を有向非巡回グラフ(DAG)として表現する。データはソースから変換を経てシンクへと流れ、各オペレーターは複数の並列インスタンスで実行される可能性がある。

FlinkStreamJob.javajava
// Basic Flink streaming application setup
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// Enable checkpointing every 10 seconds for fault tolerance
env.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE);

// Configure state backend for large state
env.setStateBackend(new EmbeddedRocksDBStateBackend());

// Define the data source - Kafka in production scenarios
DataStream<String> rawStream = env.addSource(
    new FlinkKafkaConsumer<>("events", new SimpleStringSchema(), kafkaProps)
);

// Parse and transform the stream
DataStream<Event> events = rawStream
    .map(json -> objectMapper.readValue(json, Event.class))
    .assignTimestampsAndWatermarks(
        WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
            .withTimestampAssigner((event, timestamp) -> event.getTimestamp())
    );

上記のチェックポイント設定は、10秒ごとに分散状態のスナップショットを保存する。障害が発生した場合、Flinkは最後に完了したチェックポイントから復元し、Kafkaからレコードを再生する。

イベント時間と処理時間のセマンティクス

イベント時間とは、イベントが実際に発生した時刻であり、データ自体に埋め込まれている。処理時間とは、Flinkがレコードを処理する時刻である。ネットワーク遅延、順序の乱れ、処理のバックログにより、処理時間は時間ベースの操作において信頼性が低くなるため、この区別は重要となる。

EventTimeExample.javajava
// Configuring event time with watermarks
public class EventTimeProcessor {
    
    public DataStream<AggregatedMetric> processWithEventTime(
            DataStream<SensorReading> readings) {
        
        return readings
            // Extract timestamp from the event payload
            .assignTimestampsAndWatermarks(
                WatermarkStrategy
                    .<SensorReading>forBoundedOutOfOrderness(Duration.ofMinutes(2))
                    .withTimestampAssigner((reading, ts) -> reading.getEventTime())
                    .withIdleness(Duration.ofMinutes(5)) // Handle idle partitions
            )
            // Window by event time, not wall clock
            .keyBy(SensorReading::getSensorId)
            .window(TumblingEventTimeWindows.of(Time.minutes(1)))
            .aggregate(new AverageAggregator());
    }
}

forBoundedOutOfOrderness戦略は、イベントが最大2分遅れて到着する可能性があることをFlinkに伝える。ウォーターマークは、ウォーターマーク以前のタイムスタンプを持つイベントがもう到着しないとFlinkが判断した時点で進む。

よくある面接質問

Flinkで遅延イベントはどうなるか?デフォルトでは、ウォーターマークがウィンドウの終了時刻を過ぎた後に到着したイベントは破棄される。遅延到着を処理するには.allowedLateness(Time.minutes(10))で許容遅延を設定するか、サイドアウトプットを使用して別途処理する。

リアルタイム分析のためのウィンドウ戦略

Flinkは4種類のウィンドウを提供する:タンブリング、スライディング、セッション、グローバルウィンドウ。それぞれ異なる分析ニーズに対応する。

WindowingStrategies.javajava
// Different windowing approaches for stream processing
public class WindowingStrategies {

    // Tumbling windows: fixed-size, non-overlapping
    // Use case: hourly aggregations, daily summaries
    public DataStream<Summary> tumblingAggregation(DataStream<Transaction> txns) {
        return txns
            .keyBy(Transaction::getAccountId)
            .window(TumblingEventTimeWindows.of(Time.hours(1)))
            .sum("amount");
    }

    // Sliding windows: fixed-size, overlapping
    // Use case: moving averages, rolling metrics
    public DataStream<Double> slidingAverage(DataStream<Metric> metrics) {
        return metrics
            .keyBy(Metric::getCategory)
            .window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)))
            .aggregate(new AverageAggregator());
    }

    // Session windows: activity-based, variable size
    // Use case: user sessions, conversation analysis
    public DataStream<Session> sessionAnalysis(DataStream<ClickEvent> clicks) {
        return clicks
            .keyBy(ClickEvent::getUserId)
            .window(EventTimeSessionWindows.withGap(Time.minutes(30)))
            .process(new SessionBuilder());
    }
}

セッションウィンドウは、設定可能な非アクティブ期間の後に閉じる。このパターンは、エンゲージメントに応じてセッション長が変動するユーザー行動分析に適している。

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

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

状態管理とチェックポイント

Flinkは処理全体を通じてオペレーター状態とキー付き状態を維持する。キー付き状態はキーごとにデータを分割し、関連レコードをまとめながら並列処理を可能にする。オペレーター状態はオペレーターインスタンス全体に適用される。

StatefulProcessor.javajava
// Managing state in a Flink KeyedProcessFunction
public class FraudDetector extends KeyedProcessFunction<String, Transaction, Alert> {
    
    // Keyed state: one value per key (account)
    private ValueState<Double> lastAmountState;
    private ValueState<Long> lastTransactionTimeState;
    private MapState<String, Integer> merchantCountState;
    
    @Override
    public void open(Configuration parameters) {
        // Initialize state descriptors
        lastAmountState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("lastAmount", Double.class));
        lastTransactionTimeState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("lastTime", Long.class));
        merchantCountState = getRuntimeContext().getMapState(
            new MapStateDescriptor<>("merchantCounts", String.class, Integer.class));
    }
    
    @Override
    public void processElement(Transaction txn, Context ctx, Collector<Alert> out) 
            throws Exception {
        Double lastAmount = lastAmountState.value();
        Long lastTime = lastTransactionTimeState.value();
        
        // Detect suspicious patterns
        if (lastAmount != null && lastTime != null) {
            long timeDelta = txn.getTimestamp() - lastTime;
            // Flag transactions 10x larger than previous within 1 minute
            if (txn.getAmount() > lastAmount * 10 && timeDelta < 60000) {
                out.collect(new Alert(txn.getAccountId(), "SUSPICIOUS_SPIKE", txn));
            }
        }
        
        // Update state for next transaction
        lastAmountState.update(txn.getAmount());
        lastTransactionTimeState.update(txn.getTimestamp());
        
        // Track merchant frequency
        Integer count = merchantCountState.get(txn.getMerchantId());
        merchantCountState.put(txn.getMerchantId(), (count == null ? 0 : count) + 1);
    }
}

このステートフルプロセッサは、アカウントごとのトランザクションパターンを追跡する。状態はチェックポイント間で永続化され、不正検知のコンテキストを失うことなく障害を乗り越える。

Flink 2.3は、増分ビュー保守のためのMaterialized Tablesにより SQL機能を拡張している。Table APIはバッチ処理とストリーム処理の統一インターフェースを提供する。

sql
-- flink_sql_streaming.sql
-- Create a streaming source table from Kafka
CREATE TABLE orders (
    order_id STRING,
    customer_id STRING,
    product_id STRING,
    amount DECIMAL(10, 2),
    order_time TIMESTAMP(3),
    -- Define watermark for event time processing
    WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'orders',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json',
    'scan.startup.mode' = 'earliest-offset'
);

-- Streaming aggregation with tumbling window
SELECT 
    customer_id,
    TUMBLE_START(order_time, INTERVAL '1' HOUR) AS window_start,
    COUNT(*) AS order_count,
    SUM(amount) AS total_amount
FROM orders
GROUP BY 
    customer_id,
    TUMBLE(order_time, INTERVAL '1' HOUR);

SQLアプローチは、SQLに精通したアナリストの開発を簡素化し、Flinkが裏でストリーム処理の複雑さを処理する。

両フレームワークはストリーミングデータを処理するが、アーキテクチャが根本的に異なる。Flinkは真のストリーミングでレコードを個別に処理し、Sparkはマイクロバッチを処理する。Apache Sparkとの比較において、本番環境ではレイテンシと一貫性のトレードオフが重要となる。

観点FlinkSpark Structured Streaming
処理モデル真のストリーミングマイクロバッチ
レイテンシミリ秒秒(バッチ間隔)
状態バックエンドRocksDB、HashMapインメモリ、HDFS
Exactly-Onceチェックポイントでネイティブ冪等シンクが必要
イベント時間ファーストクラスサポート2.1以降サポート
SQLサポート完全なストリーミングSQL限定的なウィンドウ
面接のポイント

ストリーミングにおけるFlink vs Sparkについて聞かれた場合、ユースケースへの適合性に焦点を当てる。Flinkは低レイテンシのイベント処理と複雑なイベントパターンに優れている。Spark Streamingは、バッチ処理ですでにSparkを使用しており、バッチとストリームの統一処理が必要な組織に適している。

よくあるFlink面接の質問と回答

FlinkはどのようにExactly-Onceセマンティクスを実現しているか?

Flinkは、トランザクションをサポートするシンクに対して、チェックポイントと2フェーズコミットを組み合わせている。チェックポイント中、Flinkはオペレーター状態のスナップショットを取り、ソースオフセットを記録する。Kafkaシンクの場合、FlinkはレコードをKafkaにプリコミットし、チェックポイントを完了してからトランザクションをコミットする。チェックポイント完了前に障害が発生した場合、コミットされていないレコードは破棄され、最後のチェックポイントから処理が再開される。

マルチソーストポロジーにおけるウォーターマークの伝播を説明せよ

ジョブが複数のパーティションまたはソースから読み取る場合、それぞれが受信イベントに基づいて独自のウォーターマークを生成する。下流オペレーターのウォーターマークは、すべての入力チャネルのウォーターマークの最小値と等しくなる。これにより、高速なパーティションが低速なパーティションより先に進んでも、ウィンドウが早期に閉じることがない。一部のパーティションがデータ送信を停止した場合にウォーターマークを進めるには、withIdleness()を設定する。

Flinkでバックプレッシャーの原因は何か、どのように診断するか?

バックプレッシャーは、下流オペレーターが上流のデータレートに追いつけない場合に発生する。Flink Web UIはオペレーターごとのバックプレッシャーステータスを表示する。一般的な原因:

  • 外部システムコールの遅延(データベースクエリ、APIコール)
  • map/process関数での高コストな計算
  • データ量に対して不十分な並列性
  • 大きな状態操作が処理をブロック

並列性の増加、遅い操作の最適化、外部コールへの非同期I/Oの使用で対処する。

セーブポイントとチェックポイントの違いは?

チェックポイントは自動的、増分的で、障害回復に最適化されている。Flinkがライフサイクルを管理し、古いものは自動的に削除される。セーブポイントはユーザーがトリガーする完全なスナップショットで、新しいコードのデプロイ、ジョブのリスケール、クラスター間の移行などの運用タスクを目的としている。セーブポイントは明示的に削除されるまで永続化され、スキーマの進化をサポートする。

KubernetesへのFlinkデプロイ

Flink Kubernetes Operator 1.15は、FlinkDeploymentカスタムリソースによりデプロイを簡素化する。ジョブのライフサイクル、アップグレード、スケーリングを処理する。

yaml
# flink-deployment.yaml
# Kubernetes deployment for a Flink application
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
  name: fraud-detection-job
spec:
  image: flink:1.20
  flinkVersion: v1_20
  flinkConfiguration:
    taskmanager.numberOfTaskSlots: "4"
    state.backend.type: rocksdb
    state.checkpoints.dir: s3://flink-checkpoints/fraud-detection
    execution.checkpointing.interval: "30s"
  serviceAccount: flink
  jobManager:
    resource:
      memory: "2048m"
      cpu: 1
  taskManager:
    resource:
      memory: "4096m"
      cpu: 2
    replicas: 3
  job:
    jarURI: s3://flink-artifacts/fraud-detection-1.0.jar
    entryClass: com.example.FraudDetectionJob
    parallelism: 12
    upgradeMode: savepoint

upgradeMode: savepoint設定により、オペレーターはアップグレード前にセーブポイントを取得し、デプロイ間で状態を保持する。

本番環境向けFlinkアプリケーションの最適化

本番デプロイでは、並列性、メモリ、状態バックエンドの設定に注意が必要である。統合の考慮事項についてはETLとデータパイプラインパターンを参照。

ProductionConfig.javajava
// Production-ready Flink configuration
public class ProductionConfig {
    
    public static void configureForProduction(StreamExecutionEnvironment env) {
        // Checkpoint configuration
        CheckpointConfig checkpointConfig = env.getCheckpointConfig();
        checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        checkpointConfig.setMinPauseBetweenCheckpoints(30000); // 30 seconds
        checkpointConfig.setCheckpointTimeout(600000); // 10 minutes
        checkpointConfig.setMaxConcurrentCheckpoints(1);
        checkpointConfig.setExternalizedCheckpointRetention(
            ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);
        
        // State backend with incremental checkpoints
        EmbeddedRocksDBStateBackend rocksDB = new EmbeddedRocksDBStateBackend(true);
        rocksDB.setDbStoragePath("/tmp/rocksdb");
        env.setStateBackend(rocksDB);
        
        // Restart strategy with exponential backoff
        env.setRestartStrategy(RestartStrategies.exponentialDelayRestart(
            Duration.ofSeconds(1),  // Initial delay
            Duration.ofMinutes(5),  // Max delay
            2.0,                    // Backoff multiplier
            Duration.ofHours(1),    // Reset backoff after
            0.1                     // Jitter
        ));
    }
}

増分チェックポイントは、最後のチェックポイント以降に変更された状態のみを書き込むことでチェックポイントサイズを削減する。この最適化は、ギガバイト単位のキー付き状態を管理する際に重要となる。

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

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

Apache Flinkストリーム処理の重要ポイント

  • Flink 2.3はミリ秒のレイテンシでイベントを個別に処理し、マイクロバッチシステムとは異なる
  • ウォーターマークを伴うイベント時間セマンティクスは順序の乱れたデータを正しく処理し、遅延イベントに関するよくある面接質問に答える
  • キー付き状態は関連レコードをまとめながら並列処理のためにデータを分割する
  • チェックポイントは分散スナップショットと2フェーズコミットを通じてExactly-Once保証を提供する
  • Kubernetes Operatorはセーブポイントベースの状態保持によりデプロイ、スケーリング、アップグレードを自動化する
  • サブ秒のレイテンシや複雑なイベント処理パターンが重要な場合、Spark StreamingよりFlinkを選択する
  • 大きな状態を持つ本番ワークロードには、増分チェックポイントを備えたRocksDBを設定する
今日のチャレンジ

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

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

Anthony Fillion-Maillet

執筆

Anthony Fillion-Maillet

SharpSkill 創業者

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

2026年8月28日 更新

タグ

#Apache Flink
#ストリーム処理
#イベント時間
#データエンジニアリング
#面接対策

共有

関連記事