探索DuckDB:使用标量 Python UDF 快速扩展 DuckDB 的功能
原文:From Waddle to Flying: Quickly expanding DuckDB's functionality with Scalar Python UDFs[1]
翻译:Gemini
校对:alitrack
作者:Pedro Holanda、Thijs Bruineman 和 Phillip Cloud
日期:2023-07-07
重点提示:DuckDB 现在支持矢量化标量 Python 用户自定义函数 (UDF)。通过实现 Python UDF,用户可以轻松扩展 DuckDB 的功能,同时利用 DuckDB 的快速执行模型、SQL 和数据安全性。
用户自定义函数 (UDF) 让用户能够方便扩展数据库管理系统 (DBMS) 的功能,以执行未作为内置函数实现的特定于域的任务。例如,经常需要导出私有数据的用户可以受益于匿名化函数,该函数在保留域的同时屏蔽电子邮件的本地部分。理想情况下,此函数将直接在 DBMS 中执行。这种方法提供了几个优点:
1. 性能。 该函数可以使用 DBMS 的相同执行模型(例如,流结果、内存外/核心外执行)执行,并且没有任何不必要的转换。
2. 易于使用。 UDF 可以无缝集成到 SQL 查询中,允许用户利用 SQL 的强大功能来调用函数。这消除了通过单独的数据库连接器传递数据并执行外部代码的需要。这些函数可以用于各种 SQL 上下文(例如,子查询、连接条件)。
3. 安全性。 敏感数据绝不会离开 DBMS 进程。
用户经常避免实现 UDF 的主要原因有两个,
1. UDF 存在安全问题。由于 UDF 是由用户创建并在 DBMS 进程中执行的自定义代码,因此存在使服务器崩溃的潜在风险。但是,对于嵌入式数据库 DuckDB,这种担忧得到了缓解,因为每个分析师都单独运行自己的 DuckDB 进程。因此,对服务器稳定性的影响并不是一个值得担心的重大问题。
2. 实现的难度是用户常见的威慑因素。高性能 UDF 通常仅在低级语言中受支持。像 Python 这样的高级语言中的 UDF 会产生巨大的性能成本。因此,许多用户无法快速实现其 UDF,而无需花费大量时间学习低级语言和了解 DBMS 的内部细节。
DuckDB 遵循类似的方法。作为专为分析任务定制的 DBMS,性能是一个关键考虑因素,导致其核心在 C++ 中实现。因此,可扩展性工作的最初重点集中在 C++ 上[2]。然而,这只鸭子不仅限于蹒跚学步;它还能飞。因此,我们很高兴地宣布最近添加[3]了标量 Python UDF 到 DuckDB。
DuckDB 提供对两种不同类型的 Python UDF 的支持,它们在用于DuckDB 的本机数据类型[4]和 Python 进程之间的通信的 Python 对象中有所不同。这些通信层包括对Python 内置类型[5]和PyArrow 表[6]的支持。
这两种方法表现出两个主要区别:
1. 零拷贝。 PyArrow 表利用我们与 Arrow 的零拷贝集成[7],能够以零拷贝成本将数据类型有效地转换为 Python 语言。
2. 矢量化。 PyArrow 表函数在块级别上运行,处理包含多达 2048 行的数据块。这种方法最大限度地提高了缓存局部性并利用了矢量化。另一方面,内置类型 UDF 实现按行操作。
本博客文章旨在演示如何使用 Python UDF 扩展 DuckDB,特别强调由 PyArrow 驱动的 UDF。在我们的快速游览部分,我们将提供使用 PyArrow UDF 类型的示例。对于有兴趣了解基准测试的人,可以跳到下面的基准测试部分。如果您想查看 Python UDF API 的详细说明,请参阅我们的文档[8]。
Python UDF
本节描述了使用 Python UDF 的几个实际示例。每个示例都使用不同类型的 Python UDF。
快速游览
为了演示在 DuckDB 中使用 Python UDF,我们考虑以下示例。我们有一个名为 world_cup_titles 的字典,它将国家/地区映射到他们赢得的世界杯数量。我们想要创建一个 Python UDF,它接受国家/地区名称作为输入,在字典中搜索相应的值,并返回该国家/地区赢得的世界杯数量。如果字典中找不到该国家/地区,则 UDF 将返回 NULL。
这是一个示例实现:
import duckdb
from duckdb.typing import *con = duckdb.connect()
# 字典,将国家/地区和他们赢得的世界杯数映射起来
world_cup_titles = {
"Brazil": 5,
"Germany": 4,
"Italy": 4,
"Argentina": 2,
"Uruguay": 2,
"France": 2,
"England": 1,
"Spain": 1
}
# 将注册为 UDF 的函数,简单地在 Python 字典中进行查找
def world_cups(x):
return world_cup_titles.get(x)
# 我们注册该函数
con.create_function("wc_titles", world_cups, [VARCHAR], INTEGER)
就是这样,然后注册该函数并准备好通过 SQL 调用。
# 让我们创建一个示例国家/地区表,其中包含我们感兴趣的国家/地区
con.execute("CREATE TABLE countries(country VARCHAR)")
con.execute("INSERT INTO countries VALUES ('Brazil'), ('Germany'), ('Italy'), ('Argentina'), ('Uruguay'), ('France'), ('England'), ('Spain'), ('Netherlands')")
# 我们可以通过 SQL 简单地调用该函数,甚至可以使用函数返回来消除从未赢得过世界杯的国家/地区
con.sql("SELECT country, wc_titles(country) as world_cups from countries").fetchall()
# [('Brazil', 5), ('Germany', 4), ('Italy', 4), ('Argentina', 2), ('Uruguay', 2), ('France', 2), ('England', 1), ('Spain', 1), ('Netherlands', None)]
使用 Faker 生成虚假数据(内置类型 UDF)
以下示例演示了在 DuckDB 中使用Faker 库[9]生成标量函数,该函数返回随机生成的日期。名为 random_date 的函数不需要任何输入,并输出一个 DATE 列。由于 Faker 利用内置 Python 类型,因此函数直接返回它们。需要注意的一件重要事情是,一个函数不能根据其输入确定性地必须标记为具有 side_effects。
import duckdb# 通过导入 duckdb.typing,我们可以直接指定 DuckDB 类型,而无需使用字符串
from duckdb.typing import *
from faker import Faker
# 我们的 Python UDF 每次调用都会生成一个随机日期
def random_date():
fake = Faker()
return fake.date_between()
然后,我们必须使用 create_function 在 DuckDB 中注册 Python 函数。由于我们的函数不需要任何输入,因此我们可以将空列表作为 argument_type_list 传递。由于函数返回日期,因此我们指定 duckdb.typing 中的 DATE 作为 return_type。请注意,由于我们的 random_date() 函数返回内置 Python 类型 (datetime.date),因此我们无需指定 UDF 类型。
# 为了举例说明副作用的影响,让我们首先在不标记的情况下运行该函数。
duckdb.create_function('random_date', random_date, [], DATE)# 注册后,我们可以直接通过 SQL 使用该函数
# 请注意,如果没有 side_effect=True,则无法保证重新评估该函数。
res = duckdb.sql('select random_date() from range (3)').fetchall()
# [(datetime.date(2003, 8, 3),), (datetime.date(2003, 8, 3),), (datetime.date(2003, 8, 3),)]
# 现在让我们重新添加标记为 true 的副作用的函数。
duckdb.remove_function('random_date')
duckdb.create_function('random_date', random_date, [], DATE, side_effects=True)
res = duckdb.sql('select random_date() from range (3)').fetchall()
# [(datetime.date(2020, 11, 29),), (datetime.date(2009, 5, 18),), (datetime.date(2018, 5, 24),)]
交换字符串大小写(PyArrow 类型 UDF)
使用内置类型的一个问题是您无法从零拷贝、矢量化和缓存 为了演示 PyArrow 函数,我们考虑一个简单的示例,我们希望将小写字符转换为大写字符,并将大写字符转换为小写字符。幸运的是,PyArrow 在计算引擎中已经为此提供了一个函数,它就像调用 pc.utf8_swapcase(x) 一样简单。
import duckdb# 通过导入 duckdb.typing,我们可以直接指定 DuckDB 类型,而无需使用字符串
from duckdb.typing import *
import pyarrow as pa
import pyarrow.compute as pc
def swap_case(x):
# 使用 utf8_swapcase 交换 'column' 的大小写并返回结果
return pc.utf8_swapcase(x)
con = duckdb.connect()
# 要注册函数,我们必须将其类型定义为 'arrow'
con.create_function('swap_case', swap_case, [VARCHAR], VARCHAR, type='arrow')
res = con.sql("select swap_case('PEDRO HOLANDA')").fetchall()
# [('pedro holanda',)]
预测出租车费用(Ibis + PyArrow UDF)
Python UDF 提供了强大的功能,因为它们使用户能够利用广泛的 Python 生态系统和工具,包括有效实现机器学习操作的库,如 PyTorch[10] 和 Tensorflow[11]。
此外,Ibis 项目[12] 提供了一个与 DuckDB 集成良好的 DataFrame API,并支持 DuckDB 的原生 Python 和 PyArrow UDF。
在此示例中,我们演示了如何使用预构建的 PyTorch 模型来估计基于行驶距离的出租车费用。您可以在 Ibis 团队的这篇博文中[13] 找到一个完整的示例。
import torch
import pyarrow as pa
import ibis
import ibis.expr.datatypes as dtfrom ibis.expr.operations import udf
# 此代码段中未指定生成模型的代码,请参阅提供的链接以获取更多信息
model = ...
# 函数使用模型和行驶距离输入张量来预测值,请参阅提供的链接以获取更多信息
def predict_linear_regression(model, tensor: torch.Tensor) -> torch.Tensor:
...
# 向 ibis 指示这是一个标量用户定义函数,其输入格式为 pyarrow
@udf.scalar.pyarrow
def predict_fare(x: dt.float64) -> dt.float32:
# `x` 是一个 pyarrow.ChunkedArray;`dt.float64` 注释指示了 ChunkedArray 的元素类型。
# 将数据从 PyArrow 转换为所需的 torch 张量格式和维度。
tensor = torch.from_numpy(x.to_numpy()[:, None]).float()
# 调用实际的预测函数,它也返回一个 torch 张量。
predicted = predict_linear_regression(model, tensor).ravel()
return pa.array(predicted.numpy())
# 在 NYC Taxi parquet 文件上执行查询以展示我们模型的预测、实际费用金额和距离。
expr = (
ibis.read_parquet('yellow_tripdata_2016-02.parquet')
.mutate(
"fare_amount",
"trip_distance",
predicted_fare=lambda t: predict_fare(t.trip_distance),
)
)
df = expr.execute()
通过在 DuckDB 中使用 Ibis 中的 Python UDF,您可以无缝地合并机器学习模型并在 Ibis 代码和 SQL 查询中直接执行预测。该示例演示了如何使用 PyTorch 模型根据距离预测出租车费用,展示了在 Ibis 驱动的 DuckDB 的 SQL 环境中集成机器学习功能。
基准
在本节中,我们将执行简单的基准比较,以演示两种不同类型的 Python UDF 之间的性能差异。基准将测量执行时间和峰值内存消耗。基准执行 5 次,并考虑中间值。基准是在具有 16GB RAM 的 Mac Apple M1 上进行的。
内置 Python 与 PyArrow
为了对这些 UDF 类型进行基准测试,我们创建了将整数列作为输入、将一添加到每个值并返回结果的 UDF。可以在 此处[14] 找到用于此基准测试部分的代码。
import pyarrow.compute as pc
import duckdb
import pyarrow as pa# 内置 UDF
def add_built_in_type(x):
return x + 1
# Arrow UDF
def add_arrow_type(x):
return pc.add(x,1)
con = duckdb.connect()
# 注册
con.create_function('built_in_types', add_built_in_type, ['BIGINT'], 'BIGINT', type='native')
con.create_function('add_arrow_type', add_arrow_type, ['BIGINT'], 'BIGINT', type='arrow')
# 具有 10,000,000 个元素的整数视图。
con.sql("""
select i
from range(10000000) tbl(i);
""").to_view("numbers")
# 调用两个 UDF
native_res = con.sql("select sum(add_built_in_type(i)) from numbers").fetchall()
arrow_res = con.sql("select sum(add_arrow_type(i)) from numbers").fetchall()
| 名称 | 时间 (秒) |
| 内置 | 5.37 |
| PyArrow | 0.35 |
我们可以在两个 UDF 之间观察到超过一个数量级的性能差异。性能差异主要归因于三个因素:
1. 在 Python 中,对象构造和一般用途相当慢。这是由于几个原因,包括自动内存管理、解释和动态类型。
2. PyArrow UDF 不需要任何数据复制。
3. PyArrow UDF 以矢量化方式执行,处理数据块而不是单独的行。
Python UDF 与外部函数
这里我们比较了 Python UDF 与外部函数的使用。在这种情况下,我们有一个函数来计算列中所有字符串长度的总和。您可以在 此处[15] 找到用于此基准测试部分的代码。
import duckdb
import pyarrow as pa# UDF 中使用的函数
def string_length_arrow(x):
tuples = len(x)
values = [len(i.as_py()) if i.as_py() != None else 0 for i in x]
array = pa.array(values, type=pa.int32(), size=tuples)
return array
# 与数据库外部相同的函数
def exec_external(con):
arrow_table = con.sql("select i from strings tbl(i)").arrow()
arrow_column = arrow_table['i']
tuples = len(arrow_column)
values = [len(i.as_py()) if i.as_py() != None else 0 for i in arrow_column]
array = pa.array(values, type=pa.int32(), size=tuples)
arrow_tbl = pa.Table.from_arrays([array], names=['i'])
return con.sql("select sum(i) from arrow_tbl").fetchall()
con = duckdb.connect()
con.create_function('strlen_arrow', string_length_arrow, ['VARCHAR'], int, type='arrow')
con.sql("""
select
case when i != 0 and i % 42 = 0
then
NULL
else
repeat(chr((65 + (i % 26))::INTEGER), (4 + (i % 12))) end
from range(10000000) tbl(i);
""").to_view("strings")
con.sql("select sum(strlen_arrow(i)) from strings tbl(i)").fetchall()
exec_external(con)
| 名称 | 时间 (秒) | 峰值内存消耗 (MB) |
| 外部 | 5.65 | 584.032 |
| UDF | 5.63 | 112.848 |
这里我们可以看到,在使用 UDF 时,性能没有明显的下降。但是,您仍然可以获得更安全的执行和 SQL 的利用。在我们的示例中,我们还可以注意到,与 UDF 方法相比,外部函数使整个查询具体化,导致峰值内存消耗高出 5 倍。
结论和进一步发展
DuckDB 现支持标量 Python UDF,这标志着在扩展数据库功能方面的一个重要里程碑。这一增强功能使用户能够使用高级语言执行复杂的计算。此外,Python UDF 可以利用 DuckDB 与 Arrow 的零拷贝集成,消除数据传输成本并确保查询高效执行。
虽然引入 Python UDF 是向前迈出的重要一步,但我们在该领域的工作仍在进行中。我们的路线图包括以下重点领域:
1. 聚合/生成表的 UDF:目前,用户可以创建标量 UDF,但我们正在积极努力支持聚合函数(对一组值执行计算并返回单个结果)和生成表的函数(无需限制列数和行数即可返回表)。
2. 类型:标量 Python UDF 目前支持大多数 DuckDB 类型,但 ENUM 类型和 BIT 类型除外。我们正在努力扩展类型支持以确保全面的功能。
引用链接
[1] From Waddle to Flying: Quickly expanding DuckDB's functionality with Scalar Python UDFs: https://duckdb.org/2023/07/07/python-udf.html[2] 集中在 C++ 上: https://www.youtube.com/watch?v=UKo_LQyLTko&ab_channel=DuckDBLabs[3] 最近添加: https://github.com/duckdb/duckdb/pull/7171[4] DuckDB 的本机数据类型: https://duckdb.org/docs/sql/data_types/overview[5] Python 内置类型: https://duckdb.org/docs/sql/data_types/overview[6] PyArrow 表: https://arrow.apache.org/docs/python/generated/pyarrow.Table.html[7] 零拷贝集成: https://duckdb.org/2021/12/03/duck-arrow.html[8] 文档: https://duckdb.org/docs/api/python/function[9] Faker 库: https://faker.readthedocs.io/en/master/[10] PyTorch: https://pytorch.org/[11] Tensorflow: https://www.tensorflow.org/[12] Ibis 项目: https://ibis-project.org/[13] Ibis 团队的这篇博文中: https://ibis-project.org/blog/rendered/torch/[14] 此处: https://gist.github.com/pdet/ebd201475581756c29e4533a8fa4106e[15] 此处: https://gist.github.com/pdet/2907290725539d390df7981e799ed593