リソースのチューニング、集計最適化の有効化、SQL パターンの書き換えにより、Flink SQL デプロイメントのスループットを向上させ、レイテンシーを削減します。
スループットのベースラインパラメーターの設定
設定 タブの パラメーター セクションにある その他の設定 フィールドに、以下のパラメーターを追加します。これにより、スループットが向上し、ホットスポットの問題が軽減されます。カスタムデプロイメントパラメーターの設定をご参照ください。
execution.checkpointing.interval: 180s
table.exec.state.ttl: 129600000
table.exec.mini-batch.enabled: true
table.exec.mini-batch.allow-latency: 5s
table.optimizer.distinct-agg.split.enabled: true
| パラメーター | 説明 |
|---|---|
execution.checkpointing.interval |
チェックポイントの間隔。値 180s は 180 秒を意味します。 |
state.backend |
状態バックエンドのタイプ。 |
table.exec.state.ttl |
状態データの Time-to-Live (TTL)。単位はミリ秒です。 |
table.exec.mini-batch.enabled |
ミニバッチ集計を有効にします。 |
table.exec.mini-batch.allow-latency |
ミニバッチがトリガーされるまでの最大レイテンシー。 |
table.exec.mini-batch.size |
ミニバッチあたりの最大レコード数。VVR はこれを自動的に最適化するため、ほとんどの場合、設定は不要です。詳細については、「主要なパラメーター」をご参照ください。 |
table.optimizer.distinct-agg.split.enabled |
COUNT DISTINCT に対する PartialFinal 最適化を有効にします。 |
デプロイメントリソースのスケーリング
Ververica Platform (VVP) は、JobManager と TaskManager の CPU を設定値で上限付けます。デプロイメントがリソースに制約される場合は、これらの値を増やしてください。
JobManager リソースのスケーリング
並列デプロイメントの場合、設定 タブの リソース セクションで JobManager のリソースを増やします。例:
-
Job Manager CPU:4
-
Job Manager Memory:8 GiB
TaskManager リソースのスケーリング
複雑なトポロジーの場合、設定 タブの リソース セクションで TaskManager のリソースを増やします。例:
-
Task Manager CPU:2
-
Task Manager Memory:4 GiB
taskmanager.numberOfTaskSlots はデフォルト値の 1 のままにしてください。
グループ集計の最適化
デフォルトでは、グループ集計演算子はレコードを 1 つずつ処理します。つまり、状態からアキュムレータを読み取り、更新し、書き戻す、という処理を繰り返します。このレコードごとのパターンは、特に RocksDB を使用している場合に状態バックエンドのオーバーヘッドを増加させ、データホットスポットの下で悪化します。
miniBatch の有効化
MiniBatch は、受信レコードをバッファリングして一括で処理することで、バッチごとの状態アクセスを削減します。これにより、わずかに高いレイテンシーと引き換えにスループットが向上します。miniBatch 機能は、指定された間隔でソースに挿入されるイベントメッセージに基づいてマイクロバッチ処理をトリガーします。
利用シーン:超低レイテンシーが要求されないデータ集計シナリオ。
有効化の方法:設定 タブの パラメーター セクションにある その他の設定 フィールドに以下を追加します。詳細については、「カスタムデプロイメントパラメーターの設定」をご参照ください。
table.exec.mini-batch.enabled: true
table.exec.mini-batch.allow-latency: 5s
| パラメーター | 説明 |
|---|---|
table.exec.mini-batch.enabled |
true に設定して miniBatch を有効にします。 |
table.exec.mini-batch.allow-latency |
ミニバッチがトリガーされるまでの最大レイテンシー。 |
table.exec.mini-batch.size |
ミニバッチあたりの最大レコード数。VVR が自動的に最適化するため、設定は不要です。詳細については、「主要なパラメーター」をご参照ください。 |
LocalGlobal の有効化
LocalGlobal は、MapReduce の combine と reduce に似て、集計を 2 つのフェーズに分割します。つまり、各上流ノードでのローカル事前集計と、その後のグローバル集計です。ローカル集計はホットスポットデータを事前にフィルタリングし、グローバルフェーズに到達するデータ量を削減します。
利用シーン:SUM、COUNT、MAX、MIN、AVG などの一般的な集計関数で、データホットスポットの問題がある場合。
前提条件:
-
MiniBatch が有効になっている必要があります。LocalGlobal は、ミニバッチの間隔を使用して、ローカル集計の前に蓄積するレコード数を決定します。
-
集計関数は
mergeメソッド (AggregateFunction) を実装している必要があります。
ステータス:miniBatch がアクティブな場合、デフォルトで有効になります。
検証:最終的なトポロジーで GlobalGroupAggregate または LocalGroupAggregate ノードを確認します。
COUNT DISTINCT に対する PartialFinal の有効化
ローカル集計では重複する distinct キーを削除できないため、LocalGlobal は COUNT DISTINCT に対しては効果が低くなります。PartialFinal は、集計を 2 つのレイヤーに分割することでこの問題に対処します。最初のレイヤーは distinct キーのハッシュによってデータを分散させ、2 番目のレイヤーが最終的な集計を実行します。
利用シーン:データ量が大きく、集計パフォーマンスがボトルネックとなっている COUNT DISTINCT クエリ。
有効化の方法:その他の設定 フィールドに以下を追加します。詳細については、「カスタムデプロイメントパラメーターの設定」をご参照ください。
table.optimizer.distinct-agg.split.enabled: true
制限事項:
-
PartialFinal は、ユーザー定義集約関数 (UDAF) を含む SQL では使用できません。
-
データ量が大きい場合にのみ有効にしてください。PartialFinal は追加のネットワークシャッフルを導入します。
検証:最終的なトポロジーで、1 レイヤー集計が 2 レイヤー集計に変わったかどうかを確認します。
CASE WHEN の代わりに AGG WITH FILTER を使用
異なる条件下で同じフィールドに対して COUNT DISTINCT を計算する場合、CASE WHEN の代わりに FILTER 句を使用します。オプティマイザーは、同じ distinct キーに対する異なるフィルター引数を認識し、単一の共有ステートインスタンスを使用するため、状態アクセスとサイズが削減されます。パフォーマンステストでは、CASE WHEN に比べて 2 倍の改善が見られました。
変更前 (非効率):
COUNT(DISTINCT visitor_id) AS uv_total,
COUNT(DISTINCT CASE WHEN is_wireless = 'y' THEN visitor_id ELSE NULL END) AS uv_wireless
変更後 (2 倍のパフォーマンス):
COUNT(DISTINCT visitor_id) AS uv_total,
COUNT(DISTINCT visitor_id) FILTER (WHERE is_wireless = 'y') AS uv_wireless
デルタ集計 (ステートレス集計) の使用
デルタ集計は、VVR 11.8 以降を使用する Realtime Compute for Apache Flink デプロイメントでのみサポートされています。
デルタ集計は、ストリーミングの非ウィンドウグループ集計に対してステートレスな集計を提供します。Flink の状態でグループごとに集計アキュムレータを継続的に維持する通常のグループ集計とは異なり、デルタ集計はデータ変更に応答します。グループ化キーを使用してグループの現在の詳細レコードを非同期でルックアップし、それらのレコードから集計結果を再計算します。
その結果、Flink 側の集計状態の量は、グループの履歴数に依存しなくなります。デプロイメントは、グループ数が増えるにつれて集計状態を継続的に拡張する必要がありません。これにより、チェックポイントに永続化されるデータ量が削減され、回復時間が短縮され、大規模なグループ集計デプロイメントのリソース使用量が安定します。
利用シーン
デルタ集計は、グループの現在の詳細レコードを非同期にルックアップし、オンデマンドで集計結果を再計算します。以下の要件の 1 つ以上を持つデプロイメントで使用します。
-
デプロイメントロジックが変更された後、通常のグループ集計では、ストリーミングデプロイメントを状態なしで開始し、その状態を再構築する必要があります。デルタ集計を使用すると、まずステートレスなバッチデプロイメントを実行して結果テーブルを更新し、その後、状態付きでストリーミングデプロイメントを開始できます。これにより、バッチ処理の高いスループットを利用して、結果のバックフィルとデプロイメントの回復を高速化します。
-
グループの状態が大きい場合のチェックポイントのレイテンシーと状態処理のバックプレッシャーを削減します。
-
デプロイメントが再起動する際に復元する必要がある集計状態を大幅に削減し、回復時間を短縮します。
-
集計状態によって使用されるコンピューティングリソースとストレージリソースを削減し、コンピューティングユニット (CU) コストを低減します。
-
複数のデプロイメントが同じ詳細データを集計する場合に、ソーステーブルの現在の詳細状態を共有します。
-
ソーステーブルから現在の詳細レコードを直接クエリして、集計結果を検証し、問題をトラブルシューティングします。
制限事項
クエリ構造
-
ストリーミングデプロイメントにおける単一レベルの非ウィンドウグループ集計のみがサポートされています。ウィンドウ集計とカスケード集計はサポートされていません。
-
GROUP BYは、ソーステーブルのフィールドに直接マッピングされる列のセットを含み、ソーステーブルの完全なプライマリキーまたはクエリ可能なインデックスを含む必要があります。 -
PROCTIME()やRAND()などの非決定論的な式はサポートされていません。
上流および下流テーブル
-
ソースコネクタは、ストリーミングスキャンと非同期ルックアップの両方をサポートしている必要があります。
-
ソーステーブルは Watermark を定義できません。
-
デルタ集計が下流に送信する変更レコードには、
UPDATE_AFTER(+U) とDELETE(-D) のみが含まれます。UPDATE_BEFORE(-U) は含まれません。
集計関数
-
サポートされている関数は次のとおりです:
SUM、COUNT、AVG、MIN、MAX、LISTAGG、FIRST_VALUE、およびLAST_VALUE。 -
DISTINCT集計、FILTER句を持つ集計、ユーザー定義集約関数 (UDAF)、Python 集計関数、およびその他リストにない集計関数はサポートされていません。
例
e コマースプラットフォームは、多数の注文に対して現在のレコード数と金額を維持する必要があります。結果をサマリーテーブルに書き込むことで、注文クエリサービスがマーチャントと注文によって直接それを取得できるようにします。次の例では、ソーステーブルの完全なプライマリキーをルックアップキーとして使用します。
-
ステップ 1: ソーステーブルと結果テーブルの作成
ソーステーブルは、ストリーミングスキャンと非同期ルックアップの両方をサポートしている必要があります。
CREATE TABLE orders ( merchant_id BIGINT, -- マーチャント ID order_id BIGINT, -- 注文 ID amount DECIMAL(18, 2), -- 現在の注文金額 PRIMARY KEY (merchant_id, order_id) NOT ENFORCED ) WITH ( 'connector' = '<source-connector>', '<source-option>' = '<source-option-value>' );すべてのグループ化フィールドを結果テーブルのプライマリキーとして使用します。
CREATE TABLE order_current_summary ( merchant_id BIGINT, -- マーチャント ID order_id BIGINT, -- 注文 ID current_record_count BIGINT, -- 現在の有効なレコード数 total_amount DECIMAL(38, 2), -- 現在の合計金額 PRIMARY KEY (merchant_id, order_id) NOT ENFORCED ) WITH ( 'connector' = '<sink-connector>', '<sink-option>' = '<sink-option-value>' ); -
ステップ 2: 集計デプロイメントの実行
INSERT INTO文の前にデルタ集計を有効にします。SET 'table.optimizer.delta-agg.strategy' = 'EVENTUAL';標準のグループ集計 SQL を使用します。専用の構文は必要ありません。
INSERT INTO order_current_summary SELECT merchant_id, order_id, COUNT(*) AS current_record_count, SUM(amount) AS total_amount FROM orders GROUP BY merchant_id, order_id;この例では、
GROUP BY (merchant_id, order_id)はソーステーブルのプライマリキーを完全にカバーし、結果テーブルのプライマリキーはグループ化フィールドと一致します。 -
ステップ 3: 実行計画の検証
デプロイメントを送信する前に SQL 実行計画を確認します。
DeltaAggregateノードは、デルタ集計が有効になっていることを示します。例:DeltaAggregate( groupBy=[merchant_id, order_id], lookupKeys=[merchant_id, order_id] )
オプションのチューニング
まず、デフォルト設定でデプロイメントを実行します。次に、ルックアップのパフォーマンスとデプロイメントのスループットに基づいて個々のパラメーターをチューニングします。以下の各 SET 文を、対応する INSERT INTO 文の前に配置します。
有効化戦略の選択
table.optimizer.delta-agg.strategy は以下の値をサポートします。
値 |
説明 |
|
デルタ集計を強制します。クエリが要件を満たさない場合、実行計画の生成時にエラーが返されます。 |
|
デルタ集計の使用を試みます。クエリが要件を満たさない場合、エンジンは通常のグループ集計にフォールバックします。 |
|
デルタ集計を無効にします。これがデフォルト値です。 |
非同期ルックアップパラメーターのチューニング
以下のパラメーターを使用して、非同期ルックアップの同時実行数とタイムアウトをチューニングします。
SET 'table.exec.async-lookup.buffer-capacity' = '100';
SET 'table.exec.async-lookup.timeout' = '3 min';
パラメーター |
デフォルト |
説明 |
|
|
各並列サブタスクが同時に実行できる非同期ルックアップリクエストの数。この値を増やすと、非同期ルックアップの同時実行数が増加します。 |
|
|
単一の非同期操作のタイムアウト期間。 |
非同期ルックアップチューニングの詳細については、「ディメンションテーブルの JOIN 文」をご参照ください。
miniBatch デルタ集計の有効化
デルタ集計とグループ集計の両方が miniBatch 最適化をサポートしています。デルタ集計の場合、miniBatch はバッチ内で同じグループ化キーを持つ変更を結合します。これにより、各グループ化キーはバッチごとに最大 1 回のルックアップに制限され、ソーステーブルのルックアップ数が削減されます。
デルタ集計で miniBatch を有効にするには、以下の設定を使用します。
SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '10s';
SET 'table.exec.mini-batch.size' = '1000';
これらの設定により、各並列サブタスクは 1,000 件の入力レコードを受信した後、またはバッチが 10 秒間待機した後に処理をトリガーします。
miniBatch 最適化とそのパラメーターに関する詳細については、前述のセクション「miniBatch の有効化」をご参照ください。
デルタ集計キャッシュの有効化
同じグループ化キーが繰り返し更新される場合は、キャッシュを有効にして、以前に取得したグループ詳細レコードを再利用します。キャッシュヒットが発生すると、Flink はキャッシュされたレコードを入力変更で更新し、再度ルックアップすることなく直接集計結果を再計算します。
キャッシュが有効な場合、ソーステーブルにはプライマリキーが必要です。Flink はソーステーブルのプライマリキーを使用して、キャッシュ内の詳細レコードを識別および更新します。table.exec.delta-agg.cache-size は 0 より大きい必要があります。
SET 'table.exec.delta-agg.cache-enabled' = 'true';
SET 'table.exec.delta-agg.cache-size' = '2000';
パラメーター |
デフォルト |
説明 |
|
|
デルタ集計キャッシュを有効にするかどうかを指定します。 |
|
|
各デルタ集計の並列サブタスクによってキャッシュされるグループ化キーの最大数。このパラメーターは、キャッシュが有効な場合にのみ効果があります。 |
集計結果フィルターの検証の緩和
デルタ集計は、以前の集計結果に対して -U レコードを送信しません。下流のフィルター HAVING SUM(amount) >= 100 を考えてみましょう。グループの SUM(amount) が 120 から 80 に変更された場合:
+U(..., 120) → 条件に一致し、結果テーブルに書き込まれる
+U(..., 80) → 条件に一致せず、フィルターで除外される
以前の結果を撤回する -U レコードがないため、SUM(amount)=120 のレコードは結果テーブルに残ります。この問題を回避するため、デルタ集計では、下流のフィルターが集計結果のユニークキーフィールド (通常は GROUP BY フィールド) のみを参照することを要求します。
フィルターが SUM のような非ユニークキーフィールドを参照する場合、EVENTUAL 戦略はエラーを返し、AUTO 戦略は通常のグループ集計にフォールバックします。GROUP BY の前の WHERE 条件は、この検証の対象外です。
ビジネス要件としてこの動作を許容できる場合は、以下の設定を使用して検証をスキップします。
SET 'table.optimizer.delta-agg.ignore-non-unique-key-filter' = 'true';
パラメーター |
デフォルト |
説明 |
|
|
デルタ集計結果に対する非ユニークキーフィルターの検証をスキップするかどうかを指定します。 |
この設定は検証をスキップするだけで、SQL 文のフィルターを変更したり、デルタ集計が -U レコードを生成したりするわけではありません。したがって、Flink は上記で説明した古い結果を削除するための撤回レコードを送信しません。
結合の最適化
2ストリーム結合におけるキーバリュー分離
VVR 6.0.1 以降では、エンジンは 2 ストリーム JOIN 演算子に対してキーバリュー分離を有効にするかどうかを自動的に推論し、典型的な結合パフォーマンスを 40% 以上向上させます。
この動作を制御するには、table.exec.join.kv-separate パラメーターを設定します。
| 値 | 動作 |
|---|---|
AUTO |
エンジンは、結合演算子の状態に基づいてキーバリュー分離を自動的に有効にします。デフォルト。 |
FORCE |
キーバリュー分離を強制的に有効にします。 |
NONE |
キーバリュー分離を強制的に無効にします。 |
キーバリュー分離は GeminiStateBackend でのみ有効です。
結合に対する miniBatch の有効化
デフォルトでは、通常の結合演算子はレコードを 1 つずつ処理します。つまり、結合キーで相手側の状態をルックアップし、状態を更新し、結果を出力します。このレコードごとのパターンは、特に RocksDB を使用している場合に状態バックエンドのオーバーヘッドを増加させ、カスケード結合では深刻なレコード増幅を引き起こす可能性があります。
結合のための MiniBatch は、2 つのコアな最適化でこれに対処します。
-
レコードの折りたたみ:バッファー内のレコードを折りたたむことで、結合プロセス前のデータ量を削減します。
-
出力の抑制:バッチ処理中に冗長な中間結果を抑制します。
MiniBatch は、最大許容レイテンシーに達したとき、バッチサイズ制限に達したとき、またはチェックポイントが発生したときに処理をトリガーします。
VVR バージョン:8.0.4 以降。
有効化の方法:その他の設定 フィールドに以下を追加します。詳細については、「カスタムデプロイメントパラメーターの設定」をご参照ください。
table.exec.mini-batch.enabled: true
table.exec.mini-batch.allow-latency: 5s
table.exec.stream.join.mini-batch-enabled: true
最適なケース:メッセージ増幅を伴うカスケード外部結合。上流の演算子はメッセージの折りたたみまたはマージを使用して、下流の演算子での増幅を抑制します。例:
SELECT a.id AS a_id, a.a_content, B.id AS b_id, B.b_content
FROM a LEFT JOIN
(SELECT * FROM b
LEFT JOIN c ON b.prd_id = c.id) B
ON a.id = B.id
TopN クエリの最適化
TopN アルゴリズム
使用されるアルゴリズムは、入力データストリームが静的か動的かによって異なります。
| 入力タイプ | 利用可能なアルゴリズム | 注意 |
|---|---|---|
| 静的データストリーム (例:Simple Log Service から) | AppendRank | 静的ストリームの唯一のオプション。 |
| 動的データストリーム (例:集計や結合から) | UpdateFastRank, RetractRank | UpdateFastRank が最適です。RetractRank はフォールバックです。 |
アルゴリズム名はトポロジーノード名に表示されます。
RetractRank から UpdateFastRank への切り替え
UpdateFastRank には、以下のすべての条件が必要です。
-
入力ストリームに DELETE または UPDATE_BEFORE メッセージが含まれていないこと。EXPLAIN で確認します。
EXPLAIN CHANGELOG_MODE <query_statement_or_insert_statement_or_statement_set> -
入力ストリームにプライマリキー情報 (例:GROUP BY 句で使用される列) が含まれていること。
-
ORDER BY フィールドが、ソート順とは逆の順序で単調に更新されること。例:ORDER BY COUNT DESC、ORDER BY COUNT_DISTINCT DESC、または ORDER BY SUM(正の値) DESC。
例:ORDER BY SUM DESC の場合、正の値にフィルタリングして単調性を保証します。
INSERT INTO print_test
SELECT
cate_id,
seller_id,
stat_date,
pay_ord_amt
FROM (
SELECT
*,
ROW_NUMBER() OVER (
-- PARTITION BY 列はサブクエリの GROUP BY に現れる必要があります。
-- 状態が期限切れになったときの順序の乱れを防ぐために時間フィールドを含めます。
PARTITION BY cate_id, stat_date
ORDER BY pay_ord_amt DESC
) AS rownum
FROM (
SELECT
cate_id,
seller_id,
stat_date,
-- 正の値のみが含まれるため、SUM は単調増加します。
SUM(total_fee) FILTER (WHERE total_fee >= 0) AS pay_ord_amt
FROM random_test
WHERE total_fee >= 0
GROUP BY seller_id, stat_date, cate_id
) a
)
WHERE rownum <= 100;
この例では、random_test テーブルには静的なデータストリームが含まれています。集計結果には DELETE または UPDATE_BEFORE メッセージが含まれていないため、単調性が維持されます。
出力ボリュームの削減
最終的な SELECT 出力に rownum を含めないでください。代わりにフロントエンドで結果をソートすることで、sink の書き込み量を削減します。詳細については、「Top-N」をご参照ください。
TopN キャッシュサイズの増加
TopN は、ディスクの読み取りを減らすために状態キャッシュを維持します。キャッシュヒット率は次の数式で計算します。
キャッシュヒット率 = キャッシュサイズ * 並列度 / top_n / パーティションキー数
例:Top100、デフォルトのキャッシュサイズ 10,000 レコード、並列度 50、パーティションキー数 100,000 の場合:
10,000 * 50 / 100 / 100,000 = 5% のヒット率
5% のヒット率は、ほとんどの読み取りがディスクに向かうことを意味し、state seek メトリックが不安定になり、パフォーマンスが低下します。キャッシュサイズを増やします。
table.exec.rank.topn-cache-size: 200000
200,000 のキャッシュエントリの場合:
200,000 * 50 / 100 / 100,000 = 100% のヒット率
パーティションキー数が多い場合は、TopN キャッシュサイズとヒープメモリも増やしてください。詳細については、「デプロイメントの設定」をご参照ください。
PARTITION BY への時間フィールドの追加
day などの時間フィールドを PARTITION BY 句に追加します。これがないと、TTL により状態データが期限切れになったときに TopN の結果が乱れます。
効率的な重複排除
入力ストリームにはしばしば重複が含まれます。Realtime Compute for Apache Flink は、2 つの重複排除ポリシーを提供します。
-
Deduplicate Keep FirstRow:キーごとに最も早いレコードを保持します。
-
重複排除 (最終行保持):キーごとに最新のレコードを保持します。
構文
重複排除は TopN の特殊なケースです。ROW_NUMBER() と OVER 句を使用してパーティションごとに行番号を割り当て、rownum = 1 のみ保持します。
SELECT *
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY col1[, col2..]
ORDER BY timeAttributeCol [ASC|DESC]) AS rownum
FROM table_name)
WHERE rownum = 1
| 要素 | 説明 |
|---|---|
ROW_NUMBER() |
各パーティション内で 1 から始まる行番号を割り当てます。 |
PARTITION BY col1[, col2..] |
重複排除キーを定義する列。 |
ORDER BY timeAttributeCol |
時間属性列 (proctime または rowtime)。ASC は最初の行を保持 (Keep FirstRow)。DESC は最後の行を保持 (Keep LastRow)。 |
rownum = 1 |
パーティションごとに最初の行のみを保持します。rownum <= 1 もサポートします。 |
時間属性の動作:
-
proctime(処理時間):重複排除はレコードが処理された時間に基づきます。結果は実行ごとに異なる場合があります。 -
rowtime(イベント時間):重複排除はレコードが生成された時間に基づきます。結果は決定的です。
FirstRow の保持
重複排除キーごとに最初のレコードを保持します。状態はプライマリキーデータのみを保存するため、状態アクセスは非常に効率的です。
SELECT *
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY b ORDER BY proctime) AS rowNum
FROM T
)
WHERE rowNum = 1
テーブル T を列 b で重複排除し、処理時間で最も早いレコードを保持します。proctime 属性を宣言する代わりに、PROCTIME() 関数を呼び出します。
LastRow の保持
重複排除キーごとに最新のレコードを保持します。これは LAST_VALUE 関数よりもわずかにパフォーマンスが優れています。
SELECT *
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY b, d ORDER BY rowtime DESC) AS rowNum
FROM T
)
WHERE rowNum = 1
テーブル T を列 b と d で重複排除し、イベント時間で最新のレコードを保持します。
ビルトイン関数の効率的な使用
UDF よりもビルトイン関数を優先
ビルトイン関数は、シリアル化、デシリアル化、およびバイトレベルの操作に最適化されています。可能な限り、UDF をビルトインの同等機能に置き換えてください。
KEYVALUE における 1 文字デリミタの使用
KEYVALUE(content, keyValueSplit, keySplit, keyName) 関数は、keyValueSplit と keySplit が : や , などの 1 文字である場合、約 30% 高速に実行されます。1 文字のデリミタを使用すると、エンジンは入力全体を解析することなく、バイナリデータ内で直接ターゲットキーを検索します。
LIKE 演算子のパターン
| パターン | 一致 | 例 |
|---|---|---|
LIKE 'xxx%' |
xxx |
LIKE 'order%' |
LIKE '%xxx' |
xxx |
LIKE '%_id' |
LIKE '%xxx%' |
xxx |
LIKE '%error%' |
LIKE 'xxx' |
完全一致 (= 'xxx' と同等) |
LIKE 'active' |
アンダースコアのエスケープ:アンダースコア (_) は SQL における 1 文字のワイルドカードです。リテラルのアンダースコアに一致させるには、エスケープ文字を使用します。
LIKE '%seller/_id%' ESCAPE '/'
エスケープしない場合、LIKE '%seller_id%' は seller#id、sellerxid、seller1id、およびその他の意図しない文字列にも一致します。
正規表現の回避
正規表現は算術演算よりも 100 倍遅くなる可能性があり、エッジケースでは無限ループに陥り、デプロイメントをブロックすることがあります。可能な場合は、代わりに LIKE を使用してください。
正規表現が必要なケースについては、「REGEXP」および「REGEXP_REPLACE」をご参照ください。
SQL ヒント
SQL ヒントは、オプティマイザーの実行計画に影響を与え、メタデータや統計をアタッチし、テーブルごとに 動的テーブルオプション を設定します。詳細については、「SQL ヒント」をご参照ください。
構文
構文は Apache Calcite SQL に従います。
SELECT /*+ hint [, hint ] */ ...
-- ここで:
-- hint: hintName(hintOption [, hintOption]*)
-- hintOption: simpleIdentifier | numericLiteral | stringLiteral
結合ヒント
クエリヒントは、現在の クエリブロック 内の実行計画を変更します。Flink は現在、ディメンションテーブル結合 と 通常結合 の両方に対して、結合ヒントのみをサポートしています。