このトピックでは、Realtime Compute for Apache Flink でのストリーミングおよびバッチ PyFlink ジョブのデプロイと開始方法について説明し、開発ワークフローを解説します。
前提条件
-
RAM ユーザーまたは RAM ロールを使用してコンソールにアクセスする場合、その ID に必要な権限が付与されていることを確認してください。詳細については、「権限管理」をご参照ください。
-
ワークスペースが作成されていること。詳細については、「Realtime Compute for Apache Flink の有効化」をご参照ください。
ステップ 1: Python コードファイルの準備
Realtime Compute for Apache Flink の管理コンソールは、Python 開発環境を提供していません。ジョブはローカルで開発してください。ジョブのデバッグとコネクタの詳細については、「PyFlink ジョブの開発」をご参照ください。
ローカル開発で使用する Flink のバージョンが、「ステップ 3: PyFlink ジョブのデプロイ」で選択するエンジンバージョンと一致していることを確認してください。カスタム Python 仮想環境、サードパーティの Python パッケージ、JAR パッケージ、データファイルなど、他の依存関係の使用方法については、「Python 依存関係の使用」をご参照ください。
すぐに始められるように、このトピックではワードカウントジョブのサンプル Python ファイルとサンプルデータファイルを提供しています。これらをダウンロードして、次のステップで使用できます。
-
適切なサンプル Python ジョブファイルをダウンロードします。
-
ストリーミングジョブ: word_count_streaming.py
-
バッチジョブ: word_count_batch.py
-
-
Shakespeare をクリックして、サンプルデータファイルをダウンロードします。
ステップ 2: Python ファイルとデータファイルのアップロード
-
Realtime Compute コンソールにログインします。
-
対象の Flink ワークスペースを見つけ、[操作] 列の [コンソール] をクリックします。
-
左側のナビゲーションウィンドウで、[アーティファクト] をクリックします。
-
[アーティファクトのアップロード] をクリックして、Python ファイルとデータファイルをアップロードします。
ステップ 1 でダウンロードしたサンプルの Python ファイルとデータファイルをアップロードします。ファイルストレージパスの詳細については、「アーティファクト」をご参照ください。
ステップ 3: PyFlink ジョブのデプロイ
ストリーミング
-
ページで、 を選択します。
-
デプロイメントパラメーターを設定します。
パラメーター
説明
例
デプロイメントモード
ストリームモードを選択します。
ストリームモード
デプロイメント名
Python デプロイメントの名前を入力します。
flink-streaming-test-python
エンジンバージョン
デプロイメント用の Flink エンジンバージョン。
信頼性とパフォーマンスを向上させるために、[推奨] または [安定] タグが付いたバージョンを使用することを推奨します。詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。
vvr-8.0.9-flink-1.17
Python URI
word_count_streaming.py サンプルファイルをダウンロードします。次に、アップロード
アイコンをクリックしてファイルを選択し、アップロードします。ファイルがすでに [アーティファクト] に存在する場合、再アップロードせずに直接選択できます。
-
エントリモジュール
プログラムのエントリポイントモジュール。
-
PyFlink ジョブが .py ファイルの場合、このパラメーターは必須ではありません。
-
PyFlink ジョブが .zip ファイルの場合、エントリモジュールを入力する必要があります。例:
word_count
必須ではありません
エントリーポイントのメイン引数
main メソッドに渡す引数。
このチュートリアルでは、入力データファイル Shakespeare のストレージパスを入力します。
--input oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/ShakespeareShakespeare ファイルの完全なパスを [アーティファクト] ページからコピーできます。
デプロイ先
ドロップダウンリストから、ターゲットの [キュー] または [セッションクラスター] を選択します。本番環境ではセッションクラスターの使用は推奨されません。詳細については、「キューの管理」および「セッションクラスターの作成」をご参照ください。
重要セッションクラスター上のデプロイメントは、監視メトリクス、アラート設定、または Autopilot をサポートしていません。セッションクラスターは開発およびテスト目的でのみ使用し、本番環境では使用しないでください。詳細については、「デプロイメントのデバッグ」をご参照ください。
default-queue
他の構成パラメーターの詳細については、「ジョブのデプロイ」をご参照ください。
-
-
[デプロイ] をクリックします。
バッチ
-
ページで、[デプロイメントの作成]をクリックし、[Python デプロイメント]を選択します。
-
デプロイメントパラメーターを設定します。
パラメーター
説明
例
デプロイメントモード
バッチモードを選択します。
バッチモード
デプロイメント名
Python デプロイメントの名前を入力します。
flink-batch-test-python
エンジンバージョン
デプロイメント用の Flink エンジンバージョン。
信頼性とパフォーマンスを向上させるために、[推奨] または [安定] タグが付いたバージョンを使用することを推奨します。詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。
vvr-8.0.9-flink-1.17
[Python URI]
word_count_batch.py サンプルファイルをダウンロードします。次に、アップロード
アイコンをクリックしてファイルを選択し、アップロードします。-
エントリモジュール
プログラムのエントリポイントモジュール。
-
PyFlink ジョブが .py ファイルの場合、このパラメーターは必須ではありません。
-
PyFlink ジョブが .zip ファイルの場合、エントリモジュールを入力する必要があります。例:
word_count
必須ではありません
エントリポイントのメイン引数
main メソッドに渡す引数。
このチュートリアルでは、入力ファイル Shakespeare と出力ディレクトリ
python-batch-quickstart-test-outputのストレージパスを入力します。説明出力ディレクトリのパスを指定するだけで済みます。出力ディレクトリは入力ファイルと同じ親ディレクトリにある必要があります。出力ディレクトリを事前に作成する必要はありません。
--input oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/Shakespeare--output oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/python-batch-quickstart-test-outputShakespeare ファイルの完全なパスを [アーティファクト] ページからコピーできます。
デプロイ先
ドロップダウンリストから、ターゲットの [キュー] または [セッションクラスター] を選択します。本番環境ではセッションクラスターの使用は推奨されません。詳細については、「キューの管理」および「セッションクラスターの作成」をご参照ください。
重要セッションクラスター上のデプロイメントは、監視メトリクス、アラート設定、または Autopilot をサポートしていません。セッションクラスターは開発およびテスト目的でのみ使用し、本番環境では使用しないでください。詳細については、「デプロイメントのデバッグ」をご参照ください。
default-queue
他の構成パラメーターの詳細については、「ジョブのデプロイ」をご参照ください。
-
-
[デプロイ] をクリックします。
ステップ 4: デプロイメントの開始と結果の表示
ストリーミング
-
ページで、対象のデプロイメントを見つけ、[操作] 列の [開始] をクリックします。
-
[ジョブの開始] ダイアログボックスで、[初期モード] を選択し、[開始] をクリックします。詳細については、「デプロイメントの開始」をご参照ください。
[開始] をクリックした後、ステータスが [実行中] または [完了] になれば、デプロイメントは正常に実行されています。このトピックのサンプルファイルを使用した場合、最終的なステータスは [完了] になります。
-
デプロイメントのステータスが [実行中] に変わったら、ストリーミングデプロイメントの結果を表示します。
重要このトピックのサンプル Python ファイルを使用した場合、ストリーミングデプロイメントが [完了] 状態になると結果が削除されます。そのため、結果はデプロイメントが [実行中] 状態のときにのみ表示できます。
末尾が .out の TaskManager ログファイルで、
shakespeareを検索して計算結果を見つけます。[ログ] タブで、[実行中の TaskManager] タブをクリックします。関連する TaskManager について、[ログリスト] サブタブをクリックします。
flink.outファイルを開き、右上の検索ボックスにshakespeareと入力して、(shakespeare,1)のようなワードカウントの結果を見つけます。
バッチ
-
ページで、目的のデプロイメントを見つけ、操作列の [開始] をクリックします。
リストをフィルターするには、タイプのドロップダウンリストから [バッチデプロイメント] を選択します。
-
[ジョブの開始] ダイアログボックスで、[開始] をクリックします。詳細については、「デプロイメントの開始」をご参照ください。
-
デプロイメントのステータスが [完了] に変わったら、バッチデプロイメントの結果を表示します。
OSS コンソールにログインします。oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/python-batch-quickstart-test-output ディレクトリに移動します。デプロイメントの開始日時で名付けられたフォルダをクリックし、対象のファイルをクリックしてから、表示されるパネルで [ダウンロード] をクリックします。
バッチデプロイメントは .ext ファイルを生成します。ファイルをダウンロードした後、テキストエディターまたは Microsoft Word で開いて結果を表示します。出力は次のようになります:
(As,40) (At,5) (Ay,1) (Be,9) (By,14) (Do,4) (He,7) (I,,4) (If,34) (In,36) (Is,10) (It,6)
(任意) ステップ 5: デプロイメントの停止
ジョブへの変更 (コードの変更、WITH パラメーターの更新、バージョンの変更など) を適用するには、ジョブを再デプロイし、停止してから再起動する必要があります。ステートレス起動や非動的構成の変更を適用する場合も再起動が必要です。ジョブの停止の詳細については、「ジョブの停止」をご参照ください。
関連トピック
-
デプロイメントのリソースは、開始前に設定することも、デプロイメントの実行後に変更することもできます。基本 (粗粒度) とエキスパート (詳細) の 2 つのリソース設定モードがサポートされています。詳細については、「デプロイメントリソースの設定」をご参照ください。
-
Realtime Compute for Apache Flink は、デプロイメントパラメーターの動的更新をサポートしています。これにより、構成がより迅速に有効になり、デプロイメントの停止と開始によるサービス停止時間が短縮されます。詳細については、「動的スケーリングとパラメーター更新」をご参照ください。
-
デプロイメントのログレベルを設定し、異なるログレベルに対して異なる出力を指定します。詳細については、「ログ出力の設定」をご参照ください。
-
SQL 開発ワークフローのウォークスルーについては、「Flink SQL ジョブ」をご参照ください。