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

Apache Flink 2.3は、Exactly-Once セマンティクスとサブ秒レイテンシでストリーム処理を大規模に実行する。バッチ指向のシステムとは異なり、Flinkはデータが到着した時点で継続的に処理を行うため、リアルタイム分析、不正検知、イベント駆動アーキテクチャにおいて最適なフレームワークとなっている。
FlinkがSpark Streamingと異なる点は真のストリーム処理にある。Flinkはイベント時間セマンティクスでイベントを1つずつ処理するが、Spark Streamingはマイクロバッチを処理時間をデフォルトとして処理する。
Flink 2.3のストリーム処理アーキテクチャ
Flinkは分散アーキテクチャ上で動作し、JobManagerが複数のTaskManager間の作業を調整する。各TaskManagerはジョブの並列オペレーターの一部を実行するタスクスロットを持つ。この分離により、分散チェックポイントを通じて耐障害性を維持しながら、水平スケールが可能となる。
Flinkのデータフローモデルは、計算を有向非巡回グラフ(DAG)として表現する。データはソースから変換を経てシンクへと流れ、各オペレーターは複数の並列インスタンスで実行される可能性がある。
// 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がレコードを処理する時刻である。ネットワーク遅延、順序の乱れ、処理のバックログにより、処理時間は時間ベースの操作において信頼性が低くなるため、この区別は重要となる。
// 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種類のウィンドウを提供する:タンブリング、スライディング、セッション、グローバルウィンドウ。それぞれ異なる分析ニーズに対応する。
// 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は処理全体を通じてオペレーター状態とキー付き状態を維持する。キー付き状態はキーごとにデータを分割し、関連レコードをまとめながら並列処理を可能にする。オペレーター状態はオペレーターインスタンス全体に適用される。
// 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 SQLとTable API
Flink 2.3は、増分ビュー保守のためのMaterialized Tablesにより SQL機能を拡張している。Table APIはバッチ処理とストリーム処理の統一インターフェースを提供する。
-- 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 vs Spark Structured Streaming
両フレームワークはストリーミングデータを処理するが、アーキテクチャが根本的に異なる。Flinkは真のストリーミングでレコードを個別に処理し、Sparkはマイクロバッチを処理する。Apache Sparkとの比較において、本番環境ではレイテンシと一貫性のトレードオフが重要となる。
| 観点 | Flink | Spark 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カスタムリソースによりデプロイを簡素化する。ジョブのライフサイクル、アップグレード、スケーリングを処理する。
# 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: savepointupgradeMode: savepoint設定により、オペレーターはアップグレード前にセーブポイントを取得し、デプロイ間で状態を保持する。
本番環境向けFlinkアプリケーションの最適化
本番デプロイでは、並列性、メモリ、状態バックエンドの設定に注意が必要である。統合の考慮事項についてはETLとデータパイプラインパターンを参照。
// 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-MailletSharpSkill 創業者
10 年以上フルスタック開発に携わっています。SharpSkill を運営し、ここで公開される内容に責任を負っています。
2026年8月28日 更新
タグ
共有
関連記事

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

Apache Spark 4.2 vs Databricks 2026年比較:アーキテクチャ、パフォーマンス、面接質問
2026年におけるApache Spark 4.2とDatabricksの徹底比較。Auto CDC、Metric Views、Photonエンジンなどの新機能、アーキテクチャの違い、パフォーマンス最適化、そして技術面接でよく問われる質問を詳しく解説します。

2026年版 Delta Lake vs Apache Iceberg:レイクハウスアーキテクチャと面接対策
Delta LakeとApache Icebergの技術的差異を詳細に解説します。パーティション進化、ACID トランザクション、クエリエンジン互換性などのレイクハウス面接頻出トピックを網羅的にカバーします。