Python UDF、UDAF、UDWF、UDTF
Python UDF/UDAF/UDTF は Apache Doris 4.1.3 でリリースされた実験的機能です。
Python UDF/UDAF/UDTF は Apache Doris が提供するカスタム関数拡張メカニズムです。Python でスカラー関数、集約関数、テーブル関数を記述できるため、SQL で組み込み関数では実装が困難な複雑な計算ロジックを表現でき、Python の豊富なエコシステムを再利用できます。
本ドキュメントでは典型的なユーザーシナリオから始まり、3つの関数タイプそれぞれについて、使用方法、パラメータ、データ型マッピング、パフォーマンス推奨事項、制限事項、およびマルチバージョン Python 環境のデプロイメントについて説明します。
Python UDF/UDAF/UDTF を選択するタイミング
| シナリオ | 推奨 | 関係性 |
|---|---|---|
| 行ごとの複雑な変換、クレンジング、マスキング、または検証 | Python UDF(スカラー関数) | 1行入力で1行出力 |
| GROUP BY またはウィンドウ(OVER句)によるカスタム集約メトリクス | Python UDAF(集約関数) | 複数行入力で1行出力 |
| CSV/JSON解析やシーケンス生成など、1行を複数行に展開 | Python UDTF(テーブル関数) | 1行入力で0行または複数行出力 |
パフォーマンスが重要な場合は、Doris 組み込み関数(C++で実装)を優先してください。Python UDF は組み込み関数では要件を満たせない場合で、データ量が中程度のシナリオに適しています。
一般的な前提条件
Python UDF/UDAF/UDTF を作成する前に、以下の準備を完了してください:
- Python UDF を有効化し Python 環境を設定:BE ノード
be.confで関連パラメータを有効化し、Conda または venv を使用してマルチバージョン Python 環境を設定します。詳細は Python UDF/UDAF/UDTF 環境設定とマルチバージョン管理 を参照してください。 - 必須依存関係:全 BE ノードの対応する Python 環境に
pandasとpyarrowを事前インストールしてください。これらは Doris Python UDF 機能の必須依存関係で、不足している場合は関数が実行できません。 - ランタイムログ:Python UDF Server のランタイムログは
output/be/log/python_udf_output.logにあります。このログを確認して、デバッグのための関数実行とエラーメッセージを表示できます。
全ての CREATE ステートメントでは、runtime_version に完全なバージョン番号("3.10.12" など)を明示的に指定する必要があります。メジャー・マイナーバージョンのみ("3.10" など)を指定することはできません。そうでなければ関数呼び出しが失敗します。
Python UDF(スカラー関数)
Python UDF(User Defined Function)はデータを行ごとに処理します。関数は行ごとに一度呼び出され、単一の結果を返します。2つの実行モードをサポートしています:
- スカラーモード:データを行ごとに処理します。シンプルな変換と計算に適しています。
- ベクトル化モード:Pandas の助けを借りてデータをバッチで処理し、高性能な計算を実現します。
Python UDF の作成
Python UDF は2つの作成方法をサポートしています:インラインモードとモジュールモードです。
file パラメータと AS $$ インライン Python コードの両方が指定されている場合、Doris はインライン Python コードを優先し、関数をインラインモードで実行します。
インラインモード
インラインモードでは、Python コードを SQL 内で直接記述できます。シンプルなロジックに適しています。
構文:
CREATE FUNCTION function_name(parameter_type1, parameter_type2, ...)
RETURNS return_type
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "entry_function_name",
"runtime_version" = "python_version",
"always_nullable" = "true|false",
"volatility" = "immutable|stable|volatile"
)
AS $$
def entry_function_name(param1, param2, ...):
# Python code here
return result
$$;
例1: 整数の加算
DROP FUNCTION IF EXISTS py_add(INT, INT);
CREATE FUNCTION py_add(INT, INT)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"volatility" = "immutable"
)
AS $$
def evaluate(a, b):
return a + b
$$;
SELECT py_add(10, 20) AS result; -- Result: 30
例2: 文字列連結(NULL処理を含む)
DROP FUNCTION IF EXISTS py_concat(STRING, STRING);
CREATE FUNCTION py_concat(STRING, STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"volatility" = "immutable"
)
AS $$
def evaluate(s1, s2):
if s1 is None or s2 is None:
return None
return s1 + s2
$$;
SELECT py_concat('Hello', ' World') AS result; -- Result: Hello World
SELECT py_concat(NULL, ' World') AS result; -- Result: NULL
SELECT py_concat('Hello', NULL) AS result; -- Result: NULL
Module Mode
Module modeは複雑なロジックに適しています。Pythonコードを.zipアーカイブとしてパッケージ化し、関数を作成する際にfileパラメータで参照します。
ステップ1: Pythonモジュールを作成する
python_udf_scalar_ops.pyという名前のファイルを作成します:
def add_three_numbers(a, b, c):
"""Add three numbers"""
if a is None or b is None or c is None:
return None
return a + b + c
def reverse_string(s):
"""Reverse a string"""
if s is None:
return None
return s[::-1]
def is_prime(n):
"""Check if a number is prime"""
if n is None or n < 2:
return False
if n == 2:
return True
if n % 2 == 0:
return False
import math
for i in range(3, int(math.sqrt(n)) + 1, 2):
if n % i == 0:
return False
return True
ステップ2: Pythonモジュールをパッケージ化する
Pythonファイルを.zip形式でパッケージ化する必要があります(ファイルが1つだけの場合でも):
zip python_udf_scalar_ops.zip python_udf_scalar_ops.py
複数のPythonファイルがある場合:
zip python_udf_scalar_ops.zip python_udf_scalar_ops.py utils.py helper.py ...
ステップ3: .zipパッケージのパスを設定する
fileパラメータを通じて.zipパッケージのパスを指定します。2つの方法がサポートされています:
| デプロイ方法 | 形式 | 適用シナリオ |
|---|---|---|
| ローカルファイルシステム | "file" = "file:///path/to/python_udf_scalar_ops.zip" | .zipパッケージがBEノードのローカルファイルシステムに保存されている場合 |
| HTTP/HTTPSリモートダウンロード | "file" = "http://example.com/udf/xx.zip" または "file" = "https://s3.amazonaws.com/bucket/xx.zip" | オブジェクトストレージ(S3、OSS、COSなど)またはHTTPサーバーから.zipパッケージをダウンロードする場合。Dorisが自動的にダウンロードしてローカルにキャッシュします |
- リモートダウンロードを使用する場合は、すべてのBEノードがURLにアクセスできることを確認してください。
- 最初の呼び出しでファイルをダウンロードするため、多少のレイテンシが発生する可能性があります。
- ファイルはキャッシュされるため、後続の呼び出しでは再度ダウンロードされません。
ステップ4: symbolパラメータを設定する
モジュールモードでは、symbolはZIPパッケージ内のターゲット関数の場所を指定します。形式は以下の通りです:
[package_name.]module_name.func_name
パラメータの説明:
package_name(オプション):ZIPパッケージ内のトップレベルPythonパッケージの名前。関数がパッケージのルートモジュールにある場合、またはZIPパッケージにパッケージが含まれていない場合は省略します。module_name(必須):対象の関数を含むPythonモジュールファイル名(.py拡張子を除く)。func_name(必須):ユーザー定義関数名。
解決ルール:
- Dorisは
symbol文字列を.で分割します:- 結果が2つの部分文字列の場合、それらは
module_nameとfunc_nameです。 - 結果が3つ以上の部分文字列の場合、最初が
package_name、中間がmodule_name、最後がfunc_nameです。
- 結果が2つの部分文字列の場合、それらは
module_name部分は、importlibによる動的インポートのモジュールパスとして機能します。package_nameが指定されている場合、パス全体が有効なPythonインポートパスを形成する必要があり、ZIPパッケージ構造がパスと一致する必要があります。
名前空間は一意である必要があります。Python標準ライブラリや一般的なサードパーティライブラリと衝突する名前は避けて、依存関係の競合やモジュールシャドウイングによるランタイム例外を防いでください。
例A:パッケージ構造なし(2部構成)
ZIP structure:
math_ops.py
symbol = "math_ops.add"
これは、関数 add が ZIP パッケージのルートにある math_ops.py で定義されていることを示しています。
例 B: パッケージ構造あり(3部構成)
ZIP structure:
mylib/
├── __init__.py
└── string_helper.py
symbol = "mylib.string_helper.split_text"
これは、関数split_textがmylib/string_helper.pyで定義されていることを示しており、ここで:
package_name=mylibmodule_name=string_helperfunc_name=split_text
例C:ネストされたパッケージ構造(4部構成)
ZIP structure:
mylib/
├── __init__.py
└── utils/
├── __init__.py
└── string_helper.py
symbol = "mylib.utils.string_helper.split_text"
これは関数 split_text が mylib/utils/string_helper.py で定義されていることを示しており、以下のようになります:
package_name=mylibmodule_name=utils.string_helperfunc_name=split_text
注意:
symbolの形式が無効な場合(関数名の欠如、空のモジュール名、空のパスコンポーネントなど)、Dorisは関数呼び出し時にエラーを報告します。- ZIPパッケージ内のディレクトリ構造は
symbolで指定されたパスと一致する必要があります。- 各パッケージディレクトリには
__init__.pyファイル(空でも可)を含める必要があります。
ステップ5: UDFを作成する
例1: ローカルファイルを使用する(パッケージ構造なし)
DROP FUNCTION IF EXISTS py_add_three(INT, INT, INT);
DROP FUNCTION IF EXISTS py_reverse(STRING);
DROP FUNCTION IF EXISTS py_is_prime(INT);
CREATE FUNCTION py_add_three(INT, INT, INT)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/python_udf_scalar_ops.zip",
"symbol" = "python_udf_scalar_ops.add_three_numbers",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE FUNCTION py_reverse(STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/python_udf_scalar_ops.zip",
"symbol" = "python_udf_scalar_ops.reverse_string",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE FUNCTION py_is_prime(INT)
RETURNS BOOLEAN
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/python_udf_scalar_ops.zip",
"symbol" = "python_udf_scalar_ops.is_prime",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
例2: HTTP/HTTPSリモートファイルを使用する
DROP FUNCTION IF EXISTS py_add_three(INT, INT, INT);
DROP FUNCTION IF EXISTS py_reverse(STRING);
DROP FUNCTION IF EXISTS py_is_prime(INT);
CREATE FUNCTION py_add_three(INT, INT, INT)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "https://your-storage.com/udf/python_udf_scalar_ops.zip",
"symbol" = "python_udf_scalar_ops.add_three_numbers",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE FUNCTION py_reverse(STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "https://your-storage.com/udf/python_udf_scalar_ops.zip",
"symbol" = "python_udf_scalar_ops.reverse_string",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE FUNCTION py_is_prime(INT)
RETURNS BOOLEAN
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "https://your-storage.com/udf/python_udf_scalar_ops.zip",
"symbol" = "python_udf_scalar_ops.is_prime",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
例3: パッケージ構造を使用する
DROP FUNCTION IF EXISTS py_multiply(INT);
-- ZIP structure: my_udf/__init__.py, my_udf/math_ops.py
CREATE FUNCTION py_multiply(INT)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/my_udf.zip",
"symbol" = "my_udf.math_ops.multiply_by_two",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
ステップ 6: 関数を使用する
SELECT py_add_three(10, 20, 30) AS sum_result; -- Result: 60
SELECT py_reverse('hello') AS reversed; -- Result: olleh
SELECT py_is_prime(17) AS is_prime; -- Result: true
Python UDFの削除
-- Syntax
DROP FUNCTION IF EXISTS function_name(parameter_type1, parameter_type2, ...);
-- Example
DROP FUNCTION IF EXISTS py_add_three(INT, INT, INT);
DROP FUNCTION IF EXISTS py_reverse(STRING);
DROP FUNCTION IF EXISTS py_is_prime(INT);
パラメータリファレンス
CREATE FUNCTION パラメータ
| パラメータ | 必須 | 説明 |
|---|---|---|
function_name | はい | 関数名。識別子の命名規則に従う必要があります |
parameter_type | はい | パラメータ型のリスト。様々なDorisデータ型をサポートします |
return_type | はい | 戻り値の型 |
PROPERTIES パラメータ
| パラメータ | 必須 | デフォルト | 説明 |
|---|---|---|---|
type | はい | - | 固定値 "PYTHON_UDF" |
symbol | はい | - | Python関数のエントリ名。 • インラインモード: "evaluate"などの関数名を直接記述• モジュールモード: 形式は [package_name.]module_name.func_name。詳細はモジュールモードの説明を参照 |
file | いいえ | - | Pythonの.zipパッケージへのパス。モジュールモードでのみ必須。3つのプロトコルをサポート:• file://: ローカルファイルシステムのパス• http://: HTTPリモートダウンロード• https://: HTTPSリモートダウンロード |
runtime_version | はい | - | "3.10.12"などのPythonランタイムバージョン。完全なバージョン番号が必要です |
always_nullable | いいえ | true | 関数が常にnull許可の結果を返すかどうか |
volatility | いいえ | volatile | 関数の揮発性、4.1.2以降でサポート。有効な値はimmutable、stable、volatileです。immutable: 同じ入力は常に同じ出力を生成します。決定論的なUDFの多くは、より良いプラン最適化のためにこの値を使用すべきです。stable: 1つのSQL文内では結果は安定していますが、文間では変わる可能性があります。now()と同様です。SQLキャッシュとマテリアライズドビューの書き換えは無効になります。volatile: 呼び出しごとに結果が変わる可能性があります。random()と同様です。SQLキャッシュ、マテリアライズドビューの書き換え、および多くのオプティマイザーの書き換えルールが無効になります。 |
ランタイムバージョンの注意事項
- Python 3.xがサポートされています。
- 完全なバージョン番号を指定する必要があります(
"3.10.12"など);メジャーとマイナーバージョンのみ("3.10"など)は許可されません。 runtime_versionが指定されていない場合、関数の呼び出しは失敗します。
データ型マッピング
次の表は、DorisデータタイプとPythonタイプ間のマッピングを示しています:
| タイプカテゴリ | Dorisタイプ | Pythonタイプ | 説明 |
|---|---|---|---|
| Null型 | NULL | None | Null値 |
| Boolean型 | BOOLEAN | bool | Boolean値 |
| Integer型 | TINYINT | int | 8ビット整数 |
SMALLINT | int | 16ビット整数 | |
INT | int | 32ビット整数 | |
BIGINT | int | 64ビット整数 | |
LARGEINT | int | 128ビット整数 | |
| 浮動小数点型 | FLOAT | float | 32ビット浮動小数点 |
DOUBLE | float | 64ビット浮動小数点 | |
TIME / TIMEV2 | float | 時間型(浮動小数点として表現) | |
| 文字列型 | CHAR | str | 固定長文字列 |
VARCHAR | str | 可変長文字列 | |
STRING | str | 文字列 | |
JSONB | str | JSONバイナリ形式(文字列に変換) | |
VARIANT | str | Variant型(文字列に変換) | |
DATE | str | 'YYYY-MM-DD'形式の日付文字列 | |
DATETIME | str | 'YYYY-MM-DD HH:MM:SS'形式の日時文字列 | |
| 日付/時刻型 | DATEV2 | datetime.date | 日付オブジェクト |
DATETIMEV2 | datetime.datetime | 日時オブジェクト | |
TIMESTAMPTZ | datetime.datetime | タイムゾーン付き日時オブジェクト | |
| 小数型 | DECIMAL / DECIMALV2 | decimal.Decimal | 高精度小数 |
DECIMAL32 | decimal.Decimal | 32ビット固定小数点 | |
DECIMAL64 | decimal.Decimal | 64ビット固定小数点 | |
DECIMAL128 | decimal.Decimal | 128ビット固定小数点 | |
DECIMAL256 | decimal.Decimal | 256ビット固定小数点 | |
| IP型 | IPV4 | ipaddress.IPv4Address | IPv4アドレス |
IPV6 | ipaddress.IPv6Address | IPv6アドレス | |
| バイナリ型 | BITMAP | bytes | Bitmapデータ(まだサポートされていません) |
HLL | bytes | HyperLogLogデータ(まだサポートされていません) | |
QUANTILE_STATE | bytes | 分位数状態データ(まだサポートされていません) | |
| 複合データ型 | ARRAY<T> | list | 要素型Tの配列 |
MAP<K,V> | dict | キー型KおよびValue型Vの辞書 | |
STRUCT<f1:T1, f2:T2, ...> | dict | フィールド名をキー、フィールド値を値とする構造体 |
NULL処理
- Dorisの
NULL値はPythonのNoneにマップされます。 - 関数の引数が
NULLの場合、Python関数はNoneを受け取ります。 - Python関数が
Noneを返す場合、DorisはそれをNULLとして扱います。 - ランタイムエラーを避けるために、関数内で
None値を明示的に処理してください。
例:
DROP FUNCTION IF EXISTS py_safe_divide(DOUBLE, DOUBLE);
CREATE FUNCTION py_safe_divide(DOUBLE, DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(a, b):
if a is None or b is None:
return None
if b == 0:
return None
return a / b
$$;
SELECT py_safe_divide(10.0, 2.0); -- Result: 5.0
SELECT py_safe_divide(10.0, 0.0); -- Result: NULL
SELECT py_safe_divide(10.0, NULL); -- Result: NULL
ベクトル化モード
ベクトル化モードはPandasを使用してデータをバッチで処理し、スカラーモードより優れたパフォーマンスを発揮します。ベクトル化モードでは、関数の引数はpandas.Seriesオブジェクトであり、戻り値もpandas.Seriesである必要があります。
システムがベクトル化モードを認識するように、関数シグネチャで型注釈(a: pd.Seriesなど)を使用し、関数内でバッチデータ構造を直接操作してください。明示的なベクトル化型がない場合、システムはスカラーモードにフォールバックします。
関数シグネチャでpd.Series型と通常の型が混在している場合、システムは通常の型パラメータに対応する入力列を定数列として扱います(同じ値がバッチ全体で再利用される)。これにより期待と一致しない結果が生成される可能性があります。ベクトル化モードでは、パラメータのスタイルを一貫させてください。すべてのパラメータにpandas.Series型注釈を使用するか、すべてに通常の型パラメータを使用する(スカラーモード)かのいずれかにしてください。
## Vectorized mode
def add(a: pd.Series, b: pd.Series) -> pd.Series:
return a + b + 1
## Scalar mode
def add(a, b):
return a + b + 1
基本的な例
例 1: ベクトル化された整数加算
DROP FUNCTION IF EXISTS py_vec_add(INT, INT);
CREATE FUNCTION py_vec_add(INT, INT)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "add",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
import pandas as pd
def add(a: pd.Series, b: pd.Series) -> pd.Series:
return a + b + 1
$$;
SELECT py_vec_add(1, 2); -- Result: 4
例2: ベクトル化された文字列処理
DROP FUNCTION IF EXISTS py_vec_upper(STRING);
CREATE FUNCTION py_vec_upper(STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "to_upper",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
import pandas as pd
def to_upper(s: pd.Series) -> pd.Series:
return s.str.upper()
$$;
SELECT py_vec_upper('hello'); -- Result: 'HELLO'
例3: ベクトル化された数学演算
DROP FUNCTION IF EXISTS py_vec_sqrt(DOUBLE);
CREATE FUNCTION py_vec_sqrt(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "sqrt",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
import pandas as pd
import numpy as np
def sqrt(x: pd.Series) -> pd.Series:
return np.sqrt(x)
$$;
SELECT py_vec_sqrt(16); -- Result: 4.0
例4: 関数シグネチャでの混在パラメータ型(pd.Seriesと通常の型の両方)
CREATE TABLE t_bug_013 (
id INT,
a INT,
b INT
) ENGINE=OLAP
DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 1
PROPERTIES ("replication_num" = "1");
INSERT INTO t_bug_013 VALUES
(1, 1, 10),
(2, 2, 20),
(3, 3, 30),
(4, 4, NULL),
(5, NULL, 50);
DROP FUNCTION IF EXISTS py_mixed_vector_add(INT, INT);
CREATE FUNCTION py_mixed_vector_add(INT, INT)
RETURNS INT
PROPERTIES (
"type"="PYTHON_UDF",
"symbol"="py_mixed_vector_add_impl",
"always_nullable"="true",
"runtime_version"="3.12.11",
"volatility"="immutable"
)
AS $$
import pandas as pd
# Keep the parameter style consistent
def py_mixed_vector_add_impl(x: pd.Series, y: int):
return x + y
$$;
SELECT
id,
a,
b,
py_mixed_vector_add(a, b) AS vector_val
FROM t_bug_013
ORDER BY id;
-- Column b is treated as a constant column
+------+------+------+------------+
| id | a | b | vector_val |
+------+------+------+------------+
| 1 | 1 | 10 | 11 |
| 2 | 2 | 20 | 12 |
| 3 | 3 | 30 | 13 |
| 4 | 4 | NULL | 14 |
| 5 | NULL | 50 | NULL |
+------+------+------+------------+
ベクトル化モードの利点
- パフォーマンスの最適化: データをバッチで処理し、PythonとDoris間のやり取りの回数を削減します。
- Pandas/NumPyの活用: ベクトル化計算を最大限に活用します。
- 簡潔なコード: Pandas APIにより複雑なロジックをより簡潔に表現できます。
ベクトル化関数の使用
DROP TABLE IF EXISTS test_table;
CREATE TABLE test_table (
id INT,
value INT,
text STRING,
score DOUBLE
) ENGINE=OLAP
DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO test_table VALUES
(1, 10, 'hello', 85.5),
(2, 20, 'world', 92.0),
(3, 30, 'python', 78.3);
SELECT
id,
py_vec_add(value, value) AS sum_result,
py_vec_upper(text) AS upper_text,
py_vec_sqrt(score) AS sqrt_score
FROM test_table;
+------+------------+------------+-------------------+
| id | sum_result | upper_text | sqrt_score |
+------+------------+------------+-------------------+
| 1 | 21 | HELLO | 9.246621004453464 |
| 2 | 41 | WORLD | 9.591663046625438 |
| 3 | 61 | PYTHON | 8.848728722251575 |
+------+------------+------------+-------------------+
複雑なデータ型の処理
ARRAY型
例: 配列要素の合計
DROP FUNCTION IF EXISTS py_array_sum(ARRAY<INT>);
CREATE FUNCTION py_array_sum(ARRAY<INT>)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(arr):
""" The Doris ARRAY type maps to a Python list """
if arr is None:
return None
return sum(arr)
$$;
SELECT py_array_sum([1, 2, 3, 4, 5]) AS result; -- Result: 15
例: 配列をフィルタする
DROP FUNCTION IF EXISTS py_array_filter_positive(ARRAY<INT>);
CREATE FUNCTION py_array_filter_positive(ARRAY<INT>)
RETURNS ARRAY<INT>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(arr):
if arr is None:
return None
return [x for x in arr if x > 0]
$$;
SELECT py_array_filter_positive([1, -2, 3, -4, 5]) AS result; -- Result: [1, 3, 5]
MAP型
例: MAPのキー数を取得する
DROP FUNCTION IF EXISTS py_map_size(MAP<STRING, INT>);
CREATE FUNCTION py_map_size(MAP<STRING, INT>)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(m):
""" The Doris MAP type maps to a Python dict """
if m is None:
return None
return len(m)
$$;
SELECT py_map_size({'a': 1, 'b': 2, 'c': 3}) AS result; -- Result: 3
例: MAPから値を取得する
DROP FUNCTION IF EXISTS py_map_get(MAP<STRING, STRING>, STRING);
CREATE FUNCTION py_map_get(MAP<STRING, STRING>, STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(m, key):
if m is None or key is None:
return None
return m.get(key)
$$;
SELECT py_map_get({'name': 'Alice', 'age': '30'}, 'name') AS result; -- Result: Alice
STRUCT型
例: STRUCTフィールドへのアクセス
DROP FUNCTION IF EXISTS py_struct_get_name(STRUCT<name: STRING, age: INT>);
CREATE FUNCTION py_struct_get_name(STRUCT<name: STRING, age: INT>)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(s):
""" The Doris STRUCT type maps to a Python dict """
if s is None:
return None
return s.get('name')
$$;
SELECT py_struct_get_name({'Alice', 30}) AS result; -- Result: Alice
実際のシナリオ
シナリオ1: データマスキング
DROP FUNCTION IF EXISTS py_mask_email(STRING);
CREATE FUNCTION py_mask_email(STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"volatility" = "immutable"
)
AS $$
def evaluate(email):
if email is None or '@' not in email:
return None
parts = email.split('@')
if len(parts[0]) <= 1:
return email
masked_user = parts[0][0] + '***'
return f"{masked_user}@{parts[1]}"
$$;
SELECT py_mask_email('user@example.com') AS masked; -- Result: u***@example.com
シナリオ2: 文字列類似度計算
DROP FUNCTION IF EXISTS py_levenshtein_distance(STRING, STRING);
CREATE FUNCTION py_levenshtein_distance(STRING, STRING)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"volatility" = "immutable"
)
AS $$
def evaluate(s1, s2):
if s1 is None or s2 is None:
return None
if len(s1) < len(s2):
return evaluate(s2, s1)
if len(s2) == 0:
return len(s1)
previous_row = range(len(s2) + 1)
for i, c1 in enumerate(s1):
current_row = [i + 1]
for j, c2 in enumerate(s2):
insertions = previous_row[j + 1] + 1
deletions = current_row[j] + 1
substitutions = previous_row[j] + (c1 != c2)
current_row.append(min(insertions, deletions, substitutions))
previous_row = current_row
return previous_row[-1]
$$;
SELECT py_levenshtein_distance('kitten', 'sitting') AS distance; -- Result: 3
シナリオ3: 日付計算
DROP FUNCTION IF EXISTS py_days_between(DATE, DATE);
CREATE FUNCTION py_days_between(DATE, DATE)
RETURNS INT
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"volatility" = "immutable"
)
AS $$
from datetime import datetime
def evaluate(date1_str, date2_str):
if date1_str is None or date2_str is None:
return None
try:
d1 = datetime.strptime(str(date1_str), '%Y-%m-%d')
d2 = datetime.strptime(str(date2_str), '%Y-%m-%d')
return abs((d2 - d1).days)
except:
return None
$$;
SELECT py_days_between('2024-01-01', '2024-12-31') AS days; -- Result: 365
シナリオ4: IDカード番号検証
DROP FUNCTION IF EXISTS py_validate_id_card(STRING);
CREATE FUNCTION py_validate_id_card(STRING)
RETURNS BOOLEAN
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.10.12",
"volatility" = "immutable"
)
AS $$
def evaluate(id_card):
if id_card is None or len(id_card) != 18:
return False
# Verify that the first 17 characters are digits
if not id_card[:17].isdigit():
return False
# Check code weights
weights = [7, 9, 10, 5, 8, 4, 2, 1, 6, 3, 7, 9, 10, 5, 8, 4, 2]
check_codes = ['1', '0', 'X', '9', '8', '7', '6', '5', '4', '3', '2']
# Compute the check code
total = sum(int(id_card[i]) * weights[i] for i in range(17))
check_code = check_codes[total % 11]
return id_card[17].upper() == check_code
$$;
SELECT py_validate_id_card('11010519491231002X') AS is_valid; -- Result: True
SELECT py_validate_id_card('110105194912310021x') AS is_valid; -- Result: False
パフォーマンス推奨事項
1. Vectorizedモードを優先する
Vectorizedモードはscalarモードよりも大幅に優れたパフォーマンスを発揮します:
# Scalar mode: row-by-row processing
def scalar_process(x):
return x * 2
# Vectorized mode: batch processing
import pandas as pd
def vector_process(x: pd.Series) -> pd.Series:
return x * 2
2. 複雑なロジックにはモジュールモードを使用する
保守性と再利用性を向上させるため、複雑な関数ロジックは独立したPythonファイルに配置してください。
3. 関数内でのI/O操作を避ける
UDF内でファイルI/O、ネットワークリクエスト、その他のI/O操作を避けてください。これらはパフォーマンスに深刻な影響を与えます。
制限事項と注意事項
1. Pythonバージョンサポート
- Python 3.xのみサポートされています。
- Python 3.10以降を推奨します。
- Dorisクラスターに対応するPythonランタイムがインストールされていることを確認してください。
2. 依存関係ライブラリ
- Python標準ライブラリはすぐに使用できます。
- サードパーティライブラリを使用する場合は、事前にクラスター環境にインストールしてください。
3. パフォーマンスに関する考慮事項
- Python UDFのパフォーマンスは、Doris組み込み関数(C++で実装)よりも低くなります。
- パフォーマンスに敏感なシナリオでは、Doris組み込み関数を優先してください。
- 大容量データの場合は、ベクトル化モードを使用してください。
4. セキュリティ
- UDFコードはDorisプロセス内で実行されるため、コードは安全で信頼できるものでなければなりません。
- UDF内で危険な操作(システムコマンドやファイル削除など)を避けてください。
- 本番環境ではUDFコードをレビューしてください。
5. リソース制限
- UDF実行はBEノード上のCPUとメモリリソースを消費します。
- UDFの多用はクラスター全体のパフォーマンスに影響を与える可能性があります。
- UDFのリソース消費を監視してください。
よくある質問
Q1: Python UDFでサードパーティライブラリを使用するにはどうすればよいですか?
A: すべてのBEノードに対応するPythonライブラリをインストールしてください。例:
pip3 install numpy pandas
conda install numpy pandas
Q2: Python UDFは再帰関数をサポートしていますか?
A: はい、ただしスタックオーバーフローを避けるために再帰の深度に注意してください。
Q3: Python UDFをデバッグするにはどうすればよいですか?
A: まずローカルのPython環境で関数のロジックをデバッグして正しいことを確認してから、UDFを作成してください。エラー情報についてはBEログを確認してください。
Q4: Python UDFはグローバル変数をサポートしていますか?
A: はい、ただし推奨されません。分散環境では、グローバル変数の動作が期待通りにならない場合があります。
Q5: 既存のPython UDFを更新するにはどうすればよいですか?
A: まず古いUDFを削除してから、新しいものを作成してください:
DROP FUNCTION IF EXISTS function_name(parameter_types);
CREATE FUNCTION function_name(...) ...;
Q6: Python UDFは外部リソースにアクセスできますか?
A: 技術的には可能ですが、強く非推奨です。Python UDF内でネットワークライブラリ(requestsなど)を使用して外部APIやデータベースにアクセスできますが、これはパフォーマンスと安定性に深刻な影響を与えます。理由は以下の通りです:
- ネットワーク遅延がクエリを遅くする。
- 外部サービスが利用できない場合にUDFが失敗する。
- 大量の同時リクエストが外部サービスに負荷をかける可能性がある。
- タイムアウトとエラー処理の制御が困難。
Python UDAF(集約関数)
Python UDAF(User Defined Aggregate Function)は、グループ化された集約とウィンドウ計算のためのカスタム集約関数を定義できます。Python UDAFを使用すると、統計分析、データ収集、カスタムメトリクス計算などの複雑な集約ロジックを柔軟に実装できます。
Python UDAFの主要な特徴:
- 分散集約: 分散環境での集約をサポートし、データパーティション分割、マージ、最終計算を自動的に処理します。
- 状態管理: クラスインスタンスを通じて集約状態を維持し、複雑な状態オブジェクトをサポートします。
- ウィンドウ関数サポート: ウィンドウ関数(OVER句)と連携して、移動集約、ランキング、その他の高度な機能を提供します。
- 高い柔軟性: 組み込み集約関数に制限されることなく、任意に複雑な集約ロジックを実装します。
UDAF基本概念
集約関数ライフサイクル
Python UDAFはクラスとして実装されます。集約関数の実行は以下の段階を含みます:
- 初期化(
__init__): 集約状態オブジェクトを作成し、状態変数を初期化します。 - 蓄積(
accumulate): 単一の行を処理し、集約状態を更新します。 - マージ(
merge): 複数のパーティションからの集約状態をマージします(分散シナリオで)。 - 完了(
finish): 最終集約結果を計算して返します。
必須クラスメソッドと属性
完全なPython UDAFクラスは以下を実装する必要があります:
| メソッド/属性 | 説明 | 必須 |
|---|---|---|
__init__(self) | 集約状態を初期化 | はい |
accumulate(self, *args) | 単一行からのデータを蓄積 | はい |
merge(self, other_state) | 他のパーティションからの状態をマージ | はい |
finish(self) | 最終集約結果を返す | はい |
aggregate_state(属性) | シリアライズ可能な集約状態を返す。pickleシリアライゼーションをサポートする必要があります | はい |
基本構文
Python UDAFの作成
Python UDAFはインラインモードとモジュールモードの2つの作成方法をサポートしています。
fileパラメータとAS $$インラインPythonコードの両方が指定されている場合、DorisはインラインPythonコードを優先し、インラインモードでPython UDAFを実行します。
インラインモード
インラインモードでは、SQL内で直接Pythonクラスを記述できます。シンプルな集約ロジックに適しています。
構文:
CREATE AGGREGATE FUNCTION function_name(parameter_type1, parameter_type2, ...)
RETURNS return_type
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "ClassName",
"runtime_version" = "python_version",
"always_nullable" = "true|false"
)
AS $$
class ClassName:
def __init__(self):
# Initialize state variables
@property
def aggregate_state(self):
# Return the serializable state
def accumulate(self, *args):
# Accumulate data
def merge(self, other_state):
# Merge state
def finish(self):
# Return the final result
$$;
例1: Sum集約
DROP TABLE IF EXISTS sales;
CREATE TABLE IF NOT EXISTS sales (
id INT,
category VARCHAR(50),
amount INT
) DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO sales VALUES
(1, 'Electronics', 1000),
(2, 'Electronics', 1500),
(3, 'Books', 200),
(4, 'Books', 300),
(5, 'Clothing', 500),
(6, 'Clothing', 800),
(7, 'Electronics', 2000),
(8, 'Books', 150);
DROP FUNCTION IF EXISTS py_sum(INT);
CREATE AGGREGATE FUNCTION py_sum(INT)
RETURNS BIGINT
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "SumUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
class SumUDAF:
def __init__(self):
self.total = 0
@property
def aggregate_state(self):
return self.total
def accumulate(self, value):
if value is not None:
self.total += value
def merge(self, other_state):
self.total += other_state
def finish(self):
return self.total
$$;
SELECT category, py_sum(amount) as total_amount
FROM sales
GROUP BY category
ORDER BY category;
+-------------+--------------+
| category | total_amount |
+-------------+--------------+
| Books | 650 |
| Clothing | 1300 |
| Electronics | 4500 |
+-------------+--------------+
例2: 平均集計
DROP TABLE IF EXISTS employees;
CREATE TABLE IF NOT EXISTS employees (
id INT,
name VARCHAR(100),
department VARCHAR(50),
salary DOUBLE
) DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO employees VALUES
(1, 'Alice', 'Engineering', 80000.0),
(2, 'Bob', 'Engineering', 90000.0),
(3, 'Charlie', 'Sales', 60000.0),
(4, 'David', 'Sales', 80000.0),
(5, 'Eve', 'HR', 50000.0),
(6, 'Frank', 'Engineering', 70000.0),
(7, 'Grace', 'HR', 70000.0);
DROP FUNCTION IF EXISTS py_avg(DOUBLE);
CREATE AGGREGATE FUNCTION py_avg(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "AvgUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
class AvgUDAF:
def __init__(self):
self.sum = 0.0
self.count = 0
@property
def aggregate_state(self):
return (self.sum, self.count)
def accumulate(self, value):
if value is not None:
self.sum += value
self.count += 1
def merge(self, other_state):
other_sum, other_count = other_state
self.sum += other_sum
self.count += other_count
def finish(self):
if self.count == 0:
return None
return self.sum / self.count
$$;
SELECT department, py_avg(salary) as avg_salary
FROM employees
GROUP BY department
ORDER BY department;
+-------------+------------+
| department | avg_salary |
+-------------+------------+
| Engineering | 80000 |
| HR | 60000 |
| Sales | 70000 |
+-------------+------------+
Module Mode
Module modeは複雑な集約ロジックに適しています。Pythonコードを.zipアーカイブとしてパッケージ化し、関数を作成する際にそれを参照します。
Step 1: Pythonモジュールを書く
stats_udaf.pyという名前のファイルを作成します:
import math
class VarianceUDAF:
"""Compute the population variance"""
def __init__(self):
self.count = 0
self.sum_val = 0.0
self.sum_sq = 0.0
@property
def aggregate_state(self):
return (self.count, self.sum_val, self.sum_sq)
def accumulate(self, value):
if value is not None:
self.count += 1
self.sum_val += value
self.sum_sq += value * value
def merge(self, other_state):
other_count, other_sum, other_sum_sq = other_state
self.count += other_count
self.sum_val += other_sum
self.sum_sq += other_sum_sq
def finish(self):
if self.count == 0:
return None
mean = self.sum_val / self.count
variance = (self.sum_sq / self.count) - (mean * mean)
return variance
class StdDevUDAF:
"""Compute the population standard deviation"""
def __init__(self):
self.count = 0
self.sum_val = 0.0
self.sum_sq = 0.0
@property
def aggregate_state(self):
return (self.count, self.sum_val, self.sum_sq)
def accumulate(self, value):
if value is not None:
self.count += 1
self.sum_val += value
self.sum_sq += value * value
def merge(self, other_state):
other_count, other_sum, other_sum_sq = other_state
self.count += other_count
self.sum_val += other_sum
self.sum_sq += other_sum_sq
def finish(self):
if self.count == 0:
return None
mean = self.sum_val / self.count
variance = (self.sum_sq / self.count) - (mean * mean)
return math.sqrt(max(0, variance))
class MedianUDAF:
"""Compute the median"""
def __init__(self):
self.values = []
@property
def aggregate_state(self):
return self.values
def accumulate(self, value):
if value is not None:
self.values.append(value)
def merge(self, other_state):
if other_state:
self.values.extend(other_state)
def finish(self):
if not self.values:
return None
sorted_vals = sorted(self.values)
n = len(sorted_vals)
if n % 2 == 0:
return (sorted_vals[n//2 - 1] + sorted_vals[n//2]) / 2.0
else:
return sorted_vals[n//2]
ステップ2: Pythonモジュールをパッケージ化する
Pythonファイルは.zip形式でパッケージ化する必要があります(ファイルが1つだけの場合でも):
zip stats_udaf.zip stats_udaf.py
ステップ 3: .zip パッケージのパスを設定する
file パラメータを通じて .zip パッケージパスを指定します:
| デプロイ方法 | 形式 |
|---|---|
ローカルファイルシステム(file:// プロトコル) | "file" = "file:///path/to/stats_udaf.zip" |
HTTP/HTTPS リモートダウンロード(http:// または https:// プロトコル) | "file" = "http://example.com/udaf/stats_udaf.zip""file" = "https://s3.amazonaws.com/bucket/stats_udaf.zip" |
注意:
- リモートダウンロードを使用する場合は、すべてのBEノードがURLにアクセスできることを確認してください。
- 最初の呼び出しではファイルをダウンロードするため、多少の遅延が発生する可能性があります。
- ファイルはキャッシュされるため、以降の呼び出しでは再度ダウンロードされません。
ステップ 4: symbol パラメータを設定する
モジュールモードでは、symbol はZIPパッケージ内のクラスの場所を指定します。形式は以下の通りです:
[package_name.]module_name.ClassName
パラメータ説明:
package_name(任意):ZIPパッケージ内のトップレベルPythonパッケージの名前。module_name(必須):対象クラスを含むPythonモジュールファイル名(.pyサフィックスを除く)。ClassName(必須):UDAFクラス名。
解決ルール:
- Dorisは
symbol文字列を.で分割します:- 結果が2つの部分文字列の場合、それらは
module_nameとClassNameになります。 - 結果が3つ以上の部分文字列の場合、最初が
package_name、中間がmodule_name、最後がClassNameになります。
- 結果が2つの部分文字列の場合、それらは
名前空間は一意である必要があります。Python標準ライブラリや一般的なサードパーティライブラリと衝突する名前は避けて、モジュールシャドウイングによる依存関係の競合やランタイム例外を防いでください。
ステップ5:UDAFを作成する
DROP FUNCTION IF EXISTS py_variance(DOUBLE);
DROP FUNCTION IF EXISTS py_stddev(DOUBLE);
DROP FUNCTION IF EXISTS py_median(DOUBLE);
CREATE AGGREGATE FUNCTION py_variance(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/stats_udaf.zip",
"symbol" = "stats_udaf.VarianceUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE AGGREGATE FUNCTION py_stddev(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/stats_udaf.zip",
"symbol" = "stats_udaf.StdDevUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE AGGREGATE FUNCTION py_median(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/stats_udaf.zip",
"symbol" = "stats_udaf.MedianUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
ステップ 6: 関数を使用する
DROP TABLE IF EXISTS exam_results;
CREATE TABLE IF NOT EXISTS exam_results (
id INT,
student_name VARCHAR(100),
category VARCHAR(50),
score DOUBLE
) DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO exam_results VALUES
(1, 'Alice', 'Math', 85.0),
(2, 'Bob', 'Math', 92.0),
(3, 'Charlie', 'Math', 78.0),
(4, 'David', 'Math', 88.0),
(5, 'Eve', 'Math', 95.0),
(6, 'Frank', 'English', 75.0),
(7, 'Grace', 'English', 82.0),
(8, 'Henry', 'English', 88.0),
(9, 'Iris', 'English', 79.0),
(10, 'Jack', 'Physics', 90.0),
(11, 'Kate', 'Physics', 85.0),
(12, 'Lily', 'Physics', 92.0),
(13, 'Mike', 'Physics', 88.0);
SELECT
category,
py_variance(score) as variance,
py_stddev(score) as std_dev,
py_median(score) as median
FROM exam_results
GROUP BY category
ORDER BY category;
+----------+-------------------+-------------------+--------+
| category | variance | std_dev | median |
+----------+-------------------+-------------------+--------+
| English | 22.5 | 4.743416490252569 | 80.5 |
| Math | 34.64000000000033 | 5.885575587824892 | 88 |
| Physics | 6.6875 | 2.58602010819715 | 89 |
+----------+-------------------+-------------------+--------+
Python UDAFの削除
-- Syntax
DROP FUNCTION IF EXISTS function_name(parameter_types);
-- Example
DROP FUNCTION IF EXISTS py_sum(INT);
DROP FUNCTION IF EXISTS py_avg(DOUBLE);
DROP FUNCTION IF EXISTS py_variance(DOUBLE);
パラメータリファレンス
CREATE AGGREGATE FUNCTION パラメータ
| パラメータ | 説明 |
|---|---|
function_name | 関数名。SQL識別子の命名規則に従います |
parameter_types | INT、DOUBLE、STRINGなどのパラメータ型リスト |
RETURNS return_type | 戻り値の型 |
PROPERTIES パラメータ
| パラメータ | 必須 | デフォルト | 説明 |
|---|---|---|---|
type | はい | - | 固定値 "PYTHON_UDF" |
symbol | はい | - | Pythonクラス名。 • インラインモード: "SumUDAF"のようにクラス名を直接記述• モジュールモード: 形式は [package_name.]module_name.ClassName |
file | いいえ | - | Python .zipパッケージへのパス。モジュールモードでのみ必須。3つのプロトコルをサポート:• file://: ローカルファイルシステムパス• http://: HTTPリモートダウンロード• https://: HTTPSリモートダウンロード |
runtime_version | はい | - | "3.10.12"などのPythonランタイムバージョン |
always_nullable | いいえ | true | 関数が常にnull許可の結果を返すかどうか |
volatility | いいえ | immutable | 関数の揮発性、4.1.2以降でサポート。有効な値は immutable、stable、volatile。immutable: 同じ入力は常に同じ出力を生成。ほとんどの決定論的UDFは、プラン最適化を向上させるためにこの値を使用する必要があります。stable: 結果は1つのSQLステートメント内では安定していますが、ステートメント間では変更される可能性があり、now()と類似。SQLキャッシュとマテリアライズドビューの書き換えは無効になります。volatile: 結果は各呼び出しで変更される可能性があり、random()と類似。SQLキャッシュ、マテリアライズドビューの書き換え、および多くのオプティマイザー書き換えルールが無効になります。 |
runtime_version 注意事項
- Pythonバージョンは
x.x.xまたはx.x.xxの形式で完全なバージョン番号として指定する必要があります。 - Dorisは、設定されたPython環境でこのバージョンに一致するインタープリターを検索します。
ウィンドウ関数
Python UDAFをウィンドウ関数(OVER句)と組み合わせることができます:
Python UDAFをウィンドウ関数(OVER句)で使用する場合、Dorisは各ウィンドウフレームの計算後にUDAFの
resetメソッドを呼び出します。クラスでこのメソッドを実装して、集約状態を初期値にリセットしてください。
DROP TABLE IF EXISTS daily_sales_data;
CREATE TABLE IF NOT EXISTS daily_sales_data (
sales_date DATE,
daily_sales DOUBLE
) DUPLICATE KEY(sales_date)
DISTRIBUTED BY HASH(sales_date) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO daily_sales_data VALUES
('2024-01-01', 1000),
('2024-01-01', 800),
('2024-01-02', 1200),
('2024-01-02', 950),
('2024-01-03', 900),
('2024-01-03', 1100),
('2024-01-04', 1500),
('2024-01-04', 850),
('2024-01-05', 1100),
('2024-01-05', 1300);
DROP FUNCTION IF EXISTS py_running_sum(DOUBLE);
CREATE AGGREGATE FUNCTION py_running_sum(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "RunningSumUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
class RunningSumUDAF:
def __init__(self):
self.total = 0.0
def reset(self):
self.total = 0.0
@property
def aggregate_state(self):
return self.total
def accumulate(self, value):
if value is not None:
self.total += value
def merge(self, other_state):
self.total += other_state
def finish(self):
return self.total
$$;
SELECT
sales_date,
daily_sales,
py_running_sum(daily_sales) OVER (
ORDER BY sales_date
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW
) as last_3_days_sum
FROM daily_sales_data
ORDER BY sales_date;
+------------+-------------+-----------------+
| sales_date | daily_sales | last_3_days_sum |
+------------+-------------+-----------------+
| 2024-01-01 | 800 | 800 |
| 2024-01-01 | 1000 | 1800 |
| 2024-01-02 | 950 | 2750 |
| 2024-01-02 | 1200 | 3150 |
| 2024-01-03 | 1100 | 3250 |
| 2024-01-03 | 900 | 3200 |
| 2024-01-04 | 850 | 2850 |
| 2024-01-04 | 1500 | 3250 |
| 2024-01-05 | 1300 | 3650 |
| 2024-01-05 | 1100 | 3900 |
+------------+-------------+-----------------+
データ型マッピング
Python UDAFは、すべての整数、浮動小数点、文字列、日時、decimal、およびboolean型を含む、Python UDFと同じデータ型マッピングルールを使用します。
詳細な型マッピングについては、以下を参照してください: Data Type Mapping。
NULL処理
- DorisはSQLの
NULL値をPythonのNoneにマップします。 accumulateメソッドで、パラメータがNoneかどうかを確認してください。- 集約関数は結果が
NULLであることを示すためにNoneを返すことができます。
実際のシナリオ
シナリオ1: パーセンタイルの計算
DROP FUNCTION IF EXISTS py_percentile(DOUBLE, INT);
CREATE AGGREGATE FUNCTION py_percentile(DOUBLE, INT)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "PercentileUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
class PercentileUDAF:
"""Compute a percentile. The second argument is the percentile (0-100)"""
def __init__(self):
self.values = []
self.percentile = 50 # Median by default
@property
def aggregate_state(self):
return self.values
def accumulate(self, value, percentile):
if value is not None:
self.values.append(value)
if percentile is not None:
self.percentile = percentile
def merge(self, other_state):
if other_state:
self.values.extend(other_state)
def finish(self):
if not self.values:
return None
sorted_vals = sorted(self.values)
n = len(sorted_vals)
k = (n - 1) * (self.percentile / 100.0)
f = int(k)
c = k - f
if f + 1 < n:
return sorted_vals[f] + (sorted_vals[f + 1] - sorted_vals[f]) * c
else:
return sorted_vals[f]
$$;
DROP TABLE IF EXISTS api_logs;
CREATE TABLE IF NOT EXISTS api_logs (
log_id INT,
api_name VARCHAR(100),
category VARCHAR(50),
response_time DOUBLE
) DUPLICATE KEY(log_id)
DISTRIBUTED BY HASH(log_id) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO api_logs VALUES
(1, '/api/users', 'User', 120.5),
(2, '/api/users', 'User', 95.3),
(3, '/api/users', 'User', 150.0),
(4, '/api/users', 'User', 80.2),
(5, '/api/users', 'User', 200.8),
(6, '/api/orders', 'Order', 250.0),
(7, '/api/orders', 'Order', 180.5),
(8, '/api/orders', 'Order', 300.2),
(9, '/api/orders', 'Order', 220.0),
(10, '/api/products', 'Product', 50.0),
(11, '/api/products', 'Product', 60.5),
(12, '/api/products', 'Product', 45.0),
(13, '/api/products', 'Product', 70.2),
(14, '/api/products', 'Product', 55.8);
SELECT
category,
py_percentile(response_time, 25) as p25,
py_percentile(response_time, 50) as p50,
py_percentile(response_time, 75) as p75,
py_percentile(response_time, 95) as p95
FROM api_logs
GROUP BY category
ORDER BY category;
+----------+-------+-------+-------+-------+
| category | p25 | p50 | p75 | p95 |
+----------+-------+-------+-------+-------+
| Order | 235 | 235 | 235 | 235 |
| Product | 55.8 | 55.8 | 55.8 | 55.8 |
| User | 120.5 | 120.5 | 120.5 | 120.5 |
+----------+-------+-------+-------+-------+
シナリオ2:重複排除された文字列コレクション
DROP FUNCTION IF EXISTS py_collect_set(STRING);
CREATE AGGREGATE FUNCTION py_collect_set(STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "CollectSetUDAF",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
class CollectSetUDAF:
"""Collect deduplicated strings and return a comma-separated string"""
def __init__(self):
self.items = set()
@property
def aggregate_state(self):
return list(self.items)
def accumulate(self, value):
if value is not None:
self.items.add(value)
def merge(self, other_state):
if other_state:
self.items.update(other_state)
def finish(self):
if not self.items:
return None
return ','.join(sorted(self.items))
$$;
DROP TABLE IF EXISTS page_views;
CREATE TABLE IF NOT EXISTS page_views (
view_id INT,
user_id INT,
page_url VARCHAR(200),
view_time DATETIME
) DUPLICATE KEY(view_id)
DISTRIBUTED BY HASH(view_id) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO page_views VALUES
(1, 1001, '/home', '2024-01-01 10:00:00'),
(2, 1001, '/products', '2024-01-01 10:05:00'),
(3, 1001, '/home', '2024-01-01 10:10:00'),
(4, 1001, '/cart', '2024-01-01 10:15:00'),
(5, 1002, '/home', '2024-01-01 11:00:00'),
(6, 1002, '/about', '2024-01-01 11:05:00'),
(7, 1002, '/products', '2024-01-01 11:10:00'),
(8, 1003, '/products', '2024-01-01 12:00:00'),
(9, 1003, '/products', '2024-01-01 12:05:00'),
(10, 1003, '/cart', '2024-01-01 12:10:00'),
(11, 1003, '/checkout', '2024-01-01 12:15:00');
SELECT
user_id,
py_collect_set(page_url) as visited_pages
FROM page_views
GROUP BY user_id
ORDER BY user_id;
+---------+---------------------------+
| user_id | visited_pages |
+---------+---------------------------+
| 1001 | /cart,/home,/products |
| 1002 | /about,/home,/products |
| 1003 | /cart,/checkout,/products |
+---------+---------------------------+
シナリオ 3: 移動平均
DROP TABLE IF EXISTS daily_sales;
CREATE TABLE IF NOT EXISTS daily_sales (
id INT,
date DATE,
sales DOUBLE
) DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 1
PROPERTIES("replication_num" = "1");
INSERT INTO daily_sales VALUES
(1, '2024-01-01', 1000.0),
(2, '2024-01-02', 1200.0),
(3, '2024-01-03', 900.0),
(4, '2024-01-04', 1500.0),
(5, '2024-01-05', 1100.0),
(6, '2024-01-06', 1300.0),
(7, '2024-01-07', 1400.0),
(8, '2024-01-08', 1000.0),
(9, '2024-01-09', 1600.0),
(10, '2024-01-10', 1250.0);
SELECT
date,
sales,
py_avg(sales) OVER (
ORDER BY date
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
) as moving_avg_7days
FROM daily_sales
ORDER BY date;
+------------+-------+-------------------+
| date | sales | moving_avg_7days |
+------------+-------+-------------------+
| 2024-01-01 | 1000 | 1000 |
| 2024-01-02 | 1200 | 1100 |
| 2024-01-03 | 900 | 1033.333333333333 |
| 2024-01-04 | 1500 | 1150 |
| 2024-01-05 | 1100 | 1140 |
| 2024-01-06 | 1300 | 1166.666666666667 |
| 2024-01-07 | 1400 | 1200 |
| 2024-01-08 | 1000 | 1200 |
| 2024-01-09 | 1600 | 1257.142857142857 |
| 2024-01-10 | 1250 | 1307.142857142857 |
+------------+-------+-------------------+
パフォーマンス推奨事項
1. stateオブジェクトのサイズを最適化する
- stateオブジェクトに大量の生データを格納することは避けてください。
- 完全なデータリストの代わりに、可能な限り集約された統計情報を使用してください。
- データを格納する必要があるシナリオ(中央値計算など)では、サンプリングやデータ量の制限を検討してください。
推奨されません:
class BadMedianUDAF:
def __init__(self):
self.all_values = [] # Can be very large
def accumulate(self, value):
if value is not None:
self.all_values.append(value)
2. オブジェクト作成の削減
- stateオブジェクトを再利用し、新しいオブジェクトを頻繁に作成することを避ける。
- 複雑なオブジェクトの代わりにプリミティブデータ型を使用する。
3. マージロジックの簡素化
mergeメソッドは分散環境で頻繁に呼び出される。- マージ操作が効率的かつ正確であることを確認する。
4. インクリメンタル計算の使用
- インクリメンタルに計算できるメトリクス(平均など)については、すべてのデータを保存する代わりにインクリメンタル計算を使用する。
5. 外部リソースの使用を避ける
- UDAFでデータベースや外部APIにアクセスしない。
- すべての計算は入力データと内部stateに基づいて行う必要がある。
制限事項と注意点
1. パフォーマンスに関する考慮事項
- Python UDAFのパフォーマンスは組み込み集計関数よりも低い。
- 複雑なロジックで中程度のデータ量のシナリオで使用する。
- 大量のデータについては、組み込み関数を優先するか、UDAF実装を最適化する。
2. State のシリアライゼーション
aggregate_stateが返すオブジェクトはpickle シリアライゼーションをサポートしている必要がある。- サポートされる型:プリミティブ型(int、float、str、bool)、list、dict、tuple、set、およびpickle シリアライゼーションをサポートするカスタムクラスのインスタンス。
- サポートされない:ファイルハンドル、データベース接続、socket接続、threadロック、およびpickleできないその他のオブジェクト。
- stateオブジェクトがpickleできない場合、関数は実行時に失敗する。
- 互換性と保守性を確保するため、stateオブジェクトには組み込み型(dict、list、tuple)を優先する。
3. メモリ制限
- stateオブジェクトはメモリを消費する。大量のデータを保存することを避ける。
- 大きなstateオブジェクトはパフォーマンスと安定性に影響する。
4. 関数の命名
- 同じ関数名を異なるデータベースで定義できる。
- 呼び出し時にデータベース名を指定(
db.func()など)して曖昧さを避ける。
5. 環境の一貫性
- すべてのBEノードのPython環境は一貫している必要がある。
- これにはPythonバージョン、依存パッケージのバージョン、環境設定が含まれる。
よくある質問
Q1: UDAFとUDFの違いは何ですか?
A: UDFは単一行を処理して単一結果を返します。関数は行ごとに一度呼び出されます。UDAFは複数行を処理して単一の集計結果を返し、GROUP BYと組み合わせて使用されます。
-- UDF: invoked for each row
SELECT id, py_upper(name) FROM users;
-- UDAF: invoked once per group
SELECT category, py_sum(amount) FROM sales GROUP BY category;
Q2: aggregate_state属性は何をしますか?
A: aggregate_stateは分散環境で集約状態をシリアル化し送信するために使用されます:
- シリアル化: pickleプロトコルを使用して状態オブジェクトを送信可能な形式に変換します。
- マージ: ノード間で部分的な集約結果をマージします。
- pickleシリアル化をサポートする必要があります: プリミティブ型、リスト、辞書、タプル、セット、およびpickleシリアル化をサポートするカスタムクラスインスタンスを返すことができます。
- 返すことが許可されていません: ファイルハンドル、データベース接続、ソケット接続、スレッドロック、またはpickleできないその他のオブジェクト。そうでなければ実行時に関数が失敗します。
Q3: UDAFはウィンドウ関数で使用できますか?
A: はい。Python UDAFはウィンドウ関数(OVER句)を完全にサポートしています。
Q4: mergeメソッドはいつ呼び出されますか?
A: mergeは以下の状況で呼び出されます:
- 分散集約: 異なるBEノードからの部分的な集約結果をマージする際。
- 並列処理: 同じノード上の異なるスレッドからの部分結果をマージする際。
- ウィンドウ関数: ウィンドウフレーム内で部分結果をマージする際。
したがって、mergeの実装は正確である必要があり、そうでなければ結果が間違ってしまいます。
Python UDTF (Table Function)
Python UDTF(User Defined Table Function)を使用すると、単一の行を複数の出力行に変換するカスタムテーブル関数を定義できます。これは、データの分割、展開、および生成に有用です。
Python UDTFの主要な特徴:
- 1行から多数行へ: 単一の行を入力として受け取り、0個、1個、または複数の出力行を生成します。
- 柔軟な出力構造: 任意の数と型の出力列を許可し、シンプル型と複雑なSTRUCT型をサポートします。
- Lateral viewサポート: データの展開と結合のために
LATERAL VIEWと連携します。 - 関数型スタイル: Python関数と
yield文を使用し、簡潔で直感的です。
UDTF基本概念
テーブル関数の実行方法
Python UDTFは関数(クラスではなく)として実装されます。実行フローは以下のとおりです:
- 入力を受け取る: 関数は単一行の列値をパラメータとして受け取ります。
- 処理してyieldする:
yield文で0個以上の出力行をyieldします。 - ステートレス: 各関数呼び出しは1つの行を独立して処理し、前の行の状態を保持しません。
関数要件
Python UDTF関数は以下の要件を満たす必要があります:
yieldで結果をyieldする:yield文を使用して出力行を生成します。- パラメータ型と一致: 関数パラメータはSQLで定義されたパラメータ型と対応します。
- 出力形式と一致: yieldされるデータの形式は
RETURNS ARRAY<...>定義と一致する必要があります。
出力方法
- 単一列出力:
yield valueで単一の値をyieldします。 - 複数列出力:
yield (value1, value2, ...)で値のタプルをyieldします。 - 条件付きスキップ:
yieldを呼び出さないと、その行に対して出力を生成しません。
基本構文
Python UDTFの作成
Python UDTFは2つの作成方法をサポートします:インラインモードとモジュールモード。
fileパラメータとAS $$インラインPythonコードの両方が指定された場合、DorisはインラインPythonコードを優先し、Python UDTFをインラインモードで実行します。
インラインモード
インラインモードでは、Python関数をSQL内に直接記述できます。シンプルなテーブル関数ロジックに適しています。
構文:
CREATE TABLES FUNCTION function_name(parameter_type1, parameter_type2, ...)
RETURNS ARRAY<return_type>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "function_name",
"runtime_version" = "python_version",
"always_nullable" = "true|false"
)
AS $$
def function_name(param1, param2, ...):
'''Function description'''
# Processing logic
yield result # Single column output
# Or
yield (result1, result2, ...) # Multi-column output
$$;
重要な構文メモ:
CREATE TABLES FUNCTIONを使用してください(TABLESは複数形であることに注意)。- 単一列出力:
ARRAY<type>、例:ARRAY<INT>。- 複数列出力:
ARRAY<STRUCT<col1:type1, col2:type2, ...>>。
例1:文字列分割(単一列出力)
DROP FUNCTION IF EXISTS py_split(STRING, STRING);
CREATE TABLES FUNCTION py_split(STRING, STRING)
RETURNS ARRAY<STRING>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "split_string_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def split_string_udtf(text, delimiter):
'''Split a string into multiple rows by the delimiter'''
if text is not None and delimiter is not None:
parts = text.split(delimiter)
for part in parts:
# yield (part.strip(),) is also supported
yield part.strip()
$$;
SELECT part
FROM (SELECT 'apple,banana,orange' as fruits) t
LATERAL VIEW py_split(fruits, ',') tmp AS part;
+--------+
| part |
+--------+
| apple |
| banana |
| orange |
+--------+
例2: 数列を生成する(単一列出力)
DROP FUNCTION IF EXISTS py_range(INT, INT);
CREATE TABLES FUNCTION py_range(INT, INT)
RETURNS ARRAY<INT>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "generate_series_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def generate_series_udtf(start, end):
'''Generate an integer sequence from start to end'''
if start is not None and end is not None:
for i in range(start, end + 1):
yield i
$$;
SELECT num
FROM (SELECT 1 as start_val, 5 as end_val) t
LATERAL VIEW py_range(start_val, end_val) tmp AS num;
+------+
| num |
+------+
| 1 |
| 2 |
| 3 |
| 4 |
| 5 |
+------+
SELECT date_add('2024-01-01', n) as date
FROM (SELECT 0 as start_val, 6 as end_val) t
LATERAL VIEW py_range(start_val, end_val) tmp AS n;
+------------+
| date |
+------------+
| 2024-01-01 |
| 2024-01-02 |
| 2024-01-03 |
| 2024-01-04 |
| 2024-01-05 |
| 2024-01-06 |
| 2024-01-07 |
+------------+
例3: 複数列出力 (STRUCT)
DROP FUNCTION IF EXISTS py_duplicate(STRING, INT);
CREATE TABLES FUNCTION py_duplicate(STRING, INT)
RETURNS ARRAY<STRUCT<output:STRING, idx:INT>>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "duplicate_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def duplicate_udtf(text, n):
'''Duplicate text n times, each with a sequence number'''
if text is not None and n is not None:
for i in range(n):
yield (text, i + 1)
$$;
SELECT output, idx
FROM (SELECT 'Hello' as text, 3 as times) t
LATERAL VIEW py_duplicate(text, times) tmp AS output, idx;
+--------+------+
| output | idx |
+--------+------+
| Hello | 1 |
| Hello | 2 |
| Hello | 3 |
+--------+------+
例4: デカルト積(複数列のSTRUCT)
DROP FUNCTION IF EXISTS py_cartesian(STRING, STRING);
CREATE TABLES FUNCTION py_cartesian(STRING, STRING)
RETURNS ARRAY<STRUCT<item1:STRING, item2:STRING>>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "cartesian_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def cartesian_udtf(list1, list2):
'''Generate the Cartesian product of two lists'''
if list1 is not None and list2 is not None:
items1 = [x.strip() for x in list1.split(',')]
items2 = [y.strip() for y in list2.split(',')]
for x in items1:
for y in items2:
yield (x, y)
$$;
SELECT item1, item2
FROM (SELECT 'A,B' as list1, 'X,Y,Z' as list2) t
LATERAL VIEW py_cartesian(list1, list2) tmp AS item1, item2;
+-------+-------+
| item1 | item2 |
+-------+-------+
| A | X |
| A | Y |
| A | Z |
| B | X |
| B | Y |
| B | Z |
+-------+-------+
例5: JSON配列の解析
DROP FUNCTION IF EXISTS py_explode_json(STRING);
CREATE TABLES FUNCTION py_explode_json(STRING)
RETURNS ARRAY<STRING>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "explode_json_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
import json
def explode_json_udtf(json_str):
'''Parse a JSON array and output one row per element'''
if json_str is not None:
try:
data = json.loads(json_str)
if isinstance(data, list):
for item in data:
yield (str(item),)
except:
pass # Skip on parse failure
$$;
SELECT element
FROM (SELECT '["apple", "banana", "cherry"]' as json_data) t
LATERAL VIEW py_explode_json(json_data) tmp AS element;
+---------+
| element |
+---------+
| apple |
| banana |
| cherry |
+---------+
Module Mode
Module modeは複雑なテーブル関数ロジックに適しています。Pythonコードを.zipアーカイブとしてパッケージ化し、関数作成時に参照します。
Step 1: Pythonモジュールを作成する
text_udtf.pyという名前のファイルを作成します:
import json
import re
def split_lines_udtf(text):
"""Split text by line"""
if text:
lines = text.split('\n')
for line in lines:
line = line.strip()
if line: # Filter out empty lines
yield (line,)
def extract_emails_udtf(text):
"""Extract all email addresses from text"""
if text:
email_pattern = r'[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}'
emails = re.findall(email_pattern, text)
for email in emails:
yield (email,)
def parse_json_object_udtf(json_str):
"""Parse a JSON object and output key-value pairs"""
if json_str:
try:
data = json.loads(json_str)
if isinstance(data, dict):
for key, value in data.items():
yield (key, str(value))
except:
pass
def expand_json_array_udtf(json_str):
"""Expand the objects in a JSON array and output structured data"""
if json_str:
try:
data = json.loads(json_str)
if isinstance(data, list):
for item in data:
if isinstance(item, dict):
# Assume each object has id, name, and score fields
item_id = item.get('id')
name = item.get('name')
score = item.get('score')
yield (item_id, name, score)
except:
pass
def ngram_udtf(text, n):
"""Generate N-grams"""
if text and n and n > 0:
words = text.split()
for i in range(len(words) - n + 1):
ngram = ' '.join(words[i:i+n])
yield (ngram,)
ステップ2: Pythonモジュールをパッケージ化する
Pythonファイルを.zip形式でパッケージ化する必要があります(ファイルが1つだけの場合でも):
zip text_udtf.zip text_udtf.py
ステップ3: .zipパッケージのパスを設定する
fileパラメータを通じて.zipパッケージのパスを指定します:
| デプロイ方法 | 形式 |
|---|---|
ローカルファイルシステム(file://プロトコル) | "file" = "file:///path/to/text_udtf.zip" |
HTTP/HTTPSリモートダウンロード(http:// またはhttps:// プロトコル) | "file" = "http://example.com/udtf/text_udtf.zip""file" = "https://s3.amazonaws.com/bucket/text_udtf.zip" |
- リモートダウンロードを使用する場合、すべてのBEノードがURLにアクセスできることを確認してください。
- 最初の呼び出しではファイルをダウンロードするため、多少の遅延が発生する可能性があります。
- ファイルはキャッシュされるため、その後の呼び出しでは再度ダウンロードされません。
ステップ4: symbolパラメータを設定する
モジュールモードでは、symbolはZIPパッケージ内の関数の場所を指定します。形式は以下の通りです:
[package_name.]module_name.function_name
パラメータの説明:
package_name(オプション):ZIPパッケージ内のトップレベルPythonパッケージの名前。module_name(必須):対象の関数を含むPythonモジュールファイル名(.py拡張子なし)。function_name(必須):UDTF関数名。
解決ルール:
- Dorisは
symbol文字列を.で分割します:- 結果が2つの部分文字列の場合、それらは
module_nameとfunction_nameです。 - 結果が3つ以上の部分文字列の場合、最初が
package_name、中間がmodule_name、最後がfunction_nameです。
- 結果が2つの部分文字列の場合、それらは
名前空間は一意である必要があります。Pythonの標準ライブラリや一般的なサードパーティライブラリと衝突する名前は避けて、依存関係の競合やモジュールのシャドウイングによる実行時例外を防いでください。
ステップ5:UDTFを作成する
DROP FUNCTION IF EXISTS py_split_lines(STRING);
DROP FUNCTION IF EXISTS py_extract_emails(STRING);
DROP FUNCTION IF EXISTS py_parse_json(STRING);
DROP FUNCTION IF EXISTS py_expand_json(STRING);
DROP FUNCTION IF EXISTS py_ngram(STRING, INT);
CREATE TABLES FUNCTION py_split_lines(STRING)
RETURNS ARRAY<STRING>
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/text_udtf.zip",
"symbol" = "text_udtf.split_lines_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE TABLES FUNCTION py_extract_emails(STRING)
RETURNS ARRAY<STRING>
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/text_udtf.zip",
"symbol" = "text_udtf.extract_emails_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE TABLES FUNCTION py_parse_json(STRING)
RETURNS ARRAY<STRUCT<k:STRING, v:STRING>>
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/text_udtf.zip",
"symbol" = "text_udtf.parse_json_object_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE TABLES FUNCTION py_expand_json(STRING)
RETURNS ARRAY<STRUCT<id:INT, name:STRING, score:DOUBLE>>
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/text_udtf.zip",
"symbol" = "text_udtf.expand_json_array_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
CREATE TABLES FUNCTION py_ngram(STRING, INT)
RETURNS ARRAY<STRING>
PROPERTIES (
"type" = "PYTHON_UDF",
"file" = "file:///path/to/text_udtf.zip",
"symbol" = "text_udtf.ngram_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
);
ステップ6: 関数を使用する
SELECT line
FROM (SELECT 'Line 1\nLine 2\nLine 3' as text) t
LATERAL VIEW py_split_lines(text) tmp AS line;
+--------+
| line |
+--------+
| Line 1 |
| Line 2 |
| Line 3 |
+--------+
SELECT email
FROM (SELECT 'Contact us at support@example.com or sales@company.org' as content) t
LATERAL VIEW py_extract_emails(content) tmp AS email;
+---------------------+
| email |
+---------------------+
| support@example.com |
| sales@company.org |
+---------------------+
SELECT k, v
FROM (SELECT '{"name": "Alice", "age": "25"}' as json_data) t
LATERAL VIEW py_parse_json(json_data) tmp AS k, v;
+------+-------+
| k | v |
+------+-------+
| name | Alice |
| age | 25 |
+------+-------+
SELECT id, name, score
FROM (
SELECT '[{"id": 1, "name": "Alice", "score": 95.5}, {"id": 2, "name": "Bob", "score": 88.0}]' as data
) t
LATERAL VIEW py_expand_json(data) tmp AS id, name, score;
+------+-------+-------+
| id | name | score |
+------+-------+-------+
| 1 | Alice | 95.5 |
| 2 | Bob | 88 |
+------+-------+-------+
SELECT ngram
FROM (SELECT 'Apache Doris is a fast database' as text) t
LATERAL VIEW py_ngram(text, 2) tmp AS ngram;
+---------------+
| ngram |
+---------------+
| Apache Doris |
| Doris is |
| is a |
| a fast |
| fast database |
+---------------+
Python UDTFの削除
-- Syntax
DROP FUNCTION IF EXISTS function_name(parameter_types);
-- Example
DROP FUNCTION IF EXISTS py_split(STRING, STRING);
DROP FUNCTION IF EXISTS py_range(INT, INT);
DROP FUNCTION IF EXISTS py_explode_json(STRING);
Python UDTFの変更
Dorisは既存の関数を直接変更することをサポートしていません。まず削除してから再作成してください:
DROP FUNCTION IF EXISTS py_split(STRING, STRING);
CREATE TABLES FUNCTION py_split(STRING, STRING) ...;
パラメータリファレンス
CREATE TABLES FUNCTION パラメータ
| パラメータ | 説明 |
|---|---|
function_name | 関数名。SQLの識別子命名規則に従います |
parameter_types | INT、STRING、DOUBLEなどのパラメータタイプのリスト |
RETURNS ARRAY<...> | 返される配列タイプ。出力構造を定義します • 単一カラム: ARRAY<type>• 複数カラム: ARRAY<STRUCT<col1:type1, col2:type2, ...>> |
PROPERTIES パラメータ
| パラメータ | 必須 | デフォルト | 説明 |
|---|---|---|---|
type | はい | - | 固定値 "PYTHON_UDF" |
symbol | はい | - | Python関数名。 • インラインモード: "split_string_udtf"のように関数名を直接記述• モジュールモード: [package_name.]module_name.function_name形式 |
file | いいえ | - | Python .zipパッケージへのパス。モジュールモードでのみ必須。3つのプロトコルをサポート:• file://: ローカルファイルシステムパス• http://: HTTPリモートダウンロード• https://: HTTPSリモートダウンロード |
runtime_version | はい | - | "3.10.12"などのPythonランタイムバージョン |
always_nullable | いいえ | true | 関数が常にnull許容結果を返すかどうか |
volatility | いいえ | immutable | 関数の変動性。4.1.2以降でサポート。有効な値はimmutable、stable、volatile。immutable: 同じ入力は常に同じ出力を生成します。ほとんどの決定論的UDFは、より良いプラン最適化のためにこの値を使用する必要があります。stable: 結果は1つのSQL文内では安定していますが、文間では変わる可能性があります。now()に似ています。SQLキャッシュとマテリアライズドビューの書き換えが無効になります。volatile: 結果は各呼び出しで変わる可能性があります。random()に似ています。SQLキャッシュ、マテリアライズドビューの書き換え、および多くのオプティマイザの書き換えルールが無効になります。 |
runtime_version 注意事項
- Pythonバージョンは
x.x.xまたはx.x.xx形式の完全なバージョン番号として指定する必要があります。 - Dorisは設定されたPython環境でこのバージョンに一致するインタープリターを検索します。
データタイプマッピング
Python UDTFはPython UDFと同じデータタイプマッピング規則を使用し、すべての整数、浮動小数点、文字列、日時、decimal、boolean、array、STRUCTタイプを含みます。
詳細なタイプマッピングについては: データタイプマッピングを参照してください。
NULL処理
- DorisはSQL
NULL値をPythonNoneにマッピングします。 - 関数内でパラメータが
Noneかどうかを確認してください。 yieldで生成される値はNoneを含む場合があり、そのカラムがNULLであることを示します。
実際の使用シナリオ
シナリオ1: CSVデータパース
DROP FUNCTION IF EXISTS py_parse_csv(STRING);
CREATE TABLES FUNCTION py_parse_csv(STRING)
RETURNS ARRAY<STRUCT<name:STRING, age:INT, city:STRING>>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "parse_csv_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def parse_csv_udtf(csv_data):
'''Parse multi-row CSV data'''
if csv_data is None:
return
lines = csv_data.strip().split('\n')
for line in lines:
parts = line.split(',')
if len(parts) >= 3:
name = parts[0].strip()
age = int(parts[1].strip()) if parts[1].strip().isdigit() else None
city = parts[2].strip()
yield (name, age, city)
$$;
SELECT name, age, city
FROM (
SELECT 'Alice,25,Beijing\nBob,30,Shanghai\nCharlie,28,Guangzhou' as data
) t
LATERAL VIEW py_parse_csv(data) tmp AS name, age, city;
+---------+------+-----------+
| name | age | city |
+---------+------+-----------+
| Alice | 25 | Beijing |
| Bob | 30 | Shanghai |
| Charlie | 28 | Guangzhou |
+---------+------+-----------+
シナリオ2:日付範囲の生成
DROP FUNCTION IF EXISTS py_date_range(STRING, STRING);
CREATE TABLES FUNCTION py_date_range(STRING, STRING)
RETURNS ARRAY<STRING>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "date_range_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
from datetime import datetime, timedelta
def date_range_udtf(start_date, end_date):
'''Generate a date range'''
if start_date is None or end_date is None:
return
try:
start = datetime.strptime(start_date, '%Y-%m-%d')
end = datetime.strptime(end_date, '%Y-%m-%d')
current = start
while current <= end:
yield (current.strftime('%Y-%m-%d'),)
current += timedelta(days=1)
except:
pass
$$;
SELECT date
FROM (SELECT '2024-01-01' as start_date, '2024-01-07' as end_date) t
LATERAL VIEW py_date_range(start_date, end_date) tmp AS date;
+------------+
| date |
+------------+
| 2024-01-01 |
| 2024-01-02 |
| 2024-01-03 |
| 2024-01-04 |
| 2024-01-05 |
| 2024-01-06 |
| 2024-01-07 |
+------------+
シナリオ 3: Text Tokenization
DROP FUNCTION IF EXISTS py_tokenize(STRING);
CREATE TABLES FUNCTION py_tokenize(STRING)
RETURNS ARRAY<STRUCT<word:STRING, position:INT>>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "tokenize_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
import re
def tokenize_udtf(text):
'''Tokenize text and output the words and their positions'''
if text is None:
return
# Use a regex to extract words
words = re.findall(r'\b\w+\b', text.lower())
for i, word in enumerate(words, 1):
if len(word) >= 2: # Filter out single characters
yield (word, i)
$$;
SELECT word, position
FROM (SELECT 'Apache Doris is a fast OLAP database' as text) t
LATERAL VIEW py_tokenize(text) tmp AS word, position;
+----------+----------+
| word | position |
+----------+----------+
| apache | 1 |
| doris | 2 |
| is | 3 |
| fast | 5 |
| olap | 6 |
| database | 7 |
+----------+----------+
シナリオ4: URLパラメータの解析
DROP FUNCTION IF EXISTS py_parse_url_params(STRING);
CREATE TABLES FUNCTION py_parse_url_params(STRING)
RETURNS ARRAY<STRUCT<param_name:STRING, param_value:STRING>>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "parse_url_params_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
from urllib.parse import urlparse, parse_qs
def parse_url_params_udtf(url):
'''Parse URL parameters'''
if url is None:
return
try:
parsed = urlparse(url)
params = parse_qs(parsed.query)
for key, values in params.items():
for value in values:
yield (key, value)
except:
pass
$$;
SELECT param_name, param_value
FROM (
SELECT 'https://example.com/page?id=123&category=tech&tag=python&tag=database' as url
) t
LATERAL VIEW py_parse_url_params(url) tmp AS param_name, param_value;
+------------+-------------+
| param_name | param_value |
+------------+-------------+
| id | 123 |
| category | tech |
| tag | python |
| tag | database |
+------------+-------------+
シナリオ5: IP範囲の拡張
DROP FUNCTION IF EXISTS py_expand_ip_range(STRING, STRING);
CREATE TABLES FUNCTION py_expand_ip_range(STRING, STRING)
RETURNS ARRAY<STRING>
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "expand_ip_range_udtf",
"runtime_version" = "3.10.12",
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def expand_ip_range_udtf(start_ip, end_ip):
'''Expand an IP address range (only the last octet is supported)'''
if start_ip is None or end_ip is None:
return
try:
# Assume the format is 192.168.1.10 to 192.168.1.20
start_parts = start_ip.split('.')
end_parts = end_ip.split('.')
if len(start_parts) == 4 and len(end_parts) == 4:
# Expand only the last octet
if start_parts[:3] == end_parts[:3]:
prefix = '.'.join(start_parts[:3])
start_num = int(start_parts[3])
end_num = int(end_parts[3])
for i in range(start_num, end_num + 1):
yield (f"{prefix}.{i}",)
except:
pass
$$;
SELECT ip
FROM (SELECT '192.168.1.10' as start_ip, '192.168.1.15' as end_ip) t
LATERAL VIEW py_expand_ip_range(start_ip, end_ip) tmp AS ip;
+--------------+
| ip |
+--------------+
| 192.168.1.10 |
| 192.168.1.11 |
| 192.168.1.12 |
| 192.168.1.13 |
| 192.168.1.14 |
| 192.168.1.15 |
+--------------+
パフォーマンス推奨事項
1. 出力行数を制御する
- 大量の出力行が生成される可能性があるシナリオでは、適切な上限を設定してください。
- 直積の爆発的増加を避けてください。
2. 重複計算を避ける
同じ計算結果を複数回使用する必要がある場合は、事前に計算してください:
# Not recommended
def bad_split_udtf(text):
for i in range(len(text.split(','))): # split is called every time
parts = text.split(',')
yield (parts[i],)
# Recommended
def good_split_udtf(text):
parts = text.split(',') # split only once
for part in parts:
yield (part,)
3. ジェネレータ式を使用する
中間リストの作成を避けるためにPythonジェネレータを活用してください:
# Not recommended
def bad_filter_udtf(text, delimiter):
parts = text.split(delimiter)
filtered = [p.strip() for p in parts if p.strip()] # Creates a list
for part in filtered:
yield (part,)
# Recommended
def good_filter_udtf(text, delimiter):
parts = text.split(delimiter)
for part in parts:
part = part.strip()
if part: # Filter directly
yield (part,)
4. 外部リソースへのアクセスを避ける
- UDTFではデータベース、ファイル、ネットワークにアクセスしないでください。
- すべての処理は入力パラメータに基づいて行う必要があります。
制限事項と注意点
1. ステートレス制限
- Python UDTFはステートレスです。各関数呼び出しは1行を独立して処理します。
- 呼び出し間で状態を保持することはできません。
- 行間集約には、UDAFを使用してください。
2. パフォーマンスに関する考慮事項
- Python UDTFのパフォーマンスは組み込みテーブル関数よりも低くなります。
- 複雑なロジックで中程度のデータ量のシナリオで使用してください。
- 大量のデータの場合は、最適化または組み込み関数を優先してください。
3. 固定出力タイプ
RETURNS ARRAY<...>で定義されたタイプは固定です。yieldで出力される値は定義と一致する必要があります。- 単一列:
yield valueまたはyield (value,)。複数列:yield (value1, value2, ...)。
4. 関数名
- 同じ関数名を異なるデータベースで定義できます。
- 曖昧さを避けるため、呼び出し時にデータベース名を指定してください。
5. 環境の一貫性
- すべてのBEノードのPython環境は一貫している必要があります。
- これには、Pythonバージョン、依存関係パッケージバージョン、環境設定が含まれます。
よくある質問
Q1:UDTFとUDFの違いは何ですか?
A:UDFは1行を入力として1行を出力する一対一の関係です。UDTFは1行を入力として0行以上を出力する一対多の関係です。
例:
SELECT py_upper(name) FROM users;
SELECT tag FROM users LATERAL VIEW py_split(tags, ',') tmp AS tag;
Q2: 複数の列を出力するにはどうすればよいですか?
A: 複数列出力の場合は、戻り値の型をSTRUCTで定義し、タプルをyieldします:
CREATE TABLES FUNCTION func(...)
RETURNS ARRAY<STRUCT<col1:INT, col2:STRING>>
...
def func(...):
yield (123, 'hello') # Corresponds to col1 and col2
Q3: UDTFが出力を生成しないのはなぜですか?
A: 考えられる理由:
yieldが呼び出されていない: 関数がyieldを呼び出していることを確認してください。- フィルタリング: すべてのデータがフィルタされています。
- 例外が隠されている: try-exceptブロックがエラーを隠していないかチェックしてください。
- NULL入力: 入力がNULLで、関数が直接returnしています。
Q4: UDTFは状態を維持できますか?
A: いいえ。Python UDTFはステートレスで、各関数の呼び出しは1つの行を独立して処理します。行をまたいだ集約や状態の維持には、Python UDAFを使用してください。
Q5: UDTFの出力行数を制限するにはどうすればよいですか?
A: 関数内にカウンターまたは条件チェックを追加してください:
def limited_udtf(data):
max_rows = 1000
count = 0
for item in data.split(','):
if count >= max_rows:
break
yield (item,)
count += 1
Q6: UDTFから出力されるデータ型に制限はありますか?
A: UDTFは、プリミティブ型(INT、STRING、DOUBLEなど)と複合型(ARRAY、STRUCT、MAPなど)を含む、すべてのDorisデータ型をサポートしています。出力型はRETURNS ARRAY<...>で明示的に定義する必要があります。
Q7: UDTFで外部リソースにアクセスできますか?
A: 技術的には可能ですが、強く非推奨です。UDTFは純粋に機能的であるべきで、入力パラメータのみを処理する必要があります。外部リソース(データベース、ファイル、ネットワーク)へのアクセスは、パフォーマンスの問題と予測不可能な動作を引き起こします。
Python UDF/UDAF/UDTF環境設定とマルチバージョン管理
Python環境管理
Python UDF/UDAF/UDTFを使用する前に、DorisのBackend(BE)ノードでPythonランタイム環境が正しく設定されていることを確認してください。DorisはCondaまたは**Virtual Environment(venv)**でのPython環境管理をサポートしており、異なるUDFが異なるバージョンのPythonインタープリターと依存関係を使用できます。
DorisはPython環境を管理する2つの方法を提供しています:
- Condaモード: Miniconda/Anacondaでマルチバージョン環境を管理します。
- Venvモード: 組み込みのPython仮想環境(venv)でマルチバージョン環境を管理します。
サードパーティライブラリのインストールと使用
Python UDF、UDAF、UDTFはすべてサードパーティライブラリを使用できます。Dorisは分散システムであるため、すべてのBEノードでサードパーティライブラリを統一してインストールする必要があります。そうしないと、一部のノードで実行に失敗します。
インストール手順
-
各BEノードで依存関係をインストール:
# Install with pip
pip install numpy pandas requests
# Or install with conda
conda install numpy pandas requests -y -
関数内でそれらをインポートして使用する:
import numpy as np
import pandas as pd
# Use them in a UDF/UDAF/UDTF function
def my_function(x):
return np.sqrt(x)
注意事項
pandasとpyarrowは必須の依存関係です。すべてのPython環境に事前にインストールしてください。そうでなければPython UDF/UDAF/UDTFが実行できません。- すべてのBEノードで同じバージョンの依存関係をインストールしてください。そうでなければ一部のノードで実行に失敗します。
- インストールパスは、対応するUDF/UDAF/UDTFが使用するPythonランタイム環境と一致する必要があります。
- 仮想環境またはConda環境を使用して依存関係を管理し、システムのPython環境との競合を避けてください。
BE設定パラメータ
すべてのBEノードのbe.conf設定ファイルで以下のパラメータを設定し、設定を有効にするためにBEを再起動してください。
設定パラメータリファレンス
| パラメータ | 型 | 許可される値 | デフォルト | 説明 |
|---|---|---|---|---|
enable_python_udf_support | bool | true / false | false | Python UDF機能を有効にするかどうか |
python_env_mode | string | conda / venv | "" | Pythonマルチバージョン環境管理モード |
python_conda_root_path | string | ディレクトリパス | "" | Minicondaのルートディレクトリpython_env_mode = condaの場合のみ有効 |
python_venv_root_path | string | ディレクトリパス | ${DORIS_HOME}/lib/udf/python | venvマルチバージョン管理のルートディレクトリpython_env_mode = venvの場合のみ有効 |
python_venv_interpreter_paths | string | パスリスト(:で区切り) | "" | 利用可能なPythonインタープリターディレクトリのリストpython_env_mode = venvの場合のみ有効 |
max_python_process_num | int32 | 整数 | 0 | Python Serverプロセスプールの最大プロセス数0はCPUコア数をデフォルトとして使用することを意味します。別の正の整数を設定してデフォルトを上書きできます |
方法1: CondaでPython環境を管理
1. BEを設定
be.confに以下の設定を追加してください:
## be.conf
enable_python_udf_support = true
python_env_mode = conda
python_conda_root_path = /path/to/miniconda3
2. Environment 検索ルール
DorisはUDFで指定されたruntime_versionに一致する${python_conda_root_path}/envs/以下のConda環境を検索します。
マッチングルール:
runtime_versionは完全なPythonバージョン番号である必要があります。形式はx.x.xまたはx.x.xxで、"3.9.18"や"3.12.11"などです。- Dorisは全てのConda環境を反復処理し、各環境内の実際のPythonインタープリターバージョンが
runtime_versionと完全に一致するかどうかをチェックします。 - 一致する環境が見つからない場合、Dorisは
Python environment with version x.x.x not foundというエラーを報告します。
例:
- UDFが
runtime_version = "3.9.18"を指定した場合、DorisはPythonバージョンが3.9.18の環境を検索します。 - 環境名は何でも構いません(
py39、my-env、data-scienceなど)。その環境のPythonバージョンが3.9.18であれば問題ありません。 - 完全なバージョン番号が必要です。
"3.9"や"3.12"などのバージョンプレフィックスは許可されません。
3. ディレクトリ構造図
## File system layout on a Doris BE node (Conda mode)
/path/to/miniconda3 ← python_conda_root_path (configured in be.conf)
│
├── bin/
│ ├── conda ← conda CLI (used for operations)
│ └── ... ← Other conda tools
│
├── envs/ ← Directory for all Conda environments
│ │
│ ├── py39/ ← Conda environment 1 (user-created)
│ │ ├── bin/
│ │ │ ├── python ← Python 3.9 interpreter (called directly by Doris)
│ │ │ ├── pip
│ │ │ └── ...
│ │ ├── lib/
│ │ │ └── python3.9/
│ │ │ └── site-packages/ ← Third-party dependencies for this environment (such as pandas, pyarrow)
│ │ └── ...
│ │
│ ├── py312/ ← Conda environment 2 (user-created)
│ │ ├── bin/
│ │ │ └── python ← Python 3.12 interpreter
│ │ └── lib/
│ │ └── python3.12/
│ │ └── site-packages/ ← Pre-installed dependencies (such as torch, sklearn)
│ │
│ └── ml-env/ ← Semantic environment name (recommended)
│ ├── bin/
│ │ └── python ← May be Python 3.12 with GPU dependencies
│ └── lib/
│ └── python3.12/
│ └── site-packages/
│
└── ...
4. Conda環境の作成
DorisのPython UDF/UDAF/UDTF機能はpandasとpyarrowに必須の依存関係があります。すべてのPython環境で両方のライブラリを必ず事前にインストールする必要があります。そうしないとUDFが正しく動作しません。
Python環境を作成するために、すべてのBEノードで以下のコマンドを実行してください:
# Install Miniconda (when not yet installed)
wget https://repo.anaconda.com/miniconda/Miniconda3-latest-Linux-x86_64.sh
bash Miniconda3-latest-Linux-x86_64.sh -b -p /opt/miniconda3
# Create a Python 3.9.18 environment and install required dependencies (the environment name can be customized)
/opt/miniconda3/bin/conda create -n py39 python=3.9.18 pandas pyarrow -y
# Create a Python 3.12.11 environment and pre-install dependencies (Important: the Python version must be specified exactly, and pandas and pyarrow must be installed)
/opt/miniconda3/bin/conda create -n py312 python=3.12.11 pandas pyarrow numpy -y
# Activate an environment and install additional dependencies
source /opt/miniconda3/bin/activate py39
conda install requests beautifulsoup4 -y
conda deactivate
# Verify the Python version in the environment
/opt/miniconda3/envs/py39/bin/python --version # Should output: Python 3.9.18
/opt/miniconda3/envs/py312/bin/python --version # Should output: Python 3.12.11
5. UDFでの使用
-- Use the Python 3.12.11 environment
CREATE FUNCTION py_ml_predict(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.12.11", -- Must specify the complete version number to match Python 3.12.11
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(x):
# Libraries installed in the Python 3.12.11 environment can be used
return x * 2
$$;
-- Note: Whether the environment is named py312 or ml-env, any environment whose Python version is 3.12.11 can be used
-- runtime_version cares only about the Python version, not the environment name
Method 2: Venvを使用したPython環境の管理
1. BEの設定
be.confに以下の設定を追加してください:
## be.conf
enable_python_udf_support = true
python_env_mode = venv
python_venv_root_path = /doris/python_envs
python_venv_interpreter_paths = /opt/python3.9/bin/python3.9:/opt/python3.12/bin/python3.12
2. 設定パラメータに関する注意事項
python_venv_root_path: 仮想環境のルートディレクトリです。すべてのvenv環境がこのディレクトリ下に作成されます。python_venv_interpreter_paths: コロン(:)で区切られたPythonインタープリターの絶対パスのリストです。Dorisはそれぞれのインタープリターのバージョンをチェックし、UDFで指定されたruntime_version("3.9.18"などの完全なバージョン番号)と照合します。
3. ディレクトリ構造図
## Doris BE configuration (be.conf)
python_venv_interpreter_paths = "/opt/python3.9/bin/python3.9:/opt/python3.12/bin/python3.12"
python_venv_root_path = /doris/python_envs
/opt/python3.9/bin/python3.9 ← System pre-installed Python 3.9
/opt/python3.12/bin/python3.12 ← System pre-installed Python 3.12
/doris/python_envs/ ← Root directory for all virtual environments (python_venv_root_path)
│
├── python3.9.18/ ← Environment ID = complete Python version
│ ├── bin/
│ │ ├── python
│ │ └── pip
│ └── lib/python3.9/site-packages/
│ ├── pandas==2.1.0
│ └── pyarrow==15.0.0
│
├── python3.12.11/ ← Python 3.12.11 environment
│ ├── bin/
│ │ ├── python
│ │ └── pip
│ └── lib/python3.12/site-packages/
│ ├── pandas==2.1.0
│ └── pyarrow==15.0.0
│
└── python3.12.10/ ← Python 3.12.10 environment
└── ...
4. Venv環境の作成
DorisのPython UDF/UDAF/UDTF機能はpandasとpyarrowに必須の依存関係があります。すべてのPython環境に両方のライブラリを事前にインストールする必要があり、そうしないとUDFが正しく動作しません。
すべてのBEノードで以下のコマンドを実行してください:
# Create the root directory for virtual environments
mkdir -p /doris/python_envs
# Create a virtual environment with Python 3.9
/opt/python3.9/bin/python3.9 -m venv /doris/python_envs/python3.9.18
# Activate the environment and install the required dependencies (pandas and pyarrow are required)
source /doris/python_envs/python3.9.18/bin/activate
pip install pandas pyarrow numpy
deactivate
# Create a virtual environment with Python 3.12
/opt/python3.12/bin/python3.12 -m venv /doris/python_envs/python3.12.11
# Activate the environment and install the required dependencies (pandas and pyarrow are required)
source /doris/python_envs/python3.12.11/bin/activate
pip install pandas pyarrow numpy scikit-learn
deactivate
5. UDFでの使用
-- Use the Python 3.9.18 environment
CREATE FUNCTION py_clean_text(STRING)
RETURNS STRING
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.9.18", -- Must specify the complete version number to match Python 3.9.18
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
def evaluate(text):
return text.strip().upper()
$$;
-- Use the Python 3.12.11 environment
CREATE FUNCTION py_calculate(DOUBLE)
RETURNS DOUBLE
PROPERTIES (
"type" = "PYTHON_UDF",
"symbol" = "evaluate",
"runtime_version" = "3.12.11", -- Must specify the complete version number to match Python 3.12.11
"always_nullable" = "true",
"volatility" = "immutable"
)
AS $$
import numpy as np
def evaluate(x):
return np.sqrt(x)
$$;
環境管理のベストプラクティス
1. 適切な管理方法の選択
| シナリオ | 推奨 | 理由 |
|---|---|---|
| Pythonバージョンの頻繁な切り替え | Conda | 強力な環境分離とシンプルな依存関係管理 |
| 既存のConda環境 | Conda | 既存の環境を直接再利用可能 |
| 限られたシステムリソース | Venv | より小さなフットプリントと高速な起動 |
| 既存のPythonシステム環境 | Venv | Condaを別途インストールする必要がない |
2. 環境一貫性要件
すべてのBEノードのPython環境は完全に同一である必要があります。以下を含みます:
- Pythonバージョンが同じである必要があります。
- インストールされた依存関係パッケージとそのバージョンが同じである必要があります。
- 環境ディレクトリパスが同じである必要があります。
注意事項
1. 設定変更の有効化
be.confを変更した後、変更を有効にするにはBEプロセスを再起動する必要があります。- サービス中断を避けるため、再起動前に設定が正しいことを確認してください。
2. パス検証
設定前に、パスが正しいことを確認してください:
# Conda mode: verify the conda path
ls -la /opt/miniconda3/bin/conda
/opt/miniconda3/bin/conda env list
# Venv mode: verify interpreter paths
/opt/python3.9/bin/python3.9 --version
/opt/python3.12/bin/python3.12 --version
3. 権限設定
Doris BEプロセスがPython環境ディレクトリにアクセスする権限を持っていることを確認してください:
# Conda mode
chmod -R 755 /opt/miniconda3
# Venv mode
chmod -R 755 /doris/python_envs
chown -R doris:doris /doris/python_envs # Assume the BE process user is doris
4. Resource Limits
実際のニーズに応じてPythonプロセスプールパラメータを調整してください:
## Use the CPU core count (recommended, max_python_process_num = 0)
max_python_process_num = 0
## High concurrency: specify the process count manually
max_python_process_num = 128
## Resource-constrained: limit the process count
max_python_process_num = 32
環境の確認
各BEノードでの環境確認
# Conda mode
/opt/miniconda3/envs/py39/bin/python --version
/opt/miniconda3/envs/py39/bin/python -c "import pandas; print(pandas.__version__)"
# Venv mode
/doris/python_envs/python3.9.18/bin/python --version
/doris/python_envs/python3.9.18/bin/python -c "import pandas; print(pandas.__version__)"
全てのBEノードに共通するPythonバージョンを表示
SHOW PYTHON VERSIONS;
+---------+---------+---------+-------------------+----------------------------------------+
| Version | EnvName | EnvType | BasePath | ExecutablePath |
+---------+---------+---------+-------------------+----------------------------------------+
| 3.9.18 | py39 | conda | path/to/miniconda | path/to/miniconda/envs/py39/bin/python |
+---------+---------+---------+-------------------+----------------------------------------+
指定されたバージョンの インストール済み依存関係を表示
SHOW PYTHON PACKAGES IN '<version>' を使用して、指定されたバージョンのインストール済み依存関係を表示します。BEノード間で依存関係が異なる場合、異なる部分が一覧表示されます。
SHOW PYTHON PACKAGES IN '3.9.18'
すべてのBEノードが同一の依存関係を持つ場合:
+-----------------+-------------+
| Package | Version |
+-----------------+-------------+
| pyarrow | 21.0.0 |
| Bottleneck | 1.4.2 |
| jieba | 0.42.1 |
| six | 1.17.0 |
| wheel | 0.45.1 |
| python-dateutil | 2.9.0.post0 |
| tzdata | 2025.3 |
| setuptools | 80.9.0 |
| numpy | 2.0.1 |
| psutil | 7.0.0 |
| pandas | 2.3.3 |
| mkl_random | 1.2.8 |
| pip | 25.3 |
| snownlp | 0.12.3 |
| pytz | 2025.2 |
| mkl_fft | 1.3.11 |
| mkl-service | 2.4.0 |
| numexpr | 2.10.1 |
+-----------------+-------------+
BE ノードが異なる依存関係を持つ場合:
+-----------------+-------------+------------+----------------+
| Package | Version | Consistent | Backends |
+-----------------+-------------+------------+----------------+
| pyarrow | 21.0.0 | Yes | |
| Bottleneck | 1.4.2 | Yes | |
| six | 1.17.0 | Yes | |
| jieba | 0.42.1 | No | 127.0.0.1:9660 |
| wheel | 0.45.1 | Yes | |
| python-dateutil | 2.9.0.post0 | Yes | |
| tzdata | 2025.3 | Yes | |
| setuptools | 80.9.0 | Yes | |
| numpy | 2.0.1 | Yes | |
| psutil | 7.0.0 | No | 127.0.0.1:9660 |
| pandas | 2.3.3 | Yes | |
| mkl_random | 1.2.8 | Yes | |
| pip | 26.0.1 | No | 127.0.0.1:9077 |
| pip | 25.3 | No | 127.0.0.1:9660 |
| snownlp | 0.12.3 | No | 127.0.0.1:9660 |
| pytz | 2025.2 | Yes | |
| numexpr | 2.10.1 | Yes | |
| mkl-service | 2.4.0 | Yes | |
| mkl_fft | 1.3.11 | Yes | |
+-----------------+-------------+------------+----------------+
一般的なトラブルシューティング
Q1: UDF呼び出しで「Python environment not found」と報告される
原因:
runtime_versionで指定されたバージョンがシステム上に存在しない。- 環境パスが正しく設定されていない。
解決策:
# Check the Conda environment list
conda env list
# Check whether the venv interpreter exists
ls -la /opt/python3.9/bin/python3.9
# Check the BE configuration
grep python /path/to/be.conf
Q2: UDF呼び出し時に「ModuleNotFoundError: No module named 'xxx'」が報告される
原因: 必要な依存関係パッケージがPython環境にインストールされていません。
Q3: 異なるBEノードが異なる結果を返す
原因: BEノード間でPython環境または依存関係のバージョンが異なります。
解決方法:
- すべてのノードでPythonと依存関係のバージョンを確認してください。
- すべてのノードで環境の一貫性を確認してください。
requirements.txt(pip)またはenvironment.yml(Conda)を使用して環境を統一的にデプロイしてください。一般的な使用例:
-
requirements.txt(pip)を使用する場合:# Export dependencies in the development environment
pip freeze > requirements.txt
# Install dependencies on a BE node using the target Python
/path/to/python -m pip install -r requirements.txt -
environment.yml(Conda) を使用する場合:# Export dependencies
conda env export --from-history -n py312 -f environment.yml
# Create the environment on a BE node
conda env create -f environment.yml -n py312
# Or update an existing environment
conda env update -f environment.yml -n py312
- 依存関係ファイルに
pandasとpyarrowが含まれており、すべてのBEノードに同じバージョンがインストールされていることを確認してください。 - インストール時は、Doris設定に一致するPythonインタープリターまたはCondaパス(
/opt/miniconda3/bin/condaや指定されたvenvインタープリターなど)を使用してください。 - 依存関係ファイルをバージョン管理下または共有ストレージに配置し、運用チームがすべてのBEノードに均一に配布できるようにしてください。
- 参考資料:pip official documentation、Conda environment export/import guide
Q4: be.confへの変更が反映されない
考えられる原因:BEプロセスが再起動されていません。
使用制限
-
パフォーマンスに関する考慮事項:
- Python UDFのパフォーマンスは組み込み関数より低くなります。複雑なロジックと小さなデータ量のシナリオで使用してください。
- 大量のデータに対しては、ベクトル化モードを優先してください。
-
型の制限:
- HLLやBitmapなどの特殊な型はサポートされていません。
-
環境の分離:
- 同じ関数名を異なるデータベースで定義できます。
- 呼び出し時にデータベース名を指定(
db.func()など)して曖昧さを避けてください。
-
同時実行の制限:
- Python UDFはプロセスプールを介して実行されます。同時実行数は
max_python_process_numによって制限されます。 - 高い同時実行が必要なシナリオでは、このパラメーターを増やしてください。
- Python UDFはプロセスプールを介して実行されます。同時実行数は