リアルタイムタスク
AEデータ開発プラットフォームは、Apache Flinkの強力なストリーム計算能力を深く統合することで、より高性能なストリームデータ同期ソリューションを提供し、ストリームデータのリアルタイム分析と探索にも対応しています。これにより、将来を見据えたストリームとバッチの統合、高いリアルタイム性、低い運用コストを備えたエンタープライズ級のデータ開発プラットフォームを企業に提供します。
1. リアルタイムタスクの作成
1.1 新規作成の入口
データ開発プラットフォームの「開発」モジュールで、「リアルタイムタスク」タブ右上にある「+ リアルタイムタスク」ボタンをクリックすると、作成ページに入ります。
1.2 基本情報の入力
リアルタイムタスクはデュアル環境(開発環境 / 本番環境)に対応しており、基本情報は2つの環境で共通です。
| フィールド | 説明 | 備考 |
|---|---|---|
| 名前 | スペース内で一意で、48文字以内です。中国語、英語、数字、アンダースコアを使用できます | 必須 |
| Flinkバージョン | リアルタイムタスクに紐付けるFlinkエンジンのバージョンで、ドロップダウンリストから選択します | 作成後は変更不可 |
| 責任者 | リアルタイムタスクの責任者です。デフォルトでは作成者が責任者になります | 必須 |
| 備考 | タスクの説明で、200文字以内です | 任意 |
1.3 開発環境の編集ページ
作成が完了すると開発環境の編集ページに入ります。主に次のエリアで構成されます:
1.3.1 コード編集エリア
-
Flink SQLの構文ハイライトと自動補完に対応しています
-
パラメータ化された設定に対応しており、
${参数名}形式でパラメータを参照します -
エディタはコード内のパラメータを自動で解析し、下部の「使用されるパラメータ」エリアに表示します
-
エディタはテーブルの血縁関係も解析し、Flink SQLに含まれるテーブルのソース、出力先、connectorタイプなどの情報を表示します
1.3.2 タスクパラメータ管理(サイドバー)
リアルタイムタスクでは、次の2種類のパラメータタイプに対応しています:
| パラメータタイプ | 説明 |
|---|---|
| 空間パラメータ | プリセット空間パラメータとカスタム空間パラメータに対応しています |
| タスクパラメータ | リアルタイムタスクレベル で分離され、「純粋なテキスト」と「式」の2種類に対応しています |
パラメータの追加方法:
-
方法1:手動作成 — サイドバーで手動でパラメータを追加します。タスクの提出時に、パラメータ値が解析されて実行エンジンに渡されます
-
方法2:コード解析 — Flink SQLで
${abc}形式を使って直接パラメータを定義すると、システムが自動で解析してサイドバーに表示します
動的パラメータ:
環境ごとに異なるパラメータ値を使う必要がある場合は、動的パラメータを使用して、実行時に現在の環境に対応する値を指定できます。
${env.dev: 開発環境の値, env.product: 本番環境の値}
パラメータの検証ルール(デバッグ / 実行時にのみ検証し、保存時とリリース時には検証しません):
- タスクパラメータ:SQLで
${abc}を宣言しているのに値が定義されていない場合、デバッグ / 実行時にブロックされ、ユーザーに値の設定が求められます
- 空間パラメータ:SQLで参照している空間パラメータが削除されている場合、デバッグ / 実行時にブロックされ、パラメータが存在しないことが通知されます
2. リアルタイムタスクのデバッグ
2.1 デバッグ機能の説明
デバッグ機能は、開発環境でリアルタイムタスクの正しさを検証するためのもので、本番環境にリリースしなくてもタスクのロジックをテストできます。
「デバッグ」ボタンをクリックすると、次のいずれかを選択できます:
- 保存のみ:保存検証のフローのみをトリガーします
- 保存してデバッグする:保存検証 + デバッグ検証のフローをトリガーします
デバッグはリリース前の必須の手順ではありません。ユーザーは保存後、先にデバッグしなくても直接リリースできます。
2.2 検証フロー
保存検証:
-
SQL構文の検証
-
テーブル権限の検証(ソーステーブルと目標テーブルの読み書き権限)
-
パラメータの解析と記録
デバッグ / 実行検証(保存検証に加えて実施):
- パラメータの完全性の検証 — すべての
${XX}パラメータに値が必要です - 空間パラメータの存在検証
- 実行チャンネルの状態の検証
2.3 デバッグ実行
デバッグに成功すると、タスクは開発環境の実行チャンネルで実行されます。次の操作ができます:
- Flink UIに移動してタスクの運行状態を確認する
- リアルタイムログを確認する
- 実行指標を確認する
- デバッグタスクを手動で停止する
3. リアルタイムタスクの保存とリリース
3.1 保存
リアルタイムタスクのリリース(オンライン化)は一方向です — 「リリース」状態のリアルタイムタスクは「未リリース」状態に戻せません(オフラインのタスクフローが双方向に切り替えられるのとは異なります)。
タスクが「リリース待ち」のときの保存:
| 項目 | 説明 |
|---|---|
| 存在する環境 | 「開発環境」にのみ保存されます |
| オンライン状態 | 「リリース待ち」 |
| バージョン | 開発環境には最新のバージョンが表示されます。本番環境にはありません |
| 可能な操作 | 編集、リリースして配信開始する、削除 |
タスクが「リリース」状態のときの保存:
| 項目 | 説明 |
|---|---|
| 存在する環境 | 新しいバージョンの内容は「開発環境」にのみ保存されるため、本番環境に再リリースする必要があります |
| オンライン状態 | 「リリース」(変更あり) |
| バージョン | 開発環境には最新のものが表示されます。本番環境は前回リリースしたバージョンのままです |
| 可能な操作 | 編集、再リリース、実行、削除 |
3.2 保存して公開する
リリース操作により、開発環境のリアルタイムタスクのバージョンが本番環境に同期されます。
初回リリース:
-
開発環境と本番環境のバージョンが同時に更新されます
-
リリースに成功すると状態が「リリース」になります。「起動」をクリックして最初のインスタンスを実行する必要があります
-
リリースに失敗した場合は「リリース待ち」状態のままです
リリース済みタスクの再リリース:
-
開発環境と本番環境のバージョンがどちらも更新されます
-
状態は「リリース」のままです
-
リリースに成功したら、「起動」をクリックして新しいバージョンでインスタンスを実行します
3.3 リリース時の検証
リリース時には次の検証が行われます:
- 保存されていないタスクはリリースできません(リリースボタンはグレーアウトされます)
- 本番環境のテーブル権限の検証
- パラメータ値の完全性の検証
- テーブル構造の一貫性の検証
4. リアルタイムタスクの実行(起動)
4.1 本番環境の詳細ページ
タスクを本番環境にリリースした後は、本番環境の詳細ページでタスクを確認・管理できます。詳細ページには次のモジュールがあります:
| モジュール | 説明 |
|---|---|
| リアルタイムタスクの基本情報 | タスク名、責任者、エンジンバージョン、状態変更、作成 / 変更日時、説明 |
| 運営設定 | 実行チャンネルの選択(リソース設定)、運営設定のパラメータ |
| インスタンス内容 | 最新インスタンスのSQL内容、運営設定、テーブルの血縁関係、タスクパラメータの値 |
| 運営記録 | インスタンスのライフサイクル全体のイベントメッセージ |
| ログ | 実行ログと異常ログの情報を含みます |
| スナップショット | Checkpoint/Savepointの総数、成功 / 失敗の統計、履歴の詳細 |
4.2 実行パラメータの設定
リアルタイムタスクを起動する前に、次の設定を完了する必要があります:
4.2.1 パラメータの値設定
「コード解析パラメータ」については起動時に値を設定する必要があり、パラメータ値はインスタンスレベルで記録されます。
4.2.2 運営設定
実行リソースの設定
| 構成項目 | 説明 |
|---|---|
チャンネル(リソースプール)の選択 | Channel 名:システムはタスクのFlinkエンジンバージョンに応じて、使用可能なチャンネルを自動で絞り込みます
|
| 計算方法 |
|
| 並列度 | リソース仕様はチャンネルレベルで定義済みのため、インスタンスを実行に提出する際に並列度の数を選択します。
|
実行パラメータの設定
| カテゴリ | パラメータ | 説明 | デフォルト値 |
|---|---|---|---|
| Checkpoint 設定 | execution.checkpointing.interval | チェックポイントの間隔 | 180秒 |
execution.checkpointing.timeout | チェックポイントのタイムアウト時間 | 3分 | |
execution.checkpointing.min-pause | チェックポイント間の最短間隔 | 60秒 | |
状態データの期限切れ | table.exec.state.ttl | 状態データの有効期限です。0は期限切れにならないことを表します データが初めてシステムに入って処理されると、状態メモリに保存されます。次に同じ主キーのデータが到着すると、システムは以前に保存された状態データを使って計算を行い、そのアクセス時刻を更新します。このプロセスはデータの継続的な流れに依存しているため、リアルタイム計算の中核となります。設定したTTLの時間ウィンドウ内にデータが再びアクセスされなかった場合、そのデータはシステムによって期限切れとみなされ、状態ストレージから削除されます。TTLの値を適切に設定することで、計算の精度を維持できるだけでなく、古いデータを適時に削除して状態メモリの使用量を効果的に減らし、システムのメモリ負荷を下げて、計算効率とシステムの安定性を高めることができます。 | 0 |
| 再起動施策 | restart-strategy.type | 再起動施策のタイプ 再起動施策が設定されていない場合にのみ、Flinkはシステムのチェックポイントが有効かどうかに応じてジョブを再起動するかどうかを決定します(システムのチェックポイントが有効な場合はFixed Delayで設定された固定間隔でジョブを再起動し、無効な場合はジョブを再起動しません)。Flinkの再起動施策を設定している場合は、設定した再起動施策に従って再起動します。 このパラメータの値は次のとおりです:
| Fixed Delay |
4.2.3 起動施策
ユーザーがタスクを起動 / 再起動する際に、保存済みのスナップショットの状態から起動するかどうかを選択できます:
| 起動施策 | 説明 |
|---|---|
| 状態付き起動 | 既存の有効な状態(Checkpoint / Savepoint)から復元します。最新の状態から復元するか、履歴スナップショットの一覧から選択できます |
| 状態なし有効 | 初期状態を含まず、まったく新しく起動します |
注意: タスクは履歴バージョンでの起動には対応していません。任意のSavepointを選択できますが、Schemaの互換性はユーザー自身の責任となり、Schemaが変更されている場合は実行時にエラーが発生する可能性があります。
4.2.4 提出して起動
起動前に次の検証が実行され、検証に成功するとリアルタイムタスクが起動します:
-
タスクの存在検証
-
内容の妥当性検証
-
バージョンの一致検証
-
チャンネルの存在検証
-
スナップショットの存在検証(状態付き起動時)
5. リアルタイムタスクの運行状態の管理
5.1 タスクインスタンスの管理状態
| ステータス | 遷移可能な状態 | 実行可能な操作 | 説明 |
|---|---|---|---|
| 未リリース | — | 編集、リリースして配信開始する、削除 | リアルタイムタスクは開発環境にのみ存在します |
| 未有効 (Not Started) | — | 編集、リリースして配信開始する、削除 | 本番環境にリリース済みですが、まだ実行されていません |
| 有効中 (Starting) | 運行中、失敗、停止中、殺害中 | 編集、リリースして配信開始する、停止、Kill | タスクを実行エンジンに提出しています |
| 運行中 (Running) | 停止しました、停止中、殺害中、失敗 | 編集、リリースして配信開始する、停止、Kill | タスクが正常に実行されています |
| 失敗 (Failed) | 有効中(新しいインスタンス) | 編集、リリースして配信開始する | インスタンスの最終状態です。新しいインスタンスを再起動できます |
| 停止中 (Stopping) | 停止しました、殺害中、失敗 | 編集、リリースして配信開始する、Kill | 停止操作を実行しています |
| 停止しました (Stopped) | 有効中(新しいインスタンス) | 編集、リリースして配信開始する | インスタンスの最終状態です。新しいインスタンスを再起動できます |
| 殺害中 (Killing) | 停止しました、失敗 | 編集、リリースして配信開始する | タスクを強制終了しています |
注意: StoppingとKillingがタイムアウトすると、アラーム情報が送信されます。Stoppingの処理中は、2回目の停止操作が可能です。
5.2 StopとKillの違い
| 操作 | 説明 |
|---|---|
| 停止 (Stop) | グレースフルに停止します。Savepointをトリガーして状態を保存するかどうかを選択でき、停止後はSavepointから復元できます |
| Kill | 強制終了し、状態は保存しません。Stopが応答しない場合や、タスクがフリーズした場合に適しています |
- タスクの停止時: Savepointが自動でトリガーされます(ストレージパスをユーザーが定義する必要はありません)
- タスクの起動 / 再起動時: Checkpoint/Savepointから復元するかどうかを選択できます
6. 本番環境の詳細ページ
6.1 インスタンス内容
最新のインスタンス情報で、次の内容が含まれます
- スクリプトのSQL内容
- インスタンスの運営設定
- インスタンスのテーブルの血縁関係
- インスタンスに含まれるタスクパラメータの値
6.2 運営記録
運営記録では、Streampark Job Instance IDごとに各インスタンスのライフサイクル全体のイベントが記録されます:
- デフォルトでは最新インスタンスのイベントメッセージが表示され、インスタンスの作成時刻の降順に並びます
- ユーザーは展開してすべてのインスタンスのイベントメッセージを確認できます
- タイムラインによるグローバル情報の表示制御に対応しています
6.3 ログ
| カテゴリー | 内容 | 説明 | |
|---|---|---|---|
実行ログ | JobManagerログ | タスクのスケジューリング、Checkpoint、状態管理 | 単一のJMインスタンスのログ |
| TaskManagerログ | データ処理、ネットワーク通信 | 複数のTMインスタンスのログを集約 | |
| エラーメッセージ | ERROR/WARNログの集約 | 問題をすばやく特定するための入口 |
6.4 スナップショット管理
| フィールド | 説明 |
|---|---|
| スナップショット ID | Streamparkから返されるCheckpoint ID |
| ステータス | スナップショットの現在の状態 |
| スナップショットタイプ | Checkpoint:Flinkが自動で作成する「自動セーブ」で、障害復旧に使われます。軽量・高速で、自動化されています Savepoint:ユーザーが手動でトリガーする「手動セーブポイント」で、計画的な停止と再起動(コードの更新、クラスタのスケールアウト / スケールイン、バージョンアップなど)に使われます |
| 作成方法 | 手動(ユーザーが開始、Savepointのみ)/ 自動(Checkpointの設定に基づいて自動生成) |
| トリガー時間 | スナップショットがトリガーされた時刻 |
| 継続時間 | スナップショットの作成にかかった時間 |
| ソースインスタンス | スナップショットが属するインスタンスのID |
7. 管理ページのメタデータ
リアルタイムタスク管理の一覧ページでは、次のメタデータ情報を確認できます:
| フィールド | 説明 |
|---|---|
| リアルタイムタスク名 | タスクの一意の識別子 |
| 備考 | タスクの説明 |
| オンライン状態 | リリース / リリース待ち |
| 運行状態 | 現在のバージョンの最新インスタンスの運行状態 |
| エンジンバージョン | 例:Flink 1.16 |
| 責任者 | タスクの責任者 |
| オンライン時間 | 開発環境から本番環境にリリースした時刻 |
| 開始時間 | 最新インスタンスの起動時刻 |
| 最終修正時間 | 最後に変更された時刻 |
| 最終変更者 | 最後に変更した操作者 |
| 操作 | 編集、状態管理の操作、タスク詳細への移動 |

