PyODPS 3 ノードを使用すると、Python で MaxCompute タスクを作成し、DataWorks で定期的に実行するようにスケジューリングできます。
概要
PyODPS は MaxCompute の Python SDK です。これにより、Python でタスクの作成、テーブルやビューのクエリ、MaxCompute リソースの管理ができます。詳細については、「PyODPSの概要」をご参照ください。DataWorks では、PyODPS ノードを使用することで、Python タスクを他のタスクと並行してスケジューリングおよび実行できます。
注意事項
-
DataWorks リソースグループ上の PyODPS ノードからサードパーティパッケージを呼び出すには、カスタムイメージを使用したサーバーレスリソースグループを使用します。
説明この方法は、コードにサードパーティパッケージを参照するユーザー定義関数 (UDF) が含まれる場合は適用されません。正しい手順については、「UDFの例:Python UDFでサードパーティパッケージを使用する」をご参照ください。
-
PyODPS のバージョンをアップグレードするには、サーバーレスリソースグループでカスタムイメージを使用して、
/home/tops/bin/pip3 install pyodps==0.12.1コマンドを実行します。0.12.1は対象の PyODPS バージョンに置き換えることができます。専用スケジューリングリソースグループでは、O&M Assistant を使用して同じコマンドを実行します。 -
PyODPS タスクが VPC や IDC ネットワーク内のデータソースやサービスなど、特定のネットワーク環境にアクセスする必要がある場合は、サーバーレスリソースグループを使用します。サーバーレスリソースグループをターゲット環境に接続する方法の詳細については、「ネットワーク接続ソリューション」をご参照ください。
-
PyODPS の構文の詳細については、「PyODPSドキュメント」をご参照ください。
-
PyODPS ノードには、PyODPS 2 (Python 2) と PyODPS 3 (Python 3) の 2 種類があります。ご使用の Python バージョンに合ったノードタイプを作成してください。
-
PyODPS ノードで SQL 文を実行しても、データマップで正しいデータリネージが生成されない場合は、タスクコード内で関連する DataWorks のスケジューリングパラメーターを手動で設定してください。データリネージの表示方法については、「データリネージの表示」をご参照ください。パラメーターの設定方法については、「ランタイムパラメーターヒントの設定」をご参照ください。次のサンプルコードは、実行時に必要なパラメーターを取得します。
import os ... # DataWorks スケジューラのランタイムパラメーターを取得します skynet_hints = {} for k, v in os.environ.items(): if k.startswith('SKYNET_'): skynet_hints[k] = v ... # タスクの送信時にヒントを設定します o.execute_sql('INSERT OVERWRITE TABLE XXXX SELECT * FROM YYYY WHERE ***', hints=skynet_hints) ... -
PyODPS ノードの出力ログの最大サイズは 4 MB です。大量のデータをログに出力することは避けてください。アラートと進捗情報の出力に重点を置いてください。
制限事項
-
専用スケジューリングリソースグループで PyODPS ノードを実行する場合、ローカルで処理されるデータは 50 MB を超えないようにしてください。この制限は、専用スケジューリングリソースグループの仕様によるものです。ローカルデータがオペレーティングシステムのしきい値を超えると、OOM (Got Killed) エラーが発生する可能性があります。PyODPS ノードで過剰なデータ処理コードを記述することは避けてください。
-
サーバーレスリソースグループで PyODPS ノードを実行する場合は、データ量に基づいて適切な数の CU を設定してください。
説明サーバーレスリソースグループでは、1 つのタスクで最大
64 CUをサポートします。ただし、タスクの起動に影響を与える可能性のあるリソース不足を防ぐため、16 CU以下を推奨します。 -
[Got killed] エラーは、プロセスがメモリ制限を超えたことを示します。ローカルでのデータ操作は避けてください。PyODPS を介して開始された SQL および DataFrame タスク (to_pandas 操作を除く) は、この制限の対象外です。
-
ユーザー定義関数を含まないコードでは、プリインストールされている NumPy および pandas ライブラリを使用できます。バイナリコードを含む他のサードパーティパッケージはサポートされていません。
-
互換性の理由から、DataWorks では
options.tunnel.use_instance_tunnelはデフォルトでFalseに設定されています。インスタンストンネルをグローバルに有効にするには、この値を手動でTrueに設定してください。 -
Python 3 のマイナーバージョン (Python 3.8 と Python 3.7 など) 間では、バイトコードの定義が異なります。
MaxCompute は現在 Python 3.7 を使用しています。Python 3.8 の
finallyブロックのような他の Python 3 バージョンの構文を使用すると、実行中にエラーが発生します。Python 3.7 の使用を推奨します。 -
PyODPS 3 は、サーバーレスリソースグループでの実行をサポートしています。購入して使用するには、「サーバーレスリソースグループの使用」をご参照ください。
-
1 つの PyODPS ノード内で複数の Python タスクを同時に実行するように設定することはできません。
-
PyODPS ノードでログを出力するには、
printを使用します。logger.infoは使用できません。
前提条件
操作手順
-
PyODPS 3 ノードのエディターページでコードを開発します。
PyODPS 3 のコード例
PyODPS ノードを作成した後、コードを編集して実行できます。PyODPS の構文の詳細については、「基本的な操作」をご参照ください。次の例では、5 つの一般的なシナリオについて説明します。ニーズに合ったものを選択してください。
ODPSエントリーポイント
DataWorks の各 PyODPS ノードには、グローバルな ODPS エントリーポイント変数
odpsまたはoが含まれており、手動で定義する必要はありません。print(odps.exist_table('PyODPS_iris'))SQLの実行
PyODPS ノードで SQL 文を実行できます。詳細については、「SQL」をご参照ください。
-
デフォルトでは、DataWorks の
instance tunnelは無効になっているため、instance.open_readerは Result インターフェースを使用し、最大 10,000 件のレコードを返します。reader.countを使用してレコード数を取得できます。すべてのデータを読み取るには、limitを無効にする必要があります。次のステートメントは、instance tunnelをグローバルに有効にし、limitを無効にします。options.tunnel.use_instance_tunnel = True options.tunnel.limit_instance_tunnel = False # 制限を無効にして、すべてのデータを読み取ります。 with instance.open_reader() as reader: # すべてのデータはインスタンストンネルを通じて読み取ることができます。 -
open_reader呼び出しにtunnel=Trueを追加することで、単一のopen_reader呼び出しに対してinstance tunnelを有効にすることもできます。また、limit=Falseを追加して、その呼び出しに対するlimit制限を無効にすることもできます。# このopen_reader呼び出しにインスタンストンネルを使用し、すべてのデータを読み取ります。 with instance.open_reader(tunnel=True, limit=False) as reader:
ランタイムパラメーターの設定
-
hintsパラメーターはdictであり、これを使用してランタイムパラメーターを設定できます。hints の詳細については、「SET 操作」をご参照ください。o.execute_sql('select * from PyODPS_iris', hints={'odps.sql.mapper.split.size': 16}) -
sql.settingsを使用してグローバル設定を行うと、関連するランタイムパラメーターがすべての実行に追加されます。from odps import options options.sql.settings = {'odps.sql.mapper.split.size': 16} o.execute_sql('select * from PyODPS_iris') # グローバル設定に基づいてヒントが追加されます。
実行結果の読み取り
SQL 実行インスタンスで
open_readerを直接呼び出すことができます。これには次の 2 つのシナリオがあります:-
SQL が構造化データを返す場合。
with o.execute_sql('select * from dual').open_reader() as reader: for record in reader: # 各レコードを処理します。 -
descなどの SQL 文を実行する場合、reader.rawプロパティを使用して SQL 実行の生の結果を取得できます。with o.execute_sql('desc dual').open_reader() as reader: print(reader.raw)説明カスタムスケジューリングパラメーターを使用し、UI から直接 PyODPS 3 ノードを実行する場合、ノードは実行時に変数を置換できないため、時間の値をハードコーディングする必要があります。
DataFrame
DataFrame (非推奨) を使用してデータを処理することもできます。
-
実行
DataWorks 環境では、DataFrame 操作は即時実行メソッドを明示的に呼び出す必要があります。
from odps.df import DataFrame iris = DataFrame(o.get_table('pyodps_iris')) for record in iris[iris.sepal_width < 3].execute(): # 即時実行メソッドを呼び出して各レコードを処理します。print() で出力する際に即時実行をトリガーする必要がある場合は、
options.interactiveを有効にする必要があります。from odps import options from odps.df import DataFrame options.interactive = True # 最初にオプションを有効にします。 iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepal_width.sum()) # これにより即時実行がトリガーされます。 -
詳細情報の出力
DataWorks では
options.verboseがデフォルトで有効になっており、実行中に Logview の URL などの詳細情報が出力されます。
PyODPS 3 のコード開発
次の例では、PyODPS ノードの使用方法を説明します:
-
pyodps_iris サンプルテーブルを作成してデータセットを準備します。詳細については、「DataFrameデータの処理」をご参照ください。
-
DataFrame を作成します。詳細については、「MaxComputeテーブルからDataFrameを作成する」をご参照ください。
-
PyODPS ノードに次のコードを入力して実行します。
from odps.df import DataFrame # ODPSテーブルからDataFrameを作成します。 iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepallength.head(5))
PyODPS タスクの実行
-
デバッグの構成 ペインの 計算リソース セクションで、計算リソース、[コンピューティングクォータ]、DataWorks リソースグループ を設定します。
説明-
パブリックネットワークまたは VPC 経由でデータソースにアクセスするには、データソースとの接続テストに合格したスケジューリングリソースグループを使用する必要があります。詳細については、「ネットワーク接続ソリューション」をご参照ください。
-
タスクの要件に基づいて イメージ 情報を設定できます。
-
-
ツールバーのパラメーターダイアログボックスで、作成した MaxCompute データソースを選択し、実行 をクリックして PyODPS タスクを実行します。
-
-
ノードを定期的に実行するには、ビジネス要件に基づいてそのスケジューリングプロパティを設定してください。詳細については、「ノードのスケジューリングの設定」をご参照ください。
DataWorks の SQL ノードとは異なり、PyODPS ノードはコード内の ${param_name} などの文字列を置き換えません。代わりに、コードが実行される前に、
argsという名前の辞書がグローバル変数に追加され、そこからスケジューリングパラメーターを取得できます。たとえば、パラメータ セクションでds=${yyyymmdd}を設定した場合、コード内で次のようにパラメーター情報を取得できます。print('ds=' + args['ds']) ds=20240930説明dsという名前のパーティションを取得する必要がある場合は、次の方法を使用できます。o.get_table('table_name').get_partition('ds=' + args['ds']) -
ノードを設定した後、デプロイする必要があります。詳細については、「ノードとワークフローのデプロイ」をご参照ください。
-
タスクがデプロイされた後、オペレーションセンターでそのステータスを表示できます。詳細については、「オペレーションセンター入門」をご参照ください。
関連付けられたロールを使用したノードの実行
RAM ロールを関連付けてノードを実行することで、特定の RAM ロールでノードタスクを実行し、詳細な権限制御とセキュリティ管理を実現できます。
次のステップ
PyODPS に関する FAQ:PyODPS の実行中に発生する一般的な問題とその解決方法。