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文では、完全なバージョン番号("3.10.12"など)でruntime_versionを明示的に指定する必要があります。メジャーバージョンとマイナーバージョンのみ("3.10"など)を指定することはできません。そうしないと関数呼び出しが失敗します。
Python UDF(スカラー関数)
Python UDF(User Defined Function)は行ごとにデータを処理します。関数は行ごとに1回呼び出され、単一の結果を返します。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パラメータで参照します。
Step 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値 |
| 整数型 | 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 | バリアント型(文字列に変換) | |
DATE | str | 'YYYY-MM-DD'形式の日付文字列 | |
DATETIME | str | 'YYYY-MM-DD HH:MM:SS'形式の日時文字列 | |
| 日付/時間型 | DATEV2 | datetime.date | 日付オブジェクト |
DATETIMEV2 | datetime.datetime | 日時オブジェクト | |
TIMESTAMPTZ | datetime.datetime | タイムゾーン付き日時オブジェクト | |
| Decimal型 | DECIMAL / DECIMALV2 | decimal.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 | ビットマップデータ(まだサポートされていません) |
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
Vectorized Mode
Vectorized modeはPandasを使用してデータをバッチで処理し、scalar modeを上回るパフォーマンスを発揮します。vectorized modeでは、関数の引数はpandas.Seriesオブジェクトであり、戻り値もpandas.Seriesである必要があります。
システムがvectorized modeを認識するように、関数シグネチャで型アノテーション(a: pd.Seriesなど)を使用し、関数内でバッチデータ構造を直接操作してください。明示的なvectorized型がない場合、システムはscalar modeにフォールバックします。
関数シグネチャでpd.Series型と通常の型が混在する場合、システムは通常の型パラメータに対応する入力列を定数列として扱い(同じ値がバッチ全体で再利用される)、期待と一致しない結果が生成される可能性があります。vectorized modeでは、パラメータスタイルを一貫させてください。すべてのパラメータにpandas.Series型アノテーションを使用するか、すべてに通常の型パラメータを使用する(scalar mode)かのいずれかにしてください。
## 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 Modeを優先する
Vectorized modeはscalar modeよりも大幅に優れたパフォーマンスを発揮します:
# 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: Average集約
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 module を書く
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クラス名。 • Inlineモード: "SumUDAF"のようにクラス名を直接記述• Moduleモード: 形式は [package_name.]module_name.ClassName |
file | いいえ | - | Python .zipパッケージへのパス。moduleモードでのみ必須。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環境でこのバージョンに一致するインタープリターを検索します。
Window関数
Python UDAFをwindow関数(OVER句)と組み合わせることができます:
Python UDAFをwindow関数(OVER句)で使用する場合、Dorisは各windowフレームの計算後に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値をPythonNoneにマップします。 accumulateメソッドでは、パラメータがNoneかどうかをチェックしてください。- 集約関数は
Noneを返すことで、結果がNULLであることを示すことができます。
実際のシナリオ
シナリオ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: Moving Average
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接続、スレッドロック、および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はwindow関数で使用できますか?
A: はい。Python UDAFはwindow関数(OVER句)を完全にサポートします。
Q4: mergeメソッドはいつ呼び出されますか?
A: mergeは以下の状況で呼び出されます:
- 分散集約: 異なるBEノードからの部分的な集約結果をマージします。
- 並列処理: 同一ノード上の異なるスレッドからの部分結果をマージします。
- Window関数: windowフレーム内の部分結果をマージします。
したがって、mergeの実装は正確でなければなりません。そうでなければ結果が間違ったものになります。
Python UDTF(テーブル関数)
Python UDTF(User Defined Table Function)を使用すると、単一の行を複数の出力行に変換するカスタムテーブル関数を定義できます。データの分割、拡張、生成に役立ちます。
Python UDTFの主要特性:
- 1行から多行へ: 単一の行を入力として受け取り、ゼロ、1つ、または多数の出力行を生成します。
- 柔軟な出力構造: 任意の数と型の出力列を許可し、単純型と複雑なSTRUCT型をサポートします。
- Lateral viewサポート: データ拡張と結合のために
LATERAL VIEWと連携します。 - 関数型スタイル: Python関数と
yield文を使用し、簡潔で直感的です。
UDTFの基本概念
テーブル関数の実行方法
Python UDTFは関数(クラスではない)として実装されます。実行フローは以下の通りです:
- 入力の受信: 関数は単一行の列値をパラメータとして受け取ります。
- 処理とyield:
yield文がゼロ個以上の出力行を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を実行します。
インラインモード
インラインモードではSQL内で直接Python関数を記述できます。単純なテーブル関数ロジックに適しています。
構文:
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アーカイブとしてパッケージ化し、関数作成時に参照します。
ステップ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 | 関数が常にnullable結果を返すかどうか |
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は、すべての整数、浮動小数点、文字列、日時、decimal、boolean、array、STRUCT型を含む、Python UDFと同じデータ型マッピングルールを使用します。
詳細な型マッピングについては: データ型マッピングを参照してください。
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はステートレスです。各関数呼び出しは一つの行を独立して処理します。
- 呼び出し間で状態を保持することはできません。
- 行をまたいだ集約には、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は一行を受け取り一行を出力する一対一の関係です。UDTFは一行を受け取りゼロ行以上を出力する一対多の関係です。
例:
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ランタイム環境と一致する必要があります。
- 依存関係を管理し、システムのPython環境との競合を避けるために、仮想環境またはConda環境を使用してください。
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. 環境検索ルール
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
方法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に必須の依存関係があります。UDFが正しく動作するために、すべてのPython環境に両方のライブラリを事前にインストールする必要があります。
すべての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. リソース制限
実際のニーズに応じて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環境が見つかりません」と報告される
原因:
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はプロセスプール経由で実行されます。並行性は