このトピックでは、MaxCompute における一般的なデータスキューのシナリオとその解決策について説明します。
MapReduce
データスキューを理解するには、まず MapReduce を理解する必要があります。MapReduce は、分割統治戦略を使用する分散コンピューティングフレームワークです。大規模または複雑な問題を、管理しやすい小さなサブ問題に分割し、これらのサブ問題を処理した後、その結果をマージして最終的な出力を生成します。従来の並列プログラミングフレームワークと比較して、MapReduce は高いフォールトトレランス、使いやすさ、優れたスケーラビリティを提供します。MapReduce を使用して並列プログラムを実装する場合、データストレージやノード間の情報交換・伝達メカニズムなど、分散クラスターでのプログラミングに直接関係のない問題を考慮する必要はありません。これにより、分散プログラミングが大幅に簡素化されます。
以下の図は、MapReduce のワークフローを示しています。
データスキュー
データスキューは、多くの場合 reducer ステージで発生します。mapper は通常、入力ファイルを均等に分割しますが、データがワーカー間で不均等に分散されるとデータスキューが発生します。この不均等な分散により、一部のワーカーはすぐに終了するのに対し、他のワーカーははるかに長い時間がかかります。本番環境では、ほとんどのデータにスキューが発生しています。この現象は、80対20の法則としても知られるパレートの法則に従います。たとえば、フォーラムのアクティブユーザーの 20% が投稿の 80% を占めたり、ウェブサイトへのトラフィックの 80% を 20% のユーザーが生成したりすることがあります。ビッグデータの時代において、データスキューは分散プログラムのパフォーマンスに深刻な影響を与える可能性があります。よくある症状として、ジョブの進捗が 99% で止まっているように見えることがあります。
データスキューの特定方法
操作手順
MaxCompute でデータスキューを特定するには、次のように Logview を使用します。
[Fuxi Jobs] タブで、ジョブを [Latency] の降順でソートし、最もランタイムが長いジョブステージを選択します。
そのステージの [Fuxi instance] リストで、インスタンスを [Latency] の降順でソートします。平均より著しくランタイムが長いインスタンス (通常はリストの最初のインスタンス) を選択します。[StdOut] 列でその出力ログを表示します。
[StdOut] ログの情報を使用して、対応するジョブ実行グラフを表示します。
ジョブ実行グラフのキー情報を使用して、データスキューを引き起こしている SQL スニペットを特定します。
例
タスクの実行ログで Logview の URL を見つけます。詳細については、「Logview のエントリーポイント」をご参照ください。

問題を迅速に特定するために、Logview ページで Fuxi タスクを [Latency] の降順でソートし、最もランタイムが長いものを選択します。

