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

MaxCompute:MaxFrame に基づく分散 Pandas 処理

最終更新日:Mar 01, 2026

MaxFrame は、MaxCompute 向けの分散データ処理フレームワークであり、Pandas 互換の API を提供します。データ分析コードを書き直す代わりに、使い慣れた Pandas の操作を使用するだけで、MaxFrame が MaxCompute クラスター全体に計算を自動的に分散させます。このアプローチにより、オープンソースの Pandas よりも数十倍高速なパフォーマンスを実現しながら、コードは標準的な Pandas のワークフローとほぼ同じ状態に保たれます。

このトピックでは、MaxFrame を使用した 3 つの一般的なデータ分析シナリオを通じて、mergegroupbyaggdrop_duplicatessort_values などの操作を解説します。

前提条件

MaxFrame がインストール済みであること。詳細については、「事前準備」をご参照ください。

共通セットアップ

このトピックの各シナリオでは、ODPS 接続と MaxFrame セッションが必要です。以下のコードブロックは、全体で使用される共通のセットアップを示しています。各シナリオのセクションには実行可能な完全なスクリプトが含まれていますが、接続パラメーターは同じです。

from odps import ODPS
from maxframe.session import new_session
import maxframe.dataframe as md
import os

o = ODPS(
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
    project='your-default-project',
    endpoint='your-endpoint',
)

session = new_session(o)
print(session.session_id)

# sales_maxframe_demo テーブルを読み取り、インデックス列として "index" を指定します。
sales = md.read_odps_table("sales_maxframe_demo", index_col="index")
# product_maxframe_demo テーブルを読み取り、インデックス列として "product_id" を指定します。
product = md.read_odps_table("product_maxframe_demo", index_col="product_id")

# sales テーブルと product テーブルを、sales の "product_id" 列と product のインデックスを基にマージします。
df = sales.merge(product, left_on="product_id", right_index=True)
# 結果から "product_name"、"year"、"price" のみを選択します。
df = df[["product_name", "year", "price"]]

# データフレームを実行し、結果をフェッチして表示します。
print(df.execute().fetch())

# 結果を MaxCompute テーブルに保存し、セッションを破棄します。
md.to_odps_table(df, "result_df", overwrite=True).execute()

session.destroy()

パラメーター

  • ALIBABA_CLOUD_ACCESS_KEY_ID:ターゲットの MaxCompute プロジェクトにアクセスするには、この環境変数に MaxCompute の権限 を持つ AccessKey ID を設定する必要があります。AccessKey ID は AccessKey 管理 ページから取得できます。

  • ALIBABA_CLOUD_ACCESS_KEY_SECRET:この環境変数に、AccessKey ID に対応する AccessKey Secret を設定します。

  • your-default-project:MaxCompute プロジェクトの名前。プロジェクト名を確認するには、MaxCompute コンソールにログインし、左側のナビゲーションウィンドウから [ワークスペース] > [プロジェクト] を選択します。

  • your-endpoint:MaxCompute プロジェクトのリージョンのエンドポイント。ネットワーク接続方法に基づいてエンドポイントを選択します。例:http://service.cn-chengdu.maxcompute.aliyun.com/api。詳細については、「エンドポイント」をご参照ください。

説明

MaxFrame は遅延実行を採用しています。計算は .execute() が呼び出されたときにのみ実行されます。すべての処理は MaxCompute クラスター内で行われるため、不要なデータ転送が回避されます。

データの準備

シナリオを実行する前に、MaxCompute プロジェクトに 2 つのテストテーブルを作成します:

テーブルスキーマレコード
product_maxframe_demoindex bigint, product_id bigint, product_name string, current_price bigint3 つのプロダクト (Nokia、Apple、Samsung)
sales_maxframe_demoindex bigint, sale_id bigint, product_id bigint, user_id bigint, year bigint, quantity bigint, price bigint6 件の販売レコード (2008--2015)

ステップ 1:テーブルの作成とデータの挿入

MaxFrame がインストールされた Python 環境で、以下のコードを実行します。このスクリプトは、ODPS SDK を介して MaxCompute に接続し、両方のテーブルにサンプルレコードを挿入します。

from odps import ODPS
from odps.df import DataFrame as ODPSDataFrame
from maxframe.session import new_session
import maxframe.dataframe as md
import pandas as pd
import os

o = ODPS(
    # 環境変数 ALIBABA_CLOUD_ACCESS_KEY_ID と ALIBABA_CLOUD_ACCESS_KEY_SECRET を、
    # ご自身の AccessKey ID と AccessKey Secret にそれぞれ設定します。
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
    project='your-default-project',
    endpoint='your-endpoint',
)

