Skip to content

Latest commit

 

History

History
314 lines (241 loc) · 11.7 KB

File metadata and controls

314 lines (241 loc) · 11.7 KB
displayed_sidebar docs
sidebar_position 0.91
description Since v3.4.0, StarRocks supports creating Python UDFs.

import Experimental from '../../_assets/commonMarkdown/_experimental.mdx'

Python UDF

This topic describes how to develop User Defined Functions (UDF) using Python.

Since v3.4.0, StarRocks supports creating Python UDFs.

Currently, StarRocks only supports scalar UDFs in Python.

Prerequisites

Make sure the following requirements are met before proceeding:

  • Install Python 3.8 or later.
  • Enable UDF in StarRocks by setting enable_udf to true in the FE configuration file fe/conf/fe.conf, and then restart the FE nodes to allow the configuration to take effect. For more information, see FE configuration - enable_udf.
  • Python UDF requires the following packages: pyarrow
  • Set the location of the Python interpreter environment in the BE instance using environment variable. Add the variable item python_envs and set it to the location of the Python interpreter installation, for example, /opt/Python-3.8/.

Develop and use Python UDFs

Syntax:

CREATE [GLOBAL] FUNCTION function_name(arg_type [, ...])
RETURNS return_type
{PROPERTIES ("key" = "value" [, ...]) | key="value" [...] }
[AS $$ $$]
Parameter Required Description
GLOBAL No Whether to create a global UDF.
function_name Yes The name of the function you want to create. You can include the name of the database in this parameter, for example,db1.my_func. If function_name includes the database name, the UDF is created in that database. Otherwise, the UDF is created in the current database. The name of the new function and its parameters cannot be the same as an existing name in the destination database. Otherwise, the function cannot be created. The creation succeeds if the function name is the same but the parameters are different.
arg_type Yes Argument type of the function. The added argument can be represented by , .... For the supported data types, see Mapping between SQL data types and Python data types.
return_type Yes The return type of the function. For the supported data types, see Mapping between SQL data types and Python data types.
PROPERTIES Yes Properties of the function, which vary depending on the type of the UDF to create.
AS $$ $$ No Specify the inline UDF code between $$ marks.

Properties include:

Property Required Description
type Yes The type of the UDF. Setting it to Python indicates to create a Python-based UDF.
symbol Yes The name of the class for the Python project to which the UDF belongs. The value of this parameter is in the <package_name>.<class_name> format.
input No Type of input. Valid values: scalar (Default) and arrow (vectorized input).
file No The HTTP URL from which you can download the Python package file that contains the code for the UDF. The value of this parameter is in the http://<http_server_ip>:<http_server_port>/<python_package_name> format. Note the package name must have the suffix .py.zip. Default value: inline, which indicates to create an inline UDF.
service_url No The URL of an external Arrow Flight worker service to run the UDF on, for example, grpc+tcp://<host>:<port>. When set, the BE connects to this worker service instead of spawning a local Python worker, and you are responsible for running and isolating that worker service yourself. Leave it unset (default) to use the built-in local worker. See Use an external worker service.

Create an inline scalar input UDF using Python

The following example creates an inline echo function with scalar input using Python.

CREATE FUNCTION python_echo(INT)
RETURNS INT
type = 'Python'
symbol = 'echo'
file = 'inline'
AS
$$
def echo(x):
    return x
$$
;

Create an inline vectorized input UDF using Python

Vectorized input is supported to enhance the UDF processing performance.

The following example creates an inline add function with vectorized input using Python.

CREATE FUNCTION python_add(INT) 
RETURNS INT
type = 'Python'
symbol = 'add'
input = "arrow"
AS
$$
import pyarrow.compute as pc
def add(x):
    return pc.add(x, 1)
$$
;

Create a packaged UDF using Python

When creating the Python package, you must package the module into .py.zip files, which need to meet the zipimport format.

> tree .
.
├── main.py
└── yaml
    ├── composer.py
    ├── constructor.py
    ├── cyaml.py
    ├── dumper.py
    ├── emitter.py
    ├── error.py
    ├── events.py
    ├── __init__.py
    ├── loader.py
    ├── nodes.py
    ├── parser.py
> cat main.py 
import numpy
import yaml

def echo(a):
    return yaml.__version__
CREATE FUNCTION py_pack(string) 
RETURNS  string 
symbol = "add"
type = "Python"
file = "http://HTTP_IP:HTTP_PORT/m1.py.zip"
symbol = "main.echo"
;

Use an external worker service

By default, the BE spawns and manages a local Python worker process to run a Python UDF. Alternatively, you can run the worker as an external Arrow Flight service yourself and have the BE connect to it by setting the service_url property. This lets you run and isolate the worker independently of the BE (for example, in a separate, sandboxed container).