タスク
R31_26_27のランタイムが最も長いです。 次の図に示すように、タスクR31_26_27をクリックしてインスタンス詳細ページに移動します。
行 Latency: {min:00:00:06, avg:00:00:13, max:00:26:40}は、インスタンスの最小ランタイムが6s、平均ランタイムが13s、最大ランタイムが26 minutes and 40 secondsであることを示します。インスタンスを
Latencyで降順にソートします。4 つのインスタンスのランタイムが長いことがわかります。MaxCompute は、ランタイムが平均値の 2 倍を超える場合、Fuxi インスタンスをロングテールと見なします。これは、ランタイムが
26sを超えるタスクインスタンスがロングテールとして識別されることを意味します。この場合、21 個のインスタンスのランタイムが26sを超えています。しかし、ロングテールインスタンスの存在は、必ずしもデータスキューを示すわけではありません。インスタンスのランタイムのavg値とmax値も比較する必要があります。max値がavg値よりもはるかに大きい場合、タスクに深刻なデータスキューがあると見なされ、最適化が必要になります。[StdOut] 列の
アイコンをクリックすると、次の例に示すように出力ログが表示されます。
問題を特定したら、[ジョブ詳細] タブで
R31_26_27を右クリックし、[すべて展開] を選択してタスクを展開します。 詳細については、「Logview 2.0 を使用してジョブ情報を表示する」をご参照ください。
StreamLineRead22の前のステップであるStreamLineWriter21を確認します。 これにより、スキューしたキー (new_uri_path_structure、cookie_x5check_userid、cookie_userid) と、データスキューの原因となっている SQL スニペットを特定できます。
データスキューのトラブルシューティングと解決
データスキューの最も一般的な原因を、頻度の高い順に以下に示します。
JOIN
GROUP BY
COUNT(DISTINCT)
ROW_NUMBER (TopN)
dynamic partition
JOIN
JOIN 操作で発生するデータスキューは、大きいテーブルと小さいテーブルの結合、大きいテーブルと中規模のテーブルの結合、またはロングテールを引き起こすホットキーなど、さまざまなシナリオによって引き起こされる可能性があります。
大きいテーブルと小さいテーブル
データスキューの例
次の例では、
t1は large テーブルで、t2とt3は small テーブルです。SELECT t1.ip ,t1.is_anon ,t1.user_id ,t1.user_agent ,t1.referer ,t2.ssl_ciphers ,t3.shop_province_name ,t3.shop_city_name FROM <viewtable> t1 LEFT OUTER JOIN <other_viewtable> t2 ON t1.header_eagleeye_traceid = t2.eagleeye_traceid LEFT OUTER JOIN ( SELECT shop_id ,city_name AS shop_city_name ,province_name AS shop_province_name FROM <tenanttable> WHERE ds = MAX_PT('<tenanttable>') AND is_valid = 1 ) t3 ON t1.shopid = t3.shop_id解決策
次のコードに示すように、MAPJOIN ヒント構文を使用します。
SELECT /*+ mapjoin(t2,t3)*/ t1.ip ,t1.is_anon ,t1.user_id ,t1.user_agent ,t1.referer ,t2.ssl_ciphers ,t3.shop_province_name ,t3.shop_city_name FROM <viewtable> t1 LEFT OUTER JOIN (<other_viewtable>) t2 ON t1.header_eagleeye_traceid = t2.eagleeye_traceid LEFT OUTER JOIN ( SELECT shop_id ,city_name AS shop_city_name ,province_name AS shop_province_name FROM <tenanttable> WHERE ds = MAX_PT('<tenanttable>') AND is_valid = 1 ) t3 ON t1.shopid = t3.shop_id注意事項
小さいテーブルまたはサブクエリを参照する場合は、そのエイリアスを使用する必要があります。
MAPJOIN は、小さいテーブルとしてサブクエリをサポートします。
MAPJOIN では、非等価結合を使用したり、
ORを使用して複数の条件を結合したりできます。ON句を省略し、mapjoin on 1 = 1を使用することで、デカルト積を計算できます。たとえば、select /*+ mapjoin(a) */ a.id from shop a join table_name b on 1=1;のようにします。ただし、この操作はデータ肥大化を引き起こす可能性があります。MAPJOIN では、複数の small テーブルをカンマ (
,) で区切ります。たとえば、/*+ mapjoin(a,b,c)*/のようにします。MAPJOIN は、map ステージで指定されたテーブルのすべてのデータをメモリにロードします。したがって、指定するテーブルは小さい必要があります。各テーブルのインメモリサイズは 512 MB を超えることはできません。この制限は、メモリにロードされた後のデータサイズに適用され、圧縮ストレージサイズよりも大幅に大きくなる可能性があります。次のパラメーターを設定することで、このメモリ制限を最大 8,192 MB まで増やすことができます。
SET odps.sql.mapjoin.memory.max=2048;MAPJOIN における JOIN 操作の制限事項:
LEFT OUTER JOINでは、左側テーブルは large テーブルである必要があります。RIGHT OUTER JOINの場合、右テーブルは large テーブルである必要があります。FULL OUTER JOINはサポートされていません。INNER JOINの場合、左テーブルまたは右テーブルのどちらかを large テーブルにすることができます。MAPJOIN は最大 128 の小さいテーブルをサポートします。この制限を超えると、構文エラーが報告されます。
大きいテーブルと中規模のテーブル
データスキューの例
以下の例では、
t0は large テーブルで、t1は中間のテーブルです。SELECT request_datetime ,host ,URI ,eagleeye_traceid FROM <viewtable> t0 LEFT JOIN ( SELECT traceid, eleme_uid, isLogin_is FROM <servicetable> WHERE ds = '${today}' AND hh = '${hour}' ) t1 ON t0.eagleeye_traceid = t1.traceid WHERE ds = '${today}' AND hh = '${hour}'解決策
次のコードに示すように、DISTRIBUTED MAPJOIN ヒントを使用してデータスキューを解決します。
SELECT /*+distmapjoin(t1)*/ request_datetime ,host ,URI ,eagleeye_traceid FROM <viewtable> t0 LEFT JOIN ( SELECT traceid, eleme_uid, isLogin_is FROM <servicetable> WHERE ds = '${today}' AND hh = '${hour}' ) t1 ON t0.eagleeye_traceid = t1.traceid WHERE ds = '${today}' AND hh = '${hour}'
ホットキー ジョイン
データスキューの例
次の表では、
eleme_uidカラムに多くのホットキーが含まれているため、データスキューが発生しやすくなっています。SELECT eleme_uid, ... FROM ( SELECT eleme_uid, ... FROM <viewtable> )t1 LEFT JOIN( SELECT eleme_uid, ... FROM <customertable> ) t2 ON t1.eleme_uid = t2.eleme_uid;解決策
この問題は、次の 3 つの方法のいずれかを使用して解決できます。
方法
名称
説明
方法 1
ホットキーの手動分割
ホットキーを特定し、メインテーブルからそれらをフィルター処理して MAPJOIN で処理します。残りのホットキー以外のレコードは MergeJoin で処理します。最後に、両方の JOIN の結果をマージします。
方法 2
SkewJoin ヒント
ヒント
/*+ skewJoin(<table_name>[(<column1_name>[,<column2_name>,...])][((<value11>,<value12>)[,(<value21>,<value22>)...])]*/を使用します。SkewJoin ヒントを使用すると、スキューキーを見つけるための追加のステップが加わり、クエリの実行時間が増加します。すでにスキューキーがわかっている場合は、SkewJoin パラメーターを設定して時間を節約できます。方法 3
モジュロ等価結合
乗数テーブルを使用してホットキーを分散させます。
ホットキーの手動分割。
ホットな値が特定された後、それらを含むレコードはメインテーブルからフィルター処理され、MapJoin の対象となります。ホットな値を含まない残りのレコードは MergeJoin で処理されます。最後に、2 つの結合の結果が結合されます。詳細については、次のコード例をご参照ください。
SELECT /*+ MAPJOIN (t2) */ eleme_uid, ... FROM ( SELECT eleme_uid, ... FROM <viewtable> WHERE eleme_uid = <skewed_value> )t1 LEFT JOIN( SELECT eleme_uid, ... FROM <customertable> WHERE eleme_uid = <skewed_value> ) t2 ON t1.eleme_uid = t2.eleme_uid UNION ALL SELECT eleme_uid, ... FROM ( SELECT eleme_uid, ... FROM <viewtable> WHERE eleme_uid != <skewed_value> )t3 LEFT JOIN( SELECT eleme_uid, ... FROM <customertable> WHERE eleme_uid != <skewed_value> ) t4 ON t3.eleme_uid = t4.eleme_uidSkewJoin ヒント。
SELECT文で、ヒント/*+ skewJoin(<table_name>[(<column1_name>[,<column2_name>,...])][((<value11>,<value12>)[,(<value21>,<value22>)...])]*/を使用してスキューを処理します。このヒントでは、table_nameはスキューのあるテーブルの名前、column_nameはスキューのある列の名前、valueはスキューキーの値です。次のコードに例を示します。-- 方法 1:テーブル名をヒントとして指定します。テーブルのエイリアスをヒントとして指定することに注意してください。 SELECT /*+ skewjoin(a) */ * FROM T0 a JOIN T1 b ON a.c0 = b.c0 AND a.c1 = b.c1; -- 方法 2:テーブル名と、スキューが疑われる列をヒントとして指定します。たとえば、テーブル 'a' の列 c0 と c1 にデータスキューがあります。 SELECT /*+ skewjoin(a(c0, c1)) */ * FROM T0 a JOIN T1 b ON a.c0 = b.c0 AND a.c1 = b.c1 AND a.c2 = b.c2; -- 方法 3:テーブル名と列をヒントとして指定し、スキューキーの値を指定します。キーの値が STRING 型の場合は、引用符で囲みます。たとえば、(a.c0=1 かつ a.c1="2") と (a.c0=3 かつ a.c1="4") の値は両方ともスキューしています。 SELECT /*+ skewjoin(a(c0, c1)((1, "2"), (3, "4"))) */ * FROM T0 a JOIN T1 b ON a.c0 = b.c0 AND a.c1 = b.c1 AND a.c2 = b.c2;説明値を直接指定する SkewJoin ヒントメソッドは、ホットキーを手動で分割したり、値を指定せずにヒントを使用したりするよりも効率的です。
SkewJoin ヒントでサポートされている JOIN タイプ:
INNER JOINの場合、結合内のどちらかのテーブルにヒントを指定できます。LEFT JOIN、SEMI JOIN、またはANTI JOINの場合、ヒントワードは左側テーブルにのみ指定できます。RIGHT JOINでは、右テーブルにのみヒントワードを指定できます。FULL JOINは SkewJoin ヒントをサポートしていません。
ヒントは集約を実行し、コストが発生するため、データスキューが確実にある JOIN にのみヒントを追加することを推奨します。
ヒントが指定された JOIN の左側の結合キーのデータ型は、右側の結合キーのデータ型と同じでなければなりません。そうでない場合、SkewJoin ヒントは有効になりません。たとえば、
a.c0のデータ型はb.c0のデータ型と同じでなければならず、a.c1のデータ型はb.c1のデータ型と同じでなければなりません。サブクエリで CAST 関数を使用して、データ型の一貫性を確保できます。以下に例を示します。CREATE TABLE T0(c0 int, c1 int, c2 int, c3 int); CREATE TABLE T1(c0 string, c1 int, c2 int); -- 方法 1: SELECT /*+ skewjoin(a) */ * FROM T0 a JOIN T1 b ON cast(a.c0 AS string) = b.c0 AND a.c1 = b.c1; -- 方法 2: SELECT /*+ skewjoin(b) */ * FROM (SELECT cast(a.c0 AS string) AS c00 FROM T0 a) b JOIN T1 c ON b.c00 = c.c0;SkewJoin ヒントを追加すると、オプティマイザーは集約を実行して上位 20 個のホットキーを取得します。
20はデフォルト値であり、set odps.optimizer.skew.join.topk.num = xx;を使用して変更できます。SkewJoin ヒントは、JOIN の片側のみのヒントをサポートします。
ヒントが指定された JOIN には
left_key = right_key条件が必要です。デカルト積 JOIN はサポートされていません。すでに MAPJOIN ヒントがある JOIN に SkewJoin ヒントを追加することはできません。
乗数テーブルを使用したモジュロ等価結合。
このアプローチは、前の 3 つの解決策とは論理的に異なります。分割統治戦略は使用しません。代わりに、スキューの度合いによって決定される N の値 (1 から N) を持つ単一の整数列を含む乗数テーブルを使用します。このテーブルは、ユーザー行動テーブルを N 倍に拡張するために使用されます。その後の JOIN 操作では、ユーザー ID と
numberの 2 つの結合キーを使用します。number結合条件を追加することにより、ユーザー ID のみに基づいてデータを分散させることによって引き起こされるデータスキューは、元のレベルの1/Nに減少します。ただし、このアプローチの欠点は、データも N 倍に膨張させることです。SELECT eleme_uid, ... FROM ( SELECT eleme_uid, ... FROM <viewtable> )t1 LEFT JOIN( SELECT /*+mapjoin(<multipletable>)*/ eleme_uid, number ... FROM <customertable> JOIN <multipletable> ) t2 ON t1.eleme_uid = t2.eleme_uid AND mod(t1.<value_col>,10)+1 = t2.number;データ肥大化に対処するために、拡張を両方のテーブルのホットキーレコードのみに制限し、他の非ホットキーレコードは変更しないままにすることができます。まず、ホットキーレコードを見つけます。次に、新しい
eleme_uid_join列を追加して、トラフィックテーブルとユーザー行動テーブルを個別に処理します。ユーザー ID がホットキーの場合、ランダムに割り当てられた正の整数 (たとえば 0 から 1,000) をCONCATします。それ以外の場合は、元のユーザー ID を保持します。2 つのテーブルを結合するときは、eleme_uid_join列を使用します。これにより、ホットキーが分散されてスキューが減少し、非ホットキーレコードの不要な拡張も回避されます。ただし、このロジックは元のビジネスロジック SQL を大幅に書き換えるため、推奨されません。
GROUP BY
次のコードは、GROUP BY 句を使用した擬似コードの例です。
SELECT shop_id
,sum(is_open) AS open_days
FROM table_xxx_di
WHERE dt BETWEEN '${bizdate_365}' AND '${bizdate}'
GROUP BY shop_id;データスキューが発生した場合、次の 3 つの解決策のいずれかを使用できます。
方法 | 名称 | 説明 |
方法 1 | GROUP BY のスキュー対策パラメーターを設定 |
|
方法 2 | 乱数の追加 | ロングテールを引き起こすキーを分割します。 |
方法 3 | ローリングテーブルの作成 | コストを削減し、効率を向上させます。 |
方法 1:GROUP BY のスキュー対策パラメーターを設定する。
SET odps.sql.groupby.skewindata=true;方法 2:乱数を追加する。
この解決策は、SQL を書き換えて乱数を追加し、ロングテールを引き起こすキーを分割します。これは、GROUP BY 操作におけるロングテールを解決するための効果的な方法です。
SQL クエリ
Select Key,Count(*) As Cnt From TableName Group By Key;の場合、コンバイナがないと、mapper ノードはデータを reducer ノードにシャッフルし、reducer ノードが COUNT 操作を実行します。対応する実行計画はM->Rです。ロングテールキーが特定されたと仮定すると、そのキーの作業を次のように再分散できます。
-- ロングテールキーが KEY001 であると仮定します。 SELECT a.Key ,SUM(a.Cnt) AS Cnt FROM(SELECT Key ,COUNT(*) AS Cnt FROM <TableName> GROUP BY Key ,CASE WHEN KEY = 'KEY001' THEN Hash(Random()) % 50 ELSE 0 END ) a GROUP BY a.Key;変更された実行計画は
M->R->Rになります。実行ステップ数は増えますが、ロングテールキーが 2 段階で処理されるため、全体のランタイムは短縮される可能性があります。リソース消費と時間効率は方法 1 と同様です。ただし、実際のシナリオでは、ロングテールキーが複数存在することがよくあります。ロングテールキーを見つけて SQL を書き換える手間を考えると、方法 1 の方がコスト効率が高い場合が多いです。方法 3:ローリングテーブルを作成する。
コストを削減し、効率を向上させるために、過去 1 年間のデータを取得する必要がある場合があります。オンラインタスクの場合、毎回
T-1からT-365までのすべてのパーティションを読み取ることは、リソースの大きな無駄です。ローリングテーブルを作成すると、過去 1 年間のデータ取得に影響を与えることなく、読み取るパーティションの数を減らすことができます。次のコードに例を示します。まず、365 日分のマーチャントビジネスデータを GROUP BY 集約で初期化し、データ更新日をマークして、テーブル
aとして保存します。その後のオンラインタスクでは、T-2日のテーブルaをtable_xxx_diテーブルと結合し、再度 GROUP BY を実行できます。これにより、毎日読み取るパーティションの数が 365 から 2 に減少します。プライマリキーshop_idの重複が大幅に減少し、リソース消費も削減されます。-- ローリングテーブルを作成します。 CREATE TABLE IF NOT EXISTS m_xxx_365_df ( shop_id STRING, last_update_ds STRING, `365d_open_days` BIGINT ) PARTITIONED BY ( ds STRING COMMENT 'Date partition' )LIFECYCLE 7; -- 365 日の期間が 2021-05-01 から 2022-05-01 であると仮定します。1 回限りの初期化を実行します。 INSERT OVERWRITE TABLE m_xxx_365_df PARTITION(ds = '20220501') SELECT shop_id, max(ds) as last_update_ds, sum(is_open) AS `365d_open_days` FROM table_xxx_di WHERE dt BETWEEN '20210501' AND '20220501' GROUP BY shop_id; -- 次に、実行される日次オンラインタスクは次のとおりです。 INSERT OVERWRITE TABLE m_xxx_365_df PARTITION(ds = '${bizdate}') SELECT aa.shop_id, aa.last_update_ds, `365d_open_days` - COALESCE(is_open, 0) AS `365d_open_days` -- 開店日数の無限ローリングを防ぎます。 FROM ( SELECT shop_id, max(last_update_ds) AS last_update_ds, sum(`365d_open_days`) AS `365d_open_days` FROM ( SELECT shop_id, ds AS last_update_ds, sum(is_open) AS `365d_open_days` FROM table_xxx_di WHERE ds = '${bizdate}' GROUP BY shop_id UNION ALL SELECT shop_id, last_update_ds, `365d_open_days` FROM m_xxx_365_df WHERE dt = '${bizdate_2}' AND last_update_ds >= '${bizdate_365}' -- ソースがすでにグループ化されている場合、ここでは GROUP BY は不要です。 ) GROUP BY shop_id ) AS aa LEFT JOIN ( SELECT shop_id, is_open FROM table_xxx_di WHERE ds = '${bizdate_366}' ) AS bb ON aa.shop_id = bb.shop_id;
COUNT(DISTINCT)
あるテーブルに次のようなデータ分布があるとします。
ds (パーティション) | cnt (レコード数) |
20220416 | 73,025,514 |
20220415 | 2,292,806 |
20220417 | 2,319,160 |
次のステートメントを使用すると、データスキューが容易に発生する可能性があります。
SELECT ds
,COUNT(DISTINCT shop_id) AS cnt
FROM demo_data0
GROUP BY ds;解決策は次のとおりです。
方法 | 名称 | 説明 |
方法 1 | パラメーターチューニング |
|
方法 2 | 汎用的な 2 段階集計 | パーティションフィールドの値に乱数を追加します。 |
方法 3 | 2 段階集計のような集計 | まず、 |
方法 1:パラメーターチューニング。
次のパラメーターを設定します。
SET odps.sql.groupby.skewindata=true;方法 2:汎用的な 2 段階集計。
shop_idフィールドのデータが不均等に分散している場合、方法 1 は効果的ではありません。より汎用的な方法は、パーティションフィールドの値に乱数を追加することです。-- 方法 A:乱数を連結します。CONCAT(ROUND(RAND(),1)*10,'_', ds) AS rand_ds SELECT SPLIT_PART(rand_ds, '_', 2) AS ds ,COUNT(DISTINCT shop_id) AS id_cnt FROM ( SELECT CONCAT(CAST(FLOOR(RAND() * 10) AS STRING), '_', ds) AS rand_ds ,shop_id FROM demo_data0 ) GROUP BY rand_ds; -- 方法 B:乱数フィールドを追加します。ROUND(RAND(),1)*10 AS randint10 SELECT ds ,COUNT(DISTINCT shop_id) AS id_cnt FROM (SELECT ds ,shop_id FROM demo_data0 ) GROUP BY ds, FLOOR(RAND() * 10);方法 3:2 段階集計のような集計。
GROUP BY フィールドと DISTINCT フィールドのデータが均等に分散している場合は、まず 2 つのグループ化フィールド (ds と shop_id) に GROUP BY を適用し、次に
count(distinct)コマンドを使用することでクエリを最適化できます。SELECT ds ,COUNT(shop_id) AS cnt FROM(SELECT ds ,shop_id FROM demo_data0 GROUP BY ds ,shop_id ) GROUP BY ds;
ROW_NUMBER (TopN)
次のコードは、Top-10 の例です。
SELECT main_id
,type
FROM (SELECT main_id
,type
,ROW_NUMBER() OVER(PARTITION BY main_id ORDER BY type DESC ) rn
FROM <data_demo2>
) A
WHERE A.rn <= 10;データスキューが発生した場合、次のいずれかの方法で解決できます。
方法 | 名称 | 説明 |
方法 1 | SQL ベースの 2 段階集計 | 乱数列を追加するか、乱数を追加して PARTITION BY 句のパラメーターとして使用します。 |
方法 2 | UDAF ベースの 2 段階集計 | UDAF を使用して、最小ヒープ優先度付きキューでクエリを最適化します。 |
方法 1:SQL ベースの 2 段階集計。
map ステージで各パーティショングループのデータをできるだけ均等に分散させるために、乱数列を追加し、それを PARTITION BY 句のパラメーターとして使用します。
-- 方法 1:乱数にモジュロを使用します。 SELECT main_id ,type FROM (SELECT main_id ,type ,ROW_NUMBER() OVER(PARTITION BY main_id ORDER BY type DESC ) rn FROM (SELECT main_id ,type FROM (SELECT main_id ,type ,ROW_NUMBER() OVER(PARTITION BY main_id,src_pt ORDER BY type DESC ) rn FROM (SELECT main_id ,type ,ceil(110 * rand()) % 11 AS src_pt FROM data_demo2 ) ) B WHERE B.rn <= 10 ) ) A WHERE A.rn <= 10; -- 方法 2:カスタムの乱数を使用します。 SELECT main_id ,type FROM (SELECT main_id ,type ,ROW_NUMBER() OVER(PARTITION BY main_id ORDER BY type DESC ) rn FROM (SELECT main_id ,type FROM(SELECT main_id ,type ,ROW_NUMBER() OVER(PARTITION BY main_id,src_pt ORDER BY type DESC ) rn FROM (SELECT main_id ,type ,ceil(10 * rand()) AS src_pt FROM data_demo2 ) ) B WHERE B.rn <= 10 ) ) A WHERE A.rn <= 10;方法 2:UDAF ベースの 2 段階集計。
SQL メソッドは、冗長で保守が困難なコードになる可能性があります。代わりに、UDAF と最小ヒープ優先度付きキューを使用して最適化できます。
iterateフェーズでは Top-N 要素のみが保持され、mergeフェーズでは N 要素のみがマージされます。プロセスは次のとおりです。iterate:最初の K 個の要素をプッシュします。K 個より後の要素については、最小ヒープの先頭要素と継続的に比較し、必要に応じて要素を交換します。merge:2 つのヒープをマージした後、上位 K 個の要素をその場で返します。terminate:ヒープを配列として返します。SQL クエリでは、配列を個別の行に分割します。
@annotate('* -> array<string>') class GetTopN(BaseUDAF): def new_buffer(self): return [[], None] def iterate(self, buffer, order_column_val, k): # heapq.heappush(buffer, order_column_val) # buffer = [heapq.nlargest(k, buffer), k] if not buffer[1]: buffer[1] = k if len(buffer[0]) < k: heapq.heappush(buffer[0], order_column_val) else: heapq.heappushpop(buffer[0], order_column_val) def merge(self, buffer, pbuffer): first_buffer, first_k = buffer second_buffer, second_k = pbuffer k = first_k or second_k merged_heap = first_buffer + second_buffer merged_heap.sort(reverse=True) merged_heap = merged_heap[0: k] if len(merged_heap) > k else merged_heap buffer[0] = merged_heap buffer[1] = k def terminate(self, buffer): return buffer[0] SET odps.sql.python.version=cp37; SELECT main_id,type_val FROM ( SELECT main_id ,get_topn(type, 10) AS type_array FROM data_demo2 GROUP BY main_id ) LATERAL VIEW EXPLODE(type_array)type_ar AS type_val;
動的パーティション
動的パーティションを使用すると、PARTITION 句でパーティション列名を指定し、特定の値を提供せずに、パーティションテーブルにデータを挿入できます。代わりに、パーティション値は SELECT 句の対応する列によって提供されます。したがって、作成される正確なパーティションは、SQL クエリの実行が終了し、パーティション列の値が決定されるまで不明です。詳細については、「動的パーティションへのデータの挿入または上書き (DYNAMIC PARTITION)」をご参照ください。次のコードは SQL の例です。
CREATE TABLE total_revenues (revenue bigint) partitioned BY (region string);
INSERT overwrite TABLE total_revenues PARTITION(region)
SELECT total_price AS revenue,region
FROM sale_detail;動的パーティションは多くのシナリオで使用され、データスキューを容易に引き起こす可能性があります。データスキューが発生した場合、次のいずれかの解決策で解決できます。
方法 | 名称 | 説明 |
方法 1 | パラメーター設定 | パラメーターを設定してクエリを最適化します。 |
方法 2 | プルーニング最適化 | レコード数が多いパーティションを見つけてプルーニングし、個別に挿入します。 |
方法 1:パラメーター設定。
動的パーティションは、異なる条件を満たすデータを異なるパーティションに配置できるため、複数の INSERT OVERWRITE ステートメントが不要になります。これにより、特にパーティションが多い場合にコードを大幅に簡素化できます。ただし、動的パーティションは、過剰な数の小さいファイルにつながる可能性もあります。
データスキューの例
次の単純な SQL を例にとります。
INSERT INTO TABLE part_test PARTITION(ds) SELECT * FROM part_test;K 個の Map インスタンスと N 個のターゲットパーティションがあると仮定します。
ds=1 cfile1 ds=2 ... X ds=3 cfilek ... ds=n最も極端なケースでは、
K*N個の小さいファイルが生成される可能性があります。過剰な数の小さいファイルは、ファイルシステムに大きな管理上の負荷をかける可能性があります。そのため、MaxCompute は、追加レベルの reducer タスクを導入することで動的パーティションを処理します。同じターゲットパーティションのデータを同じ (または少数の) reducer インスタンスによって書き込むように指示し、これにより、過剰な数の小さいファイルが作成されるのを回避します。この reducer は常にジョブの最後のタスクです。MaxCompute では、この機能はデフォルトで有効になっており、次のパラメーターが true に設定されていることを意味します。SET odps.sql.reshuffle.dynamicpt=true;この機能をデフォルトで有効にすると、小さいファイルが多すぎるという問題が解決され、単一のインスタンスによって生成されるファイルが多すぎるためにタスクが失敗するのを防ぎます。ただし、これにより、データスキューという新たな問題も発生します。さらに、追加の reducer ステージを導入すると、コンピューティングリソースが消費されます。したがって、トレードオフを慎重に検討する必要があります。
解決策
set odps.sql.reshuffle.dynamicpt=true;パラメーターを有効にして追加の reducer ステージを導入する当初の目的は、小さいファイルが多すぎるという問題を解決することでした。ただし、ターゲットパーティションの数が少なく、小さいファイルが多すぎるリスクがない場合、この機能をデフォルトで有効にすると、コンピューティングリソースを無駄にするだけでなく、パフォーマンスも低下します。この場合、set odps.sql.reshuffle.dynamicpt=false;を設定してこの機能を無効にすると、パフォーマンスが大幅に向上する可能性があります。次のコードに例を示します。INSERT overwrite TABLE ads_tb_cornucopia_pool_d PARTITION (ds, lv, tp) SELECT /*+ mapjoin(t2) */ '20150503' AS ds, t1.lv AS lv, t1.type AS tp FROM (SELECT ... FROM tbbi.ads_tb_cornucopia_user_d WHERE ds = '20150503' AND lv IN ('flat', '3rd') AND tp = 'T' AND pref_cat2_id > 0 ) t1 JOIN (SELECT ... FROM tbbi.ads_tb_cornucopia_auct_d WHERE ds = '20150503' AND tp = 'T' AND is_all = 'N' AND cat2_id > 0 ) t2 ON t1.pref_cat2_id = t2.cat2_id;上記のコードにデフォルトのパラメーターを使用した場合、ジョブの総ランタイムは約 1 時間 30 分です。最後の reducer ステージには約 1 時間 20 分かかり、これは総ランタイムの約
90%を占めます。追加の reducer ステージを導入すると、各 reducer インスタンスのデータ分布が非常に不均一になり、ロングテールが発生します。
上記の例では、生成された動的パーティションの履歴数を分析すると、毎日約 2 つの動的パーティションしか生成されていないことがわかります。したがって、安全に
set odps.sql.reshuffle.dynamicpt=false;を設定できます。これにより、ジョブはわずか 9 分で完了できます。この場合、このパラメーターをfalseに設定すると、パフォーマンスが大幅に向上し、コンピューティング時間とリソースを節約できます。この 1 つのパラメーター変更だけで、最小限の労力で大幅な改善が得られます。この最適化は、多くのリソースを消費する大規模で実行時間の長いジョブだけでなく、リソース消費の少ない通常の短時間実行ジョブにも適用されます。動的パーティションが使用され、動的パーティションの数が少ない限り、
odps.sql.reshuffle.dynamicptパラメーターをfalseに設定して、リソースを節約し、パフォーマンスを向上させることができます。ジョブの期間に関係なく、次の 3 つの条件をすべて満たすノードは最適化できます。
ジョブは動的パーティションを使用します。
動的パーティションの数は 50 以下です。
ジョブに `set odps.sql.reshuffle.dynamicpt=false;` がありません。
最後の Fuxi インスタンスの実行時間を使用して、ノードにこのパラメーターを設定する緊急度を判断できます。これは
diag_levelフィールドによって識別されます。ルールは次のとおりです。Last_Fuxi_Inst_Timeが 30 分を超える場合:Diag_Level=4 ('Critical')。Last_Fuxi_Inst_Timeが 20 分から 30 分の場合:Diag_Level=3 ('High')。Last_Fuxi_Inst_Timeが 10 分から 20 分の場合:Diag_Level=2 ('Medium')。Last_Fuxi_Inst_Timeが 10 分未満の場合:Diag_Level=1 ('Low')。
方法 2:プルーニング最適化。
動的パーティションにデータを挿入する際に map ステージですでに存在するデータスキューを解決するには、レコード数が多いパーティションを見つけてプルーニングし、それらを個別に挿入します。実際のユースケースに基づいて、map ステージのパラメーター設定を次のように変更できます。
SET odps.sql.mapper.split.size=128; INSERT OVERWRITE TABLE data_demo3 partition(ds,hh) SELECT * FROM dwd_alsc_ent_shop_info_hi;結果は、フルテーブルスキャンが実行されたことを示しています。さらに最適化するには、次のようにシステムによって導入された Reduce ジョブを無効にすることができます。
SET odps.sql.reshuffle.dynamicpt=false ; INSERT OVERWRITE TABLE data_demo3 partition(ds,hh) SELECT * FROM dwd_alsc_ent_shop_info_hi;動的パーティションにデータを挿入する際に map ステージでデータスキューを解決するには、レコード数が多いパーティションを見つけてプルーニングし、それらを個別に挿入します。具体的な手順は次のとおりです。
次のコマンドを使用して、レコード数が多い特定のパーティションをクエリします。
SELECT ds ,hh ,COUNT(*) AS cnt FROM dwd_alsc_ent_shop_info_hi GROUP BY ds ,hh ORDER BY cnt DESC;パーティションの一部を次に示します。
ds
hh
cnt
20200928
17
1052800
20191017
17
1041234
20210928
17
1034332
20190328
17
1000321
20210504
1
19
20191003
20
18
20200522
1
18
20220504
1
18
レコード数が多いパーティションをフィルターで除外し、残りのデータを挿入してから、レコード数が多いパーティションのデータを個別に挿入します。
SET odps.sql.reshuffle.dynamicpt=false ; -- レコード数が多くないパーティションのデータを挿入します。 INSERT OVERWRITE TABLE data_demo3 partition(ds,hh) SELECT * FROM dwd_alsc_ent_shop_info_hi WHERE CONCAT(ds,hh) NOT IN ('2020092817','2019101717','2021092817','2019032817'); -- レコード数が多いパーティションのデータを挿入します。 set odps.sql.reshuffle.dynamicpt=false ; INSERT OVERWRITE TABLE data_demo3 partition(ds,hh) SELECT * FROM dwd_alsc_ent_shop_info_hi WHERE CONCAT(ds,hh) IN ('2020092817','2019101717','2021092817','2019032817'); -- 結果を検証します。 SELECT ds ,hh,COUNT(*) AS cnt FROM dwd_alsc_ent_shop_info_hi GROUP BY ds,hh ORDER BY cnt desc;