Пользовательские функции на базе Python в Clickhouse
3 технологии, которыми мы воспользуемся: Docker, Clickhouse и Python
У OLAP - БД Clickhouse есть множество полезных встроенных функций, что в сочетании с высокой скоростью обработки запросов делает ее незаменимым и высокоэффективным помощником при работе с большими объемами данных.
Однако, как и в случае со всеми запросами SQL, иногда бывает достаточно сложно или даже невозможно закодировать определенную логику приложения. К счастью, у Clickhouse есть проверенное решение данной проблемы. Речь идет о пользовательских функциях (UDF).
UDF могут быть реализованы либо в виде лямбда-выражений либо путем вызова исполняемых файлов. В этой статье мы поговорим о втором способе, который открывает перед нами практически неограниченные возможности.
На самом деле, нет никаких правил относительно того, на каком языке программирования должны быть написаны UDF. Единственное требование - UDF должны быть определены в рамках одного исполняемого скрипта! Импорт пакетов разрешен.
Настройка UDF включает в себя три этапа:
- Управление одной строкой
- Добавление X минут к дате
- Массивы дат и чисел внутри, вложенные массивы строк снаружи!
Весь код, указанный в этой статье, включая код генерации рисунков, можно найти здесь.
Обязательные условия
ОС Linux OS (в моем случае это Ubuntu 20.04) и, конечно, CH.
Для получения CH нужно… установить Docker. Затем запустить контейнер Docker с образом CH.
Теперь у Вас должен быть доступ к веб-интерфейсу для выполнения запросов к Вашей БД по адресу http://localhost:8123/play
Запускаем контейнер:
sudo docker ps -a # check the container exists sudo docker container exec -it /<container ID or name> bash
Убедитесь в том, что доступна команда ‘nano’ и установите Python:
apt update apt install nano apt install python3 && python3 --version # check version
Простой пример
Начнем с определения функции под названием 'revString', которая переворачивает строку.
Наша функция будет иметь один входной аргумент, который имеет тип данных CH String. В процессе передачи данных CH в Python происходят следующие преобразования типов данных:
- строка > str
- любой числовой тип > float
- datetime > str
Поэтому в нашем случае нам не нужно беспокоиться о входящих типах данных, у нас есть строка!
Во-первых, определим xml-файл, чтобы предоставить CH инструкции по:
- названию UDF
- аргументам пользовательской функции (и их CH-типы данных)
- CH-типу данных для вывода* (может быть только один вывод)
- python-файлу для вызова, фактически выполняемого нашей функцией
*возвращаемый объект должен быть float, string или array. Строки Python отображаются на CH Strings, а плавающие числа Python - на CH Float64s. Время и дату в Python возвращать нельзя, хотя мы можем возвращать строковые представления дат в CH, а затем для их сопоставления со временем использовать функции parseDateTime или toDateTime.
xml-файл нужно сохранить в каталоге /etc/clickhouse-server/, его название должно заканчиваться на '_function.xml'. Файл должен быть исполняемым:
cd /etc/clickhouse-server/ # xml file directory for UDFs touch string_reverse_function.xml # <xyz>_function.xml chmod +x string_reverse_function.xml # make executable nano string_reverse_function.xml # edit, may require sudo nano
Заселяем файл xml содержимым:
<functions>
<function>
<type>executable</type>
<name>revString</name>
<return_type>String</return_type>
<argument>
<type>String</type>
</argument>
<format>TabSeparated</format>
<command>revString.py</command>
</function>
</functions>
Атрибут name определяет имя функции на CH. Эта функция будет принимать один строковый аргумент, а наша python-функция должна возвращать строку.
Теперь давайте создадим файл python, который должен находиться в каталоге /var/lib/clickhouse/user_scripts/. Единственное требование к имени файла - оно должно совпадать с атрибутом command в xml-файле (отражать атрибут name зеркально, как это сделали мы, не обязательно).
cd /var/lib/clickhouse/user_scripts/ touch <your_file>.py # e.g. revString.py chmod +x <your_file>.py nano <your_file>.py
Теперь содержимое…
#!/usr/bin/python3
import sys
if __name__ == '__main__':
for line in sys.stdin:
print(line.rstrip()[::-1], end='')
sys.stdout.flush()
Шебанг Python сообщает CH, где найти интерпретатор Python, необходимый для чтения этого модуля. Входной аргумент UDF передается Python в sys.stdin.
К сожалению, sys.stdin добавляет в конец строк символ новой строки, поэтому чтобы удалить его, мы вызываем метод rstrip (). Затем мы разворачиваем строку и печатаем ее уже без этого символа. sys.stdout вернет напечатанную строку обратно в терминал CH, завершая тем самым работу функции.
Чтобы проверить, что функция была создана, выполните приведенные ниже запросы (например, через веб-интерфейс). В запросе Select Вы должны увидеть результат для функции с именем 'revString'
-- system reload functions -- not usually needed select * from system.functions where origin = 'ExecutableUserDefined'
Пробуем…
Semordnilaps (полулапы) - это слова, которые при перевертывании образуют разные слова, а palindromes (палиндромы) – полные близнецы.
Пример с несколькими аргументами
Теперь давайте создадим функцию, которая к входному времени даты добавляет X минут или часов. (Поскольку мы не можем возвращать даты, нам нужно будет сначала возвращать строку, а затем разбирать ее обратно на даты в CH).
<functions>
<function>
<type>executable</type>
<name>add_x_time</name>
<return_type>String</return_type>
<return_name>result</return_name>
<argument>
<type>Float64</type>
<name>number</name>
</argument>
<argument>
<type>String</type>
<name>time_unit</name>
</argument>
<argument>
<type>DateTime64</type>
<name>date</name>
</argument>
<format>JSONEachRow</format>
<command>add_time.py</command>
</function>
</functions>
Поскольку для ввода данных мы используем формат JSONEachRow, мы должны дать каждому аргументу имя (ключи JSON). Нам также нужно будет добавить нашу выходную строку в JSON, где выходная строка должна быть значением ключа, названного по атрибуту return_name.
Удобно, что при использовании JSON нам не нужно удалять символы новой строки; мы извлекаем только небольшие части каждой строки, а не всю строку целиком.
#!/usr/bin/python3
import sys
import json
from datetime import datetime, timedelta
def add_time(x:float, unit:str, date:datetime) -> str:
valid_units = ['minutes', 'hours']
try:
assert unit in valid_units
except:
return f"Unimplemented time unit '{unit}'.\nValid units are: {valid_units}"
if unit == 'minutes':
new_date = date + timedelta(minutes=x) else:
new_date = date + timedelta(hours=x) return new_date.strftime("%Y-%m-%d %H:%M:%S")
if __name__ == '__main__':
for line in sys.stdin:
# convert input to json string dict, even floats will be str
value = json.loads(line) num = float(value['number']) # cast str as float
unit = value['time_unit'] # no cast needed as already str
date = datetime.strptime(value['date'], "%Y-%m-%d %H:%M:%S")
# json.dumps casts all data back to strings
result = json.dumps({'result': add_time(num, unit, date)})
print(result, end='\n') # json must end with a new line char
sys.stdout.flush()
-- trying the function on the web UI
select now() as in, toTypeName(now()) as type_in
, add_x_time(3, 'minutes', now()) as plus_3_mins
, toTypeName(plus_3_mins) as type_plus_3_mins
, add_x_time(-1, 'hours', now()) as minus_1_hour
, toTypeName(minus_1_hour) as type_minus_1_hour
, add_x_time(2, 'days', now()) as invalid_unit
, toTypeName(invalid_unit) as type_invalid_unit
, parseDateTime(add_x_time(3, 'minutes', now()), '%Y-%m-%d %H:%i:%s') + INTERVAL 2 DAYS as new_dt
, toTypeName(new_dt) as type_new_dt
Пример с вложенным массивом
Теперь определим более сложную функцию, которая в качестве аргументов использует массивы, подразумевает импорт нестандартной библиотеки и возвращает вложенные массивы. Затем с помощью специального запроса ClickHouse мы раскроем вложенные массивы.
Предположим, у нас есть данные о ЧСС, и мы хотим определить, на какое время приходятся самые высокие показатели.
Модуль UDF, который это делает, Вы найдете здесь. Результат работы функции показан ниже (зеленые столбики). На самом деле данная функция не сильно отличается от предыдущего примера, поскольку она тоже принимает JSON, извлекает атрибуты, применяет некоторую логику и возвращает JSON. Небольшое отличие заключается в том, что теперь результат JSON возвращается в виде списка списков.
Для запуска функции нам нужно установить все используемые ею пакеты нестандартной библиотеки. Сделаем это с помощью обычных команд терминала pip.
apt install pip pip3 install numpy pip3 install scipy # etc.
Объект, фактически возвращаемый пользовательской функцией, - это:
[['2024-05-12 05:15:00', '2024-05-12 07:25:00'], ['2024-05-12 11:25:00', '2024-05-12 12:05:00'], ['2024-05-12 16:35:00', '2024-05-12 16:55:00'], ['2024-05-12 18:45:00', '2024-05-12 19:50:00']]
Каждый период активности представлен элементом в списке, где элемент - это 2 даты, представляющие время начала и окончания активности и приведенные в виде строк.
А как же XML-файл? Обратите внимание на то, что return_type задает вложенный массив строк.
<functions>
<function>
<type>executable</type>
<name>activity</name>
<return_type>Array(Array(String))</return_type>
<return_name>result</return_name>
<argument>
<type>Array(DateTime64)</type>
<name>readingDates</name>
</argument>
<argument>
<type>Array(Float32)</type>
<name>values</name>
</argument>
<format>JSONEachRow</format>
<command>findActivity.py</command>
</function>
</functions>
Теперь запустим функцию на ClickHouse, где разложим вложенный массив с помощью arrayJoin и получим сего одну строку для каждого периода активности.
with [...] as dates, [...] as vals, -- see GitHub for expansion
base as (
select arrayMap(x -> toDateTime(x), dates) as d
, arrayMap(x -> toFloat32(x), vals) as v
, activity(d, v) as a
) select arrayMap(x -> formatDateTime(toDateTime(x), '%H:%i'), arrayJoin(a)) as periods
from base
Получим следующие данные…
Как, собственно говоря, и ожидалось!
Полезные советы
1. Наведите порядок в каталоге xml
Блок functions в xml-файлах может содержать сразу несколько блоков функций. Определение нескольких функций из одного xml-файла полезно для поддержания порядка в каталоге /etc/clickhouse-server/. Пример src/clickhouse/xml/_combined_function.xml .
2. Обработка ошибок
В производственных условиях эти функции могут применяться сразу к большому количеству записей, возможно, как часть запроса, включающего в себя группу by, где каждая группа вызывает функцию по отдельности.
Если один вызов функции не сработает, то обязательно сработает весь запрос! Вы можете обернуть всю логику UDF в блок try except. Тогда, если один бит группы by не сработает, он, по крайней мере, не будет блокировать работу других битов. Это выглядит как-то так...
if __name__ == '__main__':
for line in sys.stdin:
try: r = some_logic(*args) except: r = [] # an empty version of the return_type.. result = json.dumps({'result': r})
print(result, end='\n')
sys.stdout.flush()
3. Dockerfile для предварительно установленного python
Вместо того чтобы устанавливать пакеты python (и других нестандартных библиотек) вручную через pip, установим их как часть сборки Docker, создав Dockerfile. Например:
FROM clickhouse/clickhouse-server:23.9.3.12 RUN apt update && apt install python3 -y RUN apt install -y python-numpy RUN apt install -y python-scipy
Затем (из рабочего каталога Dockerfile) создадим собственный образ Docker:
chmod +x Dockerfile # only needed once sudo docker build -t <your_image_name> . # choose an image name sudo docker images # check image exists
Затем просто замените образ clickhouse/clickhouse-server:23.9.3.12, используемый в скрипте init.sh в начале этой статьи, на < your_image_name>.









