All Products
Search
Document Center

Realtime Compute for Apache Flink:Python UDTFs

Last Updated:Sep 23, 2026

This topic describes how to develop, register, and use Python user-defined table-valued functions (UDTFs).

Definition

A user-defined table-valued function (UDTF) takes zero, one, or more scalar values as input parameters, which can be variable-length. Unlike a scalar function, a UDTF can return any number of rows instead of a single value. Each returned row can contain one or more columns, and each call to a UDTF can output multiple rows or columns.

Limits

Due to the deployment and network environments of Realtime Compute for Apache Flink, take note of the following limits when you develop Python user-defined functions:

  • Only open source Flink V1.12 and later are supported.

  • Python is pre-installed in the Flink workspace. Develop your code in the pre-installed Python version.

Note

Realtime Compute engine VVR earlier than 8.0.11 has Python 3.7.9 pre-installed, and VVR 8.0.11 and later has Python 3.9.21 pre-installed. If you upgrade an earlier version to VVR 8.0.11 or later, you must retest, deploy, and run the PyFlink jobs that were developed in the earlier version.

  • The Flink runtime supports only JDK 8 and JDK 11. If a Python job depends on a third-party JAR package, make sure that the JAR package is compatible.

  • Only open source Scala V2.11 is supported. If a Python job depends on a third-party JAR package, use the JAR package that corresponds to Scala V2.11.

  • Inline functions are supported only in VVR 11.9.0 and later.

Develop a UDTF

You can develop a Python function in one of the following ways:

  • Package and upload the Python code, and then register the code on the platform.

  • Declare the code logic as an inline function in an SQL script.

If the UDTF logic is simple, you can develop the Python function as an inline function.

Develop a UDTF by using a code package

Develop a UDTF

Note

Flink provides a sample project of Python user-defined extensions (UDXs) to help you develop UDXs. The sample includes the implementations of Python user-defined scalar functions, user-defined aggregate functions (UDAFs), and UDTFs. This example uses the Windows operating system to describe how to develop a UDTF.

Important

  • A registered UDF name can contain only lowercase letters, digits, and hyphens (-). Underscores (\_) and other special characters are not supported.

  • The function name cannot be the same as that of a built-in function. The managed environment of Realtime Compute for Apache Flink has some built-in functions pre-installed, such as split. If a custom function has the same name as a built-in function, the registration fails or SQL validation reports an error. Use a custom prefix for the function name, such as my-split.

  1. Download and decompress the python\_demo-master sample package to your computer.

  2. In PyCharm, choose file > open to open the decompressed python\_demo-master folder.

  3. Double-click \python\_demo-master\udx\udtfs.py and modify the content of the udtfs.py file based on your business requirements.

In this example, my\_split splits the string of each row into multiple columns by vertical bars (|).


    from pyflink.table import DataTypes
    from pyflink.table.udf import udtf
    
    @udtf(result_types=[DataTypes.STRING(), DataTypes.STRING()])
    def my_split(s: str):
        splits = s.split("|")
        yield splits[0], splits[1]
        
  1. In the directory that contains the udx folder in the downloaded package (\python\_demo), run the following command to package the files.


    zip -r python_demo.zip udx
        

The python\_demo.zip package is generated in the \python\_demo\ directory, indicating that the UDTF is developed.

Register a UDTF

For more information about how to register a UDTF, see Manage UDFs.

Use a UDTF

After the UDTF is registered, you can use it. Perform the following steps:

  1. Develop a Flink SQL job. For more information, see Job development overview.

The value of the message field in each row of the ASI\_UDTF\_Source table is concatenated with the string aa by a vertical bar (|), and then split into multiple columns by vertical bars (|). The following code provides an example:


    CREATE TEMPORARY TABLE ASI_UDTF_Source (
      `message`  VARCHAR
    ) WITH (
      'connector'='datagen'
    );
    
    CREATE TEMPORARY TABLE ASI_UDTF_Sink (
      name  VARCHAR,
      place  VARCHAR
    ) WITH (
      'connector' = 'blackhole'
    );
    
    INSERT INTO ASI_UDTF_Sink
    SELECT name,place
    FROM ASI_UDTF_Source,lateral table(my_split(concat_ws('|', `message`, 'aa'))) as T(name,place);
        
  1. On the O&M > Deployments page, click Start in the Actions column of the target job.

After the job starts, the ASI\_UDTF\_Sink table contains two columns. The data is generated by concatenating the message and aa fields of each row in the ASI\_UDTF\_Source table with a vertical bar (|) and then splitting the result by vertical bars (|).

Python inline table functions

An inline function embeds the function implementation directly in a CREATE FUNCTION statement. The function definition and the registration declaration are completed in the same SQL statement. The following example defines an inline function that masks email addresses. You must declare the complete code logic between $$.

CREATE TEMPORARY FUNCTION parse_order_items(items_json STRING)
RETURNS TABLE (
    sku STRING,
    quantity INT
)
HANDLER 'ParseOrderItems'
AS $$
import json

class ParseOrderItems:
    def eval(self, items_json):
        if not items_json:
            return

        try:
            items = json.loads(items_json)
        except (TypeError, ValueError):
            return

        if not isinstance(items, list):
            return

        for item in items:
            if not isinstance(item, dict):
                continue

            sku = str(item.get("sku", "")).strip().upper()

            try:
                quantity = int(item.get("quantity"))
            except (TypeError, ValueError):
                continue

            if sku and quantity > 0:
                yield sku, quantity
$$
LANGUAGE PYTHON;
Note
  1. Python is sensitive to indentation. Write the first level of the function body at the beginning of a line, and make sure that the relative indentation inside the code block is correct.

  2. Inline functions do not support vectorized execution and cannot be declared as nondeterministic functions.

  3. Only TEMPORARY functions can be created.