Liu Liu created FLINK-40472:
-------------------------------
Summary: Add Arrow-native vectorized UDF support to DataFrame API
Key: FLINK-40472
URL: https://issues.apache.org/jira/browse/FLINK-40472
Project: Flink
Issue Type: Sub-task
Reporter: Liu Liu
Support the following usage:
{code:java}
import pyarrow as pa
import pyarrow.compute as pcfrom pyflink.dataframe import DataType, col, udf
@udf(return_dtype=DataType.string(), batch_size=1024)
def normalize_name(names: pa.Array) -> pa.Array:
# Operates directly on an Arrow array without converting to pandas.
return pc.utf8_upper(names)
result = df.with_column(
"normalized_name",
normalize_name(col("name")),
) {code}
--
This message was sent by Atlassian Jira
(v8.20.10#820010)