すべてのプロダクト
Search
ドキュメントセンター

DataWorks:PyODPS 3 ノード

最終更新日:Aug 26, 2026

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 は使用できません。

前提条件

MaxCompute コンピュートエンジンと DataWorks ワークスペースの関連付け。

操作手順

  1. 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 ノードの使用方法を説明します:

    1. pyodps_iris サンプルテーブルを作成してデータセットを準備します。詳細については、「DataFrameデータの処理」をご参照ください。

    2. DataFrame を作成します。詳細については、「MaxComputeテーブルからDataFrameを作成する」をご参照ください。

    3. PyODPS ノードに次のコードを入力して実行します。

      from odps.df import DataFrame
      
      # ODPSテーブルからDataFrameを作成します。
      iris = DataFrame(o.get_table('pyodps_iris'))
      print(iris.sepallength.head(5))

    PyODPS タスクの実行

    1. デバッグの構成 ペインの 計算リソース セクションで、計算リソース、[コンピューティングクォータ]、DataWorks リソースグループ を設定します。

      説明
      • パブリックネットワークまたは VPC 経由でデータソースにアクセスするには、データソースとの接続テストに合格したスケジューリングリソースグループを使用する必要があります。詳細については、「ネットワーク接続ソリューション」をご参照ください。

      • タスクの要件に基づいて イメージ 情報を設定できます。

    2. ツールバーのパラメーターダイアログボックスで、作成した MaxCompute データソースを選択し、実行 をクリックして PyODPS タスクを実行します。

  2. ノードを定期的に実行するには、ビジネス要件に基づいてそのスケジューリングプロパティを設定してください。詳細については、「ノードのスケジューリングの設定」をご参照ください。

    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'])
  3. ノードを設定した後、デプロイする必要があります。詳細については、「ノードとワークフローのデプロイ」をご参照ください。

  4. タスクがデプロイされた後、オペレーションセンターでそのステータスを表示できます。詳細については、「オペレーションセンター入門」をご参照ください。

関連付けられたロールを使用したノードの実行

RAM ロールを関連付けてノードを実行することで、特定の RAM ロールでノードタスクを実行し、詳細な権限制御とセキュリティ管理を実現できます。

次のステップ

PyODPS に関する FAQ:PyODPS の実行中に発生する一般的な問題とその解決方法。