When service_url is set:

  • The BE connects to the worker at that URL instead of spawning a local worker.
  • For inline UDFs, the function code is sent to the worker over the connection.
  • For packaged (file) UDFs, the worker downloads the package from the file URL itself and verifies it against the function's checksum, so the worker service must be able to reach that URL.
  • You are responsible for running, scaling, securing, and isolating the worker service. The connection is plaintext by default; restrict network access to the worker and enable TLS/authentication as needed.
  • To keep a query from hanging when a worker is unreachable or stuck, set the BE configuration item python_udf_rpc_timeout_ms (milliseconds) to bound the worker call. It applies to both local and external workers and defaults to 0 (no timeout). Because it bounds the whole call stream, set it above the longest expected UDF query.

The example below creates an inline UDF that runs on an external worker service:

CREATE FUNCTION py_echo(string)
RETURNS string
type = "Python"
symbol = "echo"
service_url = "grpc+tcp://192.168.1.10:8815"
AS
$$
def echo(x):
    return x
$$
;

Appendix

Mapping between SQL data types and Python data types

SQL Type Python 3 Type
SCALAR
TINYINT/SMALLINT/INT/BIGINT/LARGEINT INT
STRING string
DOUBLE FLOAT
BOOLEAN BOOL
DATETIME DATETIME.DATETIME
FLOAT FLOAT
CHAR STRING
VARCHAR STRING
DATE DATETIME.DATE
DECIMAL DECIMAL.DECIMAL
ARRAY list
MAP dict
STRUCT dict (keyed by field name)
JSON str (use json.loads to parse)
VECTORIZED
TYPE_BOOLEAN pyarrow.lib.BoolArray
TYPE_TINYINT pyarrow.lib.Int8Array
TYPE_SMALLINT pyarrow.lib.Int15Array
TYPE_INT pyarrow.lib.Int32Array
TYPE_BIGINT pyarrow.lib.Int64Array
TYPE_FLOAT pyarrow.FloatArray
TYPE_DOUBLE pyarrow.DoubleArray
VARCHAR pyarrow.StringArray
DECIMAL pyarrow.Decimal128Array
DATE pyarrow.Date32Array
TYPE_TIME pyarrow.TimeArray
ARRAY pyarrow.ListArray
MAP pyarrow.MapArray
STRUCT pyarrow.StructArray

Nested types (ARRAY, MAP, STRUCT) may be composed arbitrarily — for example, array<map<string,int>>, map<string,array<int>>, struct<a array<int>, b map<string,int>>, or array<struct<...>>. In scalar mode each row is delivered as plain Python list / dict values recursively. To return a nested type, return matching Python list / dict values (with dict keys equal to struct field names); pyarrow performs the recursive conversion automatically.

Compile Python

Follow these steps to compile Python:

  1. Get the OpenSSL package.

    wget 'https://github.com/openssl/openssl/archive/OpenSSL_1_1_1m.tar.gz'
  2. Decompress the package.

    tar -zxf OpenSSL_1_1_1m.tar.gz
  3. Navigate to the decompressed folder.

    cd openssl-OpenSSL_1_1_1m
  4. Set the environment variable OPENSSL_DIR.

    export OPENSSL_DIR=`pwd`/install
  5. Prepare source code for compilation.

    ./Configure --prefix=`pwd`/install
    ./config --prefix=`pwd`/install
  6. Compile OpenSSL.

    make -j 16 && make install
  7. Set the environment variable LD_LIBRARY_PATH.

    LD_LIBRARY_PATH=$OPENSSL_DIR/lib:$LD_LIBRARY_PATH
  8. Navigate back to the working directory and get the Python package.

    wget 'https://www.python.org/ftp/python/3.12.9/Python-3.12.9.tgz'
  9. Decompress the package.

    tar -zxf ./Python-3.12.9.tgz 
  10. Navigate to the decompressed folder.

cd Python-3.12.9
  1. Make and navigate to the directory build.
mkdir build && cd build
  1. Prepare source code for compilation.
../configure --prefix=`pwd`/install --with-openssl=$OPENSSL_DIR
  1. Compile Python.
make -j 16 && make install
  1. Install PyArrow and grpcio.
./install/bin/pip3 install pyarrow grpcio
  1. Compress the files into the package.
tar -zcf ./Python-3.12.9.tar.gz install
  1. Distribute the package to the target BE server, and decompress the package.
tar -zxf ./Python-3.12.9.tar.gz
  1. Modify the BE configuration file be.conf to add the following configuration item.
python_envs=/home/disk1/sr/install