メインコンテンツまでスキップ

リアルタイムタスク

最終更新 2026/10/03

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: 本番環境の値}

パラメータの検証ルール(デバッグ / 実行時にのみ検証し、保存時とリリース時には検証しません):

  1. タスクパラメータ:SQLで ${abc} を宣言しているのに値が定義されていない場合、デバッグ / 実行時にブロックされ、ユーザーに値の設定が求められます
  1. 空間パラメータ:SQLで参照している空間パラメータが削除されている場合、デバッグ / 実行時にブロックされ、パラメータが存在しないことが通知されます

2. リアルタイムタスクのデバッグ​

2.1 デバッグ機能の説明​

デバッグ機能は、開発環境でリアルタイムタスクの正しさを検証するためのもので、本番環境にリリースしなくてもタスクのロジックをテストできます。

「デバッグ」ボタンをクリックすると、次のいずれかを選択できます:

  • 保存のみ:保存検証のフローのみをトリガーします
  • 保存してデバッグする:保存検証 + デバッグ検証のフローをトリガーします

デバッグはリリース前の必須の手順ではありません。ユーザーは保存後、先にデバッグしなくても直接リリースできます。

2.2 検証フロー​

保存検証:

  1. SQL構文の検証

  2. テーブル権限の検証(ソーステーブルと目標テーブルの読み書き権限)

  3. パラメータの解析と記録

デバッグ / 実行検証(保存検証に加えて実施):

  1. パラメータの完全性の検証 — すべての ${XX} パラメータに値が必要です
  2. 空間パラメータの存在検証
  3. 実行チャンネルの状態の検証

2.3 デバッグ実行​

デバッグに成功すると、タスクは開発環境の実行チャンネルで実行されます。次の操作ができます:

  • Flink UIに移動してタスクの運行状態を確認する
  • リアルタイムログを確認する
  • 実行指標を確認する
  • デバッグタスクを手動で停止する

3. リアルタイムタスクの保存とリリース​

3.1 保存​

リアルタイムタスクのリリース(オンライン化)は一方向です — 「リリース」状態のリアルタイムタスクは「未リリース」状態に戻せません(オフラインのタスクフローが双方向に切り替えられるのとは異なります)。

タスクが「リリース待ち」のときの保存:

項目説明
存在する環境「開発環境」にのみ保存されます
オンライン状態「リリース待ち」
バージョン開発環境には最新のバージョンが表示されます。本番環境にはありません
可能な操作編集、リリースして配信開始する、削除

タスクが「リリース」状態のときの保存:

項目説明
存在する環境新しいバージョンの内容は「開発環境」にのみ保存されるため、本番環境に再リリースする必要があります
オンライン状態「リリース」(変更あり)
バージョン開発環境には最新のものが表示されます。本番環境は前回リリースしたバージョンのままです
可能な操作編集、再リリース、実行、削除

3.2 保存して公開する​

リリース操作により、開発環境のリアルタイムタスクのバージョンが本番環境に同期されます。

初回リリース:

  • 開発環境と本番環境のバージョンが同時に更新されます

  • リリースに成功すると状態が「リリース」になります。「起動」をクリックして最初のインスタンスを実行する必要があります

  • リリースに失敗した場合は「リリース待ち」状態のままです

リリース済みタスクの再リリース:

  • 開発環境と本番環境のバージョンがどちらも更新されます

  • 状態は「リリース」のままです

  • リリースに成功したら、「起動」をクリックして新しいバージョンでインスタンスを実行します

3.3 リリース時の検証​

リリース時には次の検証が行われます:

  1. 保存されていないタスクはリリースできません(リリースボタンはグレーアウトされます)
  2. 本番環境のテーブル権限の検証
  3. パラメータ値の完全性の検証
  4. テーブル構造の一貫性の検証

4. リアルタイムタスクの実行(起動)​

4.1 本番環境の詳細ページ​

タスクを本番環境にリリースした後は、本番環境の詳細ページでタスクを確認・管理できます。詳細ページには次のモジュールがあります:

モジュール説明
リアルタイムタスクの基本情報タスク名、責任者、エンジンバージョン、状態変更、作成 / 変更日時、説明
運営設定実行チャンネルの選択(リソース設定)、運営設定のパラメータ
インスタンス内容最新インスタンスのSQL内容、運営設定、テーブルの血縁関係、タスクパラメータの値
運営記録インスタンスのライフサイクル全体のイベントメッセージ
ログ実行ログと異常ログの情報を含みます
スナップショットCheckpoint/Savepointの総数、成功 / 失敗の統計、履歴の詳細

