MaxFrame は、MaxCompute 向けの分散データ処理フレームワークであり、Pandas 互換の API を提供します。データ分析コードを書き直す代わりに、使い慣れた Pandas の操作を使用するだけで、MaxFrame が MaxCompute クラスター全体に計算を自動的に分散させます。このアプローチにより、オープンソースの Pandas よりも数十倍高速なパフォーマンスを実現しながら、コードは標準的な Pandas のワークフローとほぼ同じ状態に保たれます。
このトピックでは、MaxFrame を使用した 3 つの一般的なデータ分析シナリオを通じて、merge、groupby、agg、drop_duplicates、sort_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_demo | index bigint, product_id bigint, product_name string, current_price bigint | 3 つのプロダクト (Nokia、Apple、Samsung) |
sales_maxframe_demo | index bigint, sale_id bigint, product_id bigint, user_id bigint, year bigint, quantity bigint, price bigint | 6 件の販売レコード (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_name、year、price を取得することです。
使用する 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 |
| MaxFrame | 22 |
シナリオ 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 |
| MaxFrame | 21 |
シナリオ 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 |
| MaxFrame | 85 |
パフォーマンスのまとめ
以下の表は、3 つのシナリオすべてにおけるパフォーマンスの向上をまとめたものです。
| シナリオ | ローカル Pandas (V1.3.5) | MaxFrame | 速度向上 |
|---|---|---|---|
| 1:2 つのテーブルのマージ | 65.8s | 22s | 約 3 倍 |
| 2:プロダクトごとの最初の販売年 | 186s | 21s | 約 9 倍 |
| 3:ユーザーごとの最大消費プロダクト | 176s | 85s | 約 2 倍 |
まとめ
MaxFrame は Pandas API とシームレスに統合され、MaxCompute 上での自動分散処理を可能にします。計算を MaxCompute クラスターにオフロードすることで、MaxFrame は堅牢なデータ処理能力を維持しつつ、データ分析の規模と効率を大幅に向上させます。上記の 3 つのシナリオが示すように、MaxFrame はローカル Pandas よりも 2 倍から約 9 倍の速度向上を実現し、最小限の変更で標準的な Pandas コードを記述できます。これにより、MaxFrame は、データ量や計算時間のためにローカル Pandas での処理が非現実的になる大規模なデータ分析ワークロードに特に適しています。