Skip to main content

RisingWave Python UDF SDK

This library provides a Python SDK for creating user-defined functions (UDF) in RisingWave.

For a detailed guide on how to use Python UDF in RisingWave, please refer to this doc.

Introduction

RisingWave supports user-defined functions implemented as external functions. With the RisingWave Python UDF SDK, users can define custom UDFs using Python and start a Python process as a UDF server. RisingWave can then remotely access the UDF server to execute the defined functions.

Installation

pip install risingwave

Usage

Define functions in a Python file:

# udf.py
from risingwave.udf import udf, udtf, UdfServer
import struct
import socket

# Define a scalar function
@udf(input_types=['INT', 'INT'], result_type='INT')
def gcd(x, y):
    while y != 0:
        (x, y) = (y, x % y)
    return x

# Define a scalar function that returns multiple values (within a struct)
@udf(input_types=['BYTEA'], result_type='STRUCT<VARCHAR, VARCHAR, SMALLINT, SMALLINT>')
def extract_tcp_info(tcp_packet: bytes):
    src_addr, dst_addr = struct.unpack('!4s4s', tcp_packet[12:20])
    src_port, dst_port = struct.unpack('!HH', tcp_packet[20:24])
    src_addr = socket.inet_ntoa(src_addr)
    dst_addr = socket.inet_ntoa(dst_addr)
    return src_addr, dst_addr, src_port, dst_port

# Define a table function
@udtf(input_types='INT', result_types='INT')
def series(n):
    for i in range(n):
        yield i

# Start a UDF server
if __name__ == '__main__':
    server = UdfServer(location="0.0.0.0:8815")
    server.add_function(gcd)
    server.add_function(series)
    server.serve()

Start the UDF server:

python3 udf.py

To create functions in RisingWave, use the following syntax:

create function <name> ( <arg_type>[, ...] )
    [ returns <ret_type> | returns table ( <column_name> <column_type> [, ...] ) ]
    as <name_defined_in_server> using link '<udf_server_address>';
  • The as parameter specifies the function name defined in the UDF server.
  • The link parameter specifies the address of the UDF server.

For example:

create function gcd(int, int) returns int
as gcd using link 'http://localhost:8815';

create function series(int) returns table (x int)
as series using link 'http://localhost:8815';

select gcd(25, 15);

select * from series(10);

Data Types

The RisingWave Python UDF SDK supports the following data types:

SQL Type Python Type Notes
BOOLEAN bool
SMALLINT int
INT int
BIGINT int
REAL float
DOUBLE PRECISION float
DECIMAL decimal.Decimal
DATE datetime.date
TIME datetime.time
TIMESTAMP datetime.datetime
INTERVAL MonthDayNano / (int, int, int) Fields can be obtained by months(), days() and nanoseconds() from MonthDayNano
VARCHAR str
BYTEA bytes
JSONB any
T[] list[T]
STRUCT<> tuple
...others Not supported yet.

Metadata

Release files for risingwave 0.1.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for risingwave 0.1.1
File Size Uploaded
risingwave-0.1.1.tar.gz 10.5 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for risingwave 0.1.1
File Interpreter ABI Platform
risingwave-0.1.1-py3-none-any.whl Python 3 none any Details

Total release size: 21.6 kB

Release files / risingwave-0.1.1.tar.gz

Download URL risingwave-0.1.1.tar.gz
Size 10.5 kB
Tags Source
SHA-256 checksum
How to use checksums
ddfb031c3582852f077c36f99dcd540f5fa4b73e44f950c0d926bdb59795095a
BLAKE2b-256 checksum
How to use checksums
e0aa25dce50fde98c1973380bf58bff18d206fa3ece5a0e6dcdfdceedc1df267
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.11.6

Release files / risingwave-0.1.1-py3-none-any.whl

Download URL risingwave-0.1.1-py3-none-any.whl
Size 11.2 kB
Tags Python 3
SHA-256 checksum
How to use checksums
236b58f2cc8cb5525baec6e6710d1ce9aedad0212b00d1d4dce275dae2ddd379
BLAKE2b-256 checksum
How to use checksums
587a31a1f80960031357f0f0a90f234c547e77e0076493eeb294dffa5b7dda2b
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.11.6

Release history Release notifications | RSS feed

This release

0.1.1 This release

2 release files

0.1.0

2 release files

0.0.11

2 release files

0.0.10

2 release files

0.0.9

1 release file

0.0.8

1 release file

0.0.7

1 release file

0.0.6

1 release file

0.0.5

1 release file

0.0.4

1 release file

0.0.3

1 release file

0.0.2

1 release file

0.0.1

1 release file

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page