4.2 実行パラメータの設定​

リアルタイムタスクを起動する前に、次の設定を完了する必要があります:

4.2.1 パラメータの値設定​

「コード解析パラメータ」については起動時に値を設定する必要があり、パラメータ値はインスタンスレベルで記録されます。

4.2.2 運営設定​

実行リソースの設定

構成項目説明

チャンネル(リソースプール)の選択

Channel 名:システムはタスクのFlinkエンジンバージョンに応じて、使用可能なチャンネルを自動で絞り込みます

「有効」のチャンネルのみ選択できます

リアルタイムタスクの実行では、実行モード = Application のチャンネル(リソースプール)のみ選択できます

計算方法
  • リソースモードは計算型 / 汎用型の2種類です。
  • 計算型(1:2):1コアあたり2Gのメモリを割り当てます。CPU性能が高く、集約的な計算に適しています。汎用型(1:4):1コアあたり4Gのメモリを割り当てます。CPUとメモリのバランスが取れており、ほとんどの業務に適しています
  • 選択できるリソースモードは、選択したリソースプールに紐付けられた選択可能な仕様によって決まります
並列度

リソース仕様はチャンネルレベルで定義済みのため、インスタンスを実行に提出する際に並列度の数を選択します。

並列度は、Flinkにおけるタスクの並列分割を表すものです。Slot数は、単一のTaskManagerに対するリソース分割の粒度です。

実行パラメータの設定

カテゴリパラメータ説明デフォルト値
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の再起動施策を設定している場合は、設定した再起動施策に従って再起動します。

このパラメータの値は次のとおりです:

  • Failure Rate(失敗率リスタート):検出の時間間隔、最大失敗回数、各リスタート間のインターバルを設定します
  • Fixed Delay(固定時間インターバル再起動):再起動試行回数、各リスタート間のインターバルを設定します
  • No Restarts(再起動なし):ジョブの失敗後、自動で再起動しません
Fixed Delay

4.2.3 起動施策​

ユーザーがタスクを起動 / 再起動する際に、保存済みのスナップショットの状態から起動するかどうかを選択できます:

起動施策説明
状態付き起動既存の有効な状態(Checkpoint / Savepoint)から復元します。最新の状態から復元するか、履歴スナップショットの一覧から選択できます
状態なし有効初期状態を含まず、まったく新しく起動します

注意: タスクは履歴バージョンでの起動には対応していません。任意のSavepointを選択できますが、Schemaの互換性はユーザー自身の責任となり、Schemaが変更されている場合は実行時にエラーが発生する可能性があります。

4.2.4 提出して起動​

起動前に次の検証が実行され、検証に成功するとリアルタイムタスクが起動します:

  1. タスクの存在検証

  2. 内容の妥当性検証

  3. バージョンの一致検証

  4. チャンネルの存在検証

  5. スナップショットの存在検証(状態付き起動時)

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 スナップショット管理​

フィールド説明
スナップショット IDStreamparkから返されるCheckpoint ID
ステータススナップショットの現在の状態
スナップショットタイプ

Checkpoint:Flinkが自動で作成する「自動セーブ」で、障害復旧に使われます。軽量・高速で、自動化されています

Savepoint:ユーザーが手動でトリガーする「手動セーブポイント」で、計画的な停止と再起動(コードの更新、クラスタのスケールアウト / スケールイン、バージョンアップなど)に使われます

作成方法手動(ユーザーが開始、Savepointのみ)/ 自動(Checkpointの設定に基づいて自動生成)
トリガー時間スナップショットがトリガーされた時刻
継続時間スナップショットの作成にかかった時間
ソースインスタンススナップショットが属するインスタンスのID

7. 管理ページのメタデータ​

リアルタイムタスク管理の一覧ページでは、次のメタデータ情報を確認できます:

フィールド説明
リアルタイムタスク名タスクの一意の識別子
備考タスクの説明
オンライン状態リリース / リリース待ち
運行状態現在のバージョンの最新インスタンスの運行状態
エンジンバージョン例:Flink 1.16
責任者タスクの責任者
オンライン時間開発環境から本番環境にリリースした時刻
開始時間最新インスタンスの起動時刻
最終修正時間最後に変更された時刻
最終変更者最後に変更した操作者
操作編集、状態管理の操作、タスク詳細への移動
このページは役に立ちましたか?