data_sets = [{
    "table_name": "product",
    "table_schema" : "index bigint, product_id bigint, product_name string, current_price bigint",
    "source_type": "records",
    "records" : [
        [1, 100, 'Nokia', 1000],
        [2, 200, 'Apple', 5000],
        [3, 300, 'Samsung', 9000]
    ],
},
{
    "table_name" : "sales",
    "table_schema" : "index bigint, sale_id bigint, product_id bigint, user_id bigint, year bigint, quantity bigint, price bigint",
    "source_type": "records",
    "records" : [
        [1, 1, 100, 101, 2008, 10, 5000],
        [2, 2, 300, 101, 2009, 7, 4000],
        [3, 4, 100, 102, 2011, 9, 4000],
        [4, 5, 200, 102, 2013, 6, 6000],
        [5, 8, 300, 102, 2015, 10, 9000],
        [6, 9, 100, 102, 2015, 6, 2000]
    ],
    "lifecycle": 5
}]

def prepare_data(o: ODPS, data_sets, suffix="", drop_if_exists=False):
    for index, data in enumerate(data_sets):
        table_name = data.get("table_name")
        table_schema = data.get("table_schema")
        source_type = data.get("source_type")

        if not table_name or not table_schema or not source_type:
            raise ValueError(f"インデックス {index} のデータセットに、必須キー ('table_name'、'table_schema'、'source_type') の 1 つ以上がありません。")

        lifecycle = data.get("lifecycle", 5)
        table_name += suffix

        print(f"{table_name} を処理中...")
        if drop_if_exists:
            print(f"{table_name} を削除中...")
            o.delete_table(table_name, if_exists=True)

        o.create_table(name=table_name, table_schema=table_schema, lifecycle=lifecycle, if_not_exists=True)

        if source_type == "local_file":
            file_path = data.get("file")
            if not file_path:
                raise ValueError(f"source_type が 'local_file' のインデックス {index} のデータセットに、'file' キーがありません。")
            sep = data.get("sep", ",")
            pd_df = pd.read_csv(file_path, sep=sep)
            ODPSDataFrame(pd_df).persist(table_name, drop_table=True)
        elif source_type == 'records':
            records = data.get("records")
            if not records:
                raise ValueError(f"source_type が 'records' のインデックス {index} のデータセットに、'records' キーがありません。")
            with o.get_table(table_name).open_writer() as writer:
                writer.write(records)
        else:
            raise ValueError(f"不明なデータセットの source_type: {source_type}")

        print(f"{table_name} の処理が完了しました")

prepare_data(o, data_sets, "_maxframe_demo", True)

ステップ 2:データの検証

以下の SQL クエリを実行して、両方のテーブルに期待されるレコードが含まれていることを確認します。

-- sales_maxframe_demo テーブルをクエリします
SELECT * FROM sales_maxframe_demo;

期待される出力:

+------------+------------+------------+------------+------------+------------+------------+
| index      | sale_id    | product_id | user_id    | year       | quantity   | price      |
+------------+------------+------------+------------+------------+------------+------------+
| 1          | 1          | 100        | 101        | 2008       | 10         | 5000       |
| 2          | 2          | 300        | 101        | 2009       | 7          | 4000       |
| 3          | 4          | 100        | 102        | 2011       | 9          | 4000       |
| 4          | 5          | 200        | 102        | 2013       | 6          | 6000       |
| 5          | 8          | 300        | 102        | 2015       | 10         | 9000       |
| 6          | 9          | 100        | 102        | 2015       | 6          | 2000       |
+------------+------------+------------+------------+------------+------------+------------+
-- product_maxframe_demo テーブルをクエリします
SELECT * FROM product_maxframe_demo;

期待される出力:

+------------+------------+--------------+---------------+
| index      | product_id | product_name | current_price |
+------------+------------+--------------+---------------+
| 1          | 100        | Nokia        | 1000          |
| 2          | 200        | Apple        | 5000          |
| 3          | 300        | Samsung      | 9000          |
+------------+------------+--------------+---------------+

MaxFrame を使用したデータ分析

以下の 3 つのシナリオでは、より複雑なデータ分析パイプラインを解説します。すべてのパフォーマンス比較では、同じベンチマークデータセットを使用します。具体的には、5,000 万レコード (1.96 GB) の販売テーブルと 10 万レコード (3 MB) のプロダクトテーブルです。

シナリオ 1:2 つのテーブルをマージしてプロダクト名を含む販売詳細を取得

このシナリオでは、merge() を使用して 2 つのテーブルを結合する方法を解説します。目的は、sales_maxframe_demo テーブル内のすべての sale_id 値と、各プロダクトに対応する product_nameyearprice を取得することです。

使用する Pandas 操作merge()、列選択 (df[columns])

サンプルコード

from odps import ODPS
from maxframe.session import new_session
import maxframe.dataframe as md
import os

o = ODPS(
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
    project='your-default-project',
    endpoint='your-endpoint',
)

session = new_session(o)
print(session.session_id)

sales = md.read_odps_table("sales_maxframe_demo", index_col="index")
product = md.read_odps_table("product_maxframe_demo", index_col="product_id")

df = sales.merge(product, left_on="product_id", right_index=True)
df = df[["product_name", "year", "price"]]

print(df.execute().fetch())

# 結果を MaxCompute テーブルに保存し、セッションを破棄します。
md.to_odps_table(df, "result_df", overwrite=True).execute()

session.destroy()

出力

index product_name  year  price
1            Nokia  2008   5000
2          Samsung  2009   4000
3            Nokia  2011   4000
4            Apple  2013   6000
5          Samsung  2015   9000
6            Nokia  2015   2000

パフォーマンス比較

環境時間 (秒)
ローカル Pandas (V1.3.5)65.8
MaxFrame22

シナリオ 2:各プロダクトの最初の販売年を検索

このシナリオでは、groupby()agg() を使用してプロダクトごとの最初の販売年を検索し、次にマルチインデックスで merge() を使用してその年の完全な販売レコードを取得する方法を解説します。結果には、プロダクト ID、最初の販売年、数量、価格が含まれます。

使用する Pandas 操作groupby()agg()、マルチインデックスでの merge()、列選択

サンプルコード

from odps import ODPS
from maxframe.session import new_session
import maxframe.dataframe as md
import os

o = ODPS(
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
    project='your-default-project',
    endpoint='your-endpoint',
)

session = new_session(o)
print(session.session_id)

# 集計して各プロダクトの最初の販売年を検索します。
min_year_df = md.read_odps_table("sales_maxframe_demo", index_col="index")
min_year_df = min_year_df.groupby('product_id', as_index=False).agg(first_year=('year', 'min'))

# マージして対応する販売レコードを取得します。
sales = md.read_odps_table("sales_maxframe_demo", index_col=['product_id', 'year'])
result_df = md.merge(sales, min_year_df,
                        left_index=True,
                        right_on=['product_id','first_year'],
                        how='inner')

result_df = result_df[['product_id', 'first_year', 'quantity', 'price']]

print(result_df.execute().fetch())

session.destroy()

出力

product_id  first_year  quantity  price
100         100        2008        10   5000
300         300        2009         7   4000
200         200        2013         6   6000

パフォーマンス比較

環境時間 (秒)
ローカル Pandas (V1.3.5)186
MaxFrame21

シナリオ 3:各ユーザーが最も多く消費したプロダクトを特定

このシナリオでは、複数の操作を連結した、より複雑な分析パイプラインを解説します。各ユーザーについて、合計で最も多く消費したプロダクトを見つけることが目的です。

使用する Pandas 操作groupby()agg()merge()drop_duplicates()sort_values()

サンプルコード

from odps import ODPS
from maxframe.session import new_session
import maxframe.dataframe as md
import os

o = ODPS(
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
    project='your-default-project',
    endpoint='your-endpoint',
)

session = new_session(o)
print(session.session_id)

sales = md.read_odps_table("sales_maxframe_demo", index_col="index")
product = md.read_odps_table("product_maxframe_demo", index_col="product_id")

sales['total'] = sales['price'] * sales['quantity']

product_cost_df = sales.groupby(['product_id', 'user_id'], as_index=False).agg(user_product_total=('total','sum'))
product_cost_df = product_cost_df.merge(product, left_on="product_id", right_index=True, how='right')

user_cost_df = product_cost_df.groupby('user_id').agg(max_total=('user_product_total', 'max'))
merge_df = product_cost_df.merge(user_cost_df, left_on='user_id', right_index=True)

result_df = merge_df[merge_df['user_product_total'] == merge_df['max_total']][['user_id', 'product_id']].drop_duplicates().sort_values(['user_id'], ascending = [1])

print(result_df.execute().fetch())

session.destroy()

出力

user_id  product_id
100      101         100
300      102         300

パフォーマンス比較

環境時間 (秒)
ローカル Pandas (V1.3.5)176
MaxFrame85

パフォーマンスのまとめ

以下の表は、3 つのシナリオすべてにおけるパフォーマンスの向上をまとめたものです。

シナリオローカル Pandas (V1.3.5)MaxFrame速度向上
1:2 つのテーブルのマージ65.8s22s約 3 倍
2:プロダクトごとの最初の販売年186s21s約 9 倍
3:ユーザーごとの最大消費プロダクト176s85s約 2 倍

まとめ

MaxFrame は Pandas API とシームレスに統合され、MaxCompute 上での自動分散処理を可能にします。計算を MaxCompute クラスターにオフロードすることで、MaxFrame は堅牢なデータ処理能力を維持しつつ、データ分析の規模と効率を大幅に向上させます。上記の 3 つのシナリオが示すように、MaxFrame はローカル Pandas よりも 2 倍から約 9 倍の速度向上を実現し、最小限の変更で標準的な Pandas コードを記述できます。これにより、MaxFrame は、データ量や計算時間のためにローカル Pandas での処理が非現実的になる大規模なデータ分析ワークロードに特に適しています。