I have been attempting to use all cores available on a node when calling a UDx function. We have some need to run intensive machine learning computation at scale, but don't want to use unfenced C++ mode yet. The usual way to do parallel processing using partitions given here: https://forum.vertica.com/discussion/241038/parallel-processing-using-partitions-with-vertica-udx does not exercise all cores on the node. Maybe because is it running in fenced mode and is limited to one thread.
An alternate way would be to spawn process within the UDx and collect all the results in the Python code and emit the output. So we want to run a Multiprocess Python program within a Transform function. Here is the code for a simple use case.
This is inspired from the Python example given here (Look at the very bottom of the page to see the example):
https://docs.python.org/3/library/multiprocessing.html#multiprocessing-programming
This works fine when running on the node in standalone mode.
Now we will try to use the same principles within a UDx function.
The premise is that we a table as follows:
CREATE TABLE IF NOT EXISTS dev.basic_queue_records
(
a int,
b int,
operand varchar(10)
);
We would just want to write a UDX that takes the two integers and the operand, and produces an additional column (result) based on the operand. If the operand is 'plus' it adds the two integers. If the operand is 'mul' it multiples the two numbers. The UDX will emit four fields viz. (a, b, operand, result). This is the same problem given in the Python example above but done on an SQL table.
Here is the UDX code in Python:
import vertica_sdk
import time
import random
from multiprocessing import Process, Queue, freeze_support
def worker(input, output):
for func, args in iter(input.get, 'STOP'):
result = calculate(func, args)
output.put(result)
def calculate(func, args):
result = func(*args)
return args[0], args[1], func.__name__, result
def mul(a, b):
time.sleep(0.5*random.random())
return a * b
def plus(a, b):
time.sleep(0.5*random.random())
return a + b
class BasicQueueUDX(vertica_sdk.TransformFunction):
def processPartition(self, server_interface, input, output):
"""
:param server_interface:
:param input:
:param output:
:return:
"""
freeze_support()
NUMBER_OF_PROCESSES = 4
TASKS = []
while True:
a = input.getInt(0)
b = input.getInt(1)
operand = input.getString(2)
TASKS.append((operand, (a, b)))
if not input.next():
break
# Create queues
task_queue = Queue()
done_queue = Queue()
# Submit tasks
for task in TASKS:
task_queue.put(task)
# Start worker processes
for i in range(NUMBER_OF_PROCESSES):
Process(target=worker, args=(task_queue, done_queue)).start()
# Get and print results
server_interface.log(f"About to process {len(TASKS)} tasks using {NUMBER_OF_PROCESSES} processes")
for i in range(len(TASKS)):
out_row = done_queue.get()
server_interface.log(f"{int(out_row[0])}, {int(out_row[1])}, {out_row[2]}, {int(out_row[3])}")
output.setInt(0, int(out_row[0])) # a
output.setInt(1, int(out_row[1])) # b
output.setString(2, out_row[2]) # operand
output.setInt(3, int(out_row[3])) # result
output.next()
# Tell child processes to stop
for i in range(NUMBER_OF_PROCESSES):
task_queue.put('STOP')
server_interface.log(f"All records processed. Exiting.")
class BasicQueueUDXFactory(vertica_sdk.TransformFunctionFactory):
def getPrototype(self, server_interface, arg_types, return_type):
arg_types.addInt()
arg_types.addInt()
arg_types.addVarchar()
return_type.addInt()
return_type.addInt()
return_type.addVarchar()
return_type.addInt()
def getReturnType(self, server_interface, arg_types, return_type):
return_type.addColumn(arg_types.getColumnType(0), "a")
return_type.addColumn(arg_types.getColumnType(1), "b")
return_type.addColumn(arg_types.getColumnType(2), "operand")
return_type.addColumn(arg_types.getColumnType(1), "result")
def createTransformFunction(cls, server_interface):
return BasicQueueUDX()
Assume that the file exists in path '/home/dbadmin/udx/BasicQueueUDX.py'.
To install the library, we do the following:
DROP LIBRARY IF EXISTS pymultiprocess CASCADE; CREATE LIBRARY pymultiprocess AS '/home/dbadmin/udx/BasicQueueUDX.py' DEPENDS '/usr/local/lib/python3.7/site-packages/*:/usr/local/lib64/python3.7/site-packages/*' LANGUAGE 'Python'; CREATE TRANSFORM FUNCTION BASIC_QUEUE AS NAME 'BasicQueueUDXFactory' LIBRARY pymultiprocess;
Populate the table with some rows:
INSERT INTO dev.basic_queue_records (a, b, operand) VALUES (7, 3, 'plus'); INSERT INTO dev.basic_queue_records (a, b, operand) VALUES (7, 4, 'mul'); INSERT INTO dev.basic_queue_records (a, b, operand) VALUES (9, 7, 'plus'); INSERT INTO dev.basic_queue_records (a, b, operand) VALUES (10, 3, 'mul');
Then call the UDx function.
SELECT BASIC_QUEUE(a, b, operand) OVER () FROM dev.basic_queue_records;
This essentially locks with the following output on the server logs. It runs for ever. When the process is cancelled, it calls the Python destructor with no output.
2021-09-20 20:34:31.751 [Python-v_testdb2_node0001-5537:0x13ed80-64744] 0x7fa400b74740 UDx side process started 2021-09-20 20:34:31.751 [Python-v_testdb2_node0001-5537:0x13ed80-64744] 0x7fa400b74740 My port: 59201 2021-09-20 20:34:31.751 [Python-v_testdb2_node0001-5537:0x13ed80-64744] 0x7fa400b74740 My address: 127.0.0.1 2021-09-20 20:34:31.751 [Python-v_testdb2_node0001-5537:0x13ed80-64744] 0x7fa400b74740 Vertica port: 45475 2021-09-20 20:34:31.751 [Python-v_testdb2_node0001-5537:0x13ed80-64744] 0x7fa400b74740 Vertica address: 127.0.0.1 2021-09-20 20:34:31.751 [Python-v_testdb2_node0001-5537:0x13ed80-64744] 0x7fa400b74740 Vertica address family: 2 2021-09-20 20:34:43.829 [Python-v_testdb2_node0001-5537:0x13ed80-64787] 0x7fa400b74740 PythonInterface destructor called 2021-09-20 20:34:43.983 [Python-v_testdb2_node0001-5537:0x13ed80-64788] 0x7fa400b74740 [UserMessage] BASIC_QUEUE - About to process 4 tasks using 4 processes
Anyone has any clues why this in locking within a Vertica context, but seems to run fine when called as an application?