DataWorks は、PyODPS 構文を使用して PyODPS タスクを開発するための PyODPS 2 ノードタイプを提供します。PyODPS は MaxCompute の Python SDK であり、PyODPS 2 ノードで Python コードを記述および編集して、MaxCompute と直接やり取りできます。
概要
PyODPS は MaxCompute の Python SDK です。簡潔なプログラミングインターフェイスを提供し、Python を使用してジョブの記述、テーブルやビューのクエリ、MaxCompute リソースの管理ができます。詳細については、「PyODPS 概要」をご参照ください。DataWorks では、PyODPS ノードを使用して Python タスクをスケジューリングおよび実行し、他のジョブと統合できます。
注意事項
-
DataWorks リソースグループで PyODPS ノードを実行する際にコードがサードパーティパッケージを必要とする場合は、サーバーレスリソースグループと カスタムイメージ を使用してパッケージをインストールしてください。
説明この方法は、コードでサードパーティパッケージを参照する UDF を使用している場合はサポートされません。このシナリオの対処方法については、「UDFの例:Python UDFでサードパーティパッケージを使用する」をご参照ください。
-
PyODPS タスクが VPC 内のデータソースやサービス、またはオンプレミスデータセンターなどの特定のネットワーク環境にアクセスする必要がある場合は、サーバーレスリソースグループを使用し、「ネットワーク接続ソリューション」を参照して必要なネットワーク接続を確立してください。
-
PyODPS の構文に関する詳細については、「PyODPSドキュメント」をご参照ください。
-
PyODPS ノードには、PyODPS 2 と PyODPS 3 の 2 つのタイプがあります。これらは基盤となる Python のバージョンが異なります。PyODPS 2 ノードは Python 2 を使用し、PyODPS 3 ノードは Python 3 を使用します。使用する Python のバージョンに対応するノードタイプを作成してください。
-
PyODPS ノードで SQL を実行してもデータリネージが生成されず、データマップにデータリネージが表示されない場合、タスクコードで DataWorks のスケジューリングパラメーターを手動で設定することでこの問題を解決できます。データリネージの表示については、「データリネージの表示」をご参照ください。パラメーター設定については、「ランタイムパラメーターのヒント設定」をご参照ください。以下のサンプルコードを使用して、必要なランタイムパラメーターを取得できます。
import os # ... # DataWorksスケジューリングのランタイムパラメーターを取得 skynet_hints = {} for k, v in os.environ.items(): if k.startswith('SKYNET_'): skynet_hints[k] = v # ... # タスク送信時にhintsを設定 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 ノードを実行する場合、処理するデータ量に基づいて PyODPS ノードの CU 割り当てを設定できます。
説明サーバーレスリソースグループでこのタスクを実行する場合、単一タスクの最大構成は
64CUですが、過剰な CU 割り当てによるリソース不足がタスクの起動に影響を与える可能性があるため、16CUを超えないことを推奨します。 -
[Got Killed] エラーが発生した場合、メモリ使用量が制限を超え、プロセスが終了したことを示します。したがって、ローカルでのデータ操作は可能な限り避けてください。PyODPS を介して開始される SQL および DataFrame タスク (to_pandas を除く) は、この制限の対象外です。
-
UDF 以外のコードでは、プリインストールされている NumPy および Pandas パッケージを使用できます。バイナリコードを含む他のサードパーティパッケージはサポートされていません。
-
互換性の理由から、DataWorks では、options.tunnel.use_instance_tunnel はデフォルトで False に設定されています。インスタンストンネルをグローバルに有効にする必要がある場合は、この値を手動で True に設定する必要があります。
-
PyODPS 2 ノードの基盤となる Python のバージョンは 2.7 です。
-
PyODPS ノード内で複数の Python タスクを同時に実行することはサポートされていません。
-
PyODPS ノードでログを出力するには、
printを使用します。logger.infoはサポートされていません。
前提条件
操作手順
-
PyODPS 2 ノードの編集ページで、以下の開発操作を実行してください。
PyODPS 2 コードの例
PyODPS ノードを作成した後、コードを編集して実行できます。PyODPS の構文に関する詳細については、「PyODPSドキュメント」をご参照ください。このトピックでは、以下の 5 つのコード例について説明します。ビジネスニーズに応じて例を選択できます。
ODPS エントリ
DataWorks の PyODPS ノードでは、グローバル変数
odpsまたはoが ODPS エントリとして含まれています。ODPS エントリを手動で定義する必要はありません。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 # limit 制限を無効にし、すべてのデータを読み込みます。 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:
ランタイムパラメーターの設定
-
dict型のhintsパラメーターを設定することで、実行時パラメーターを設定できます。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') # グローバル設定に基づいて hints を追加します。
実行結果の読み取り
SQL を実行するインスタンスは、
open_reader操作を直接実行できます。以下の 2 つのシナリオが考えられます。-
SQL ステートメントが構造化データを返します。
with o.execute_sql('select * from dual').open_reader() as reader: for record in reader: # 各レコードを処理します。 -
SQL ステートメントは、desc などのステートメントである場合があります。
reader.rawプロパティを使用すると、生の SQL 実行結果を取得できます。with o.execute_sql('desc dual').open_reader() as reader: print(reader.raw)説明カスタムスケジューリングパラメーターを使用する場合、PyODPS ノードは直接パラメーターを置換できないため、ページ上で PyODPS 2 ノードを直接トリガーする際には、時間値をハードコーディングする必要があります。
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()) # print 時に即時実行されます。 -
詳細情報の出力
options.verboseオプションを設定できます。DataWorks では、このオプションはデフォルトで有効になっており、実行中に Logview URL などの詳細情報が出力されます。
PyODPS 2 コード開発
次の簡単な例は、PyODPS ノードの使用方法を示します。
-
データセットを準備し、pyodps_iris サンプルテーブルを作成してください。詳細については、「DataFrameの作成」をご参照ください。
-
DataFrame を作成してください。詳細については、「DataFrameの作成」をご参照ください。
-
PyODPS ノードに次のコードを入力してください。
from odps.df import DataFrame # ODPS テーブルから DataFrame を作成します。 iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepallength.head(5))
PyODPS タスクの実行
-
デバッグの構成で、計算リソース、[コンピューティングリソースとクォータ]、およびDataWorks リソースグループを設定します。
説明-
パブリックネットワークまたは VPC ネットワーク環境のデータソースにアクセスするには、データソースの接続テストに合格したスケジューリングリソースグループを使用してください。詳細については、「ネットワーク接続ソリューション」をご参照ください。
-
タスク要件に応じて、イメージ 情報を設定できます。
-
-
ツールバーで実行をクリックすると、PyODPS タスクが実行されます。
-
-
ノードタスクを定期的に実行する必要がある場合は、ビジネスニーズに基づいてスケジュール設定を構成してください。詳細については、「スケジュール設定の構成」をご参照ください。
DataWorks の SQL ノードとは異なり、PyODPS ノードでは、コードへの影響を避けるため、コード内の ${param_name} のような文字列は置換されません。代わりに、コードを実行する前に、
argsという名前の dict がグローバル変数に追加され、そこからスケジューリングパラメーターを取得できます。たとえば、パラメータ でds=${yyyymmdd}を設定した場合、次のようにコードでパラメーターを取得できます。print('ds=' + args['ds'])説明dsという名前のパーティションを取得するには、以下の方法を使用します。o.get_table('table_name').get_partition('ds=' + args['ds']) -
ノードタスクの設定が完了したら、ノードをデプロイする必要があります。詳細については、「ノードのデプロイ」をご参照ください。
-
タスクをデプロイした後、オペレーションセンターでスケジュールされたタスクの実行ステータスを表示できます。詳細については、「スケジュールされたタスクの表示」をご参照ください。
関連付けられたロールを使用したノードの実行
RAM ロールを関連付けてノードを実行することで、特定の RAM ロールでノードタスクを実行し、詳細な権限制御とセキュリティ管理を実現できます。
次のステップ
PyODPS FAQ:PyODPS の実行中によくある問題について確認し、迅速なトラブルシューティングに役立てることができます。