I have a transform UDF that I want to run per zip_prefix (first 2 digits of zipcode). The problem is that the data distribution per this zip prefix is quite skewed. So I'd like to evenly distribute subsets of zip_prefixes to each node and run the UDF on each node (e.g., partition nodes order by zip_prefix). In this case, my UDF would handle splitting up the data by zip_prefix as it streams in to the UDF.
The problem is I do not see how I can control the segmentation in such a manner.
For instance, you can see if I just segment by hash(zip_prefix), the amount of data processed per node will be uneven:
create local temp table t on commit preserve rows as
select zip_prefix, 1 matrix, site_idx, max(site_idx) over (partition by zip_prefix) site_len, aid_idx, max(aid_idx) over ( partition by zip_prefix) aid_len, weight
from ci_dm.ul_site_edge
where batch_id = 1
union all
select zip_prefix, 2 matrix, aid1_idx, 0, aid2_idx, 0, weight
from ci_dm.ul_aid_edge
where batch_id = 1
order by zip_prefix, matrix
segmented by hash(zip_prefix) all nodes
;
dbadmin=> select node_name, (sum(row_count)/sum(sum(row_count)) over()*100)::numeric(10,2) row_distribution_perc from projection_storage where projection_name = 't_b0' group by 1 order by 1;
node_name | row_distribution_perc
----------------+-----------------------
v_udb_node0001 | 1.43
v_udb_node0002 | 3.62
v_udb_node0003 | 4.16
v_udb_node0004 | 8.24
v_udb_node0005 | 9.61
v_udb_node0006 | 0.24
v_udb_node0007 | 7.90
v_udb_node0008 | 0.68
v_udb_node0009 | 1.78
v_udb_node0010 | 5.22
v_udb_node0011 | 4.30
v_udb_node0012 | 6.40
v_udb_node0013 | 0.50
v_udb_node0014 | 5.23
v_udb_node0015 | 5.21
v_udb_node0016 | 2.31
v_udb_node0017 | 5.52
v_udb_node0018 | 2.19
v_udb_node0019 | 5.83
v_udb_node0020 | 4.66
v_udb_node0021 | 3.34
v_udb_node0022 | 5.10
v_udb_node0023 | 2.65
v_udb_node0024 | 3.87
On a 24-node cluster, I can figure out a decent way to create subsets of zip_prefixes and assign them to nodes:
dbadmin=> select node_id, count(distinct zip_prefix) zips, (sum(cnt)/sum(sum(cnt)) over()*100)::numeric(10,2) row_distribution_perc from ( select zip_prefix, count(*) cnt, (row_number() over ( order by count(*) desc )-1) % 24 node_id from t group by 1 ) v group by 1order by 1;
node_id | zips | row_distribution_perc
---------+------+-----------------------
0 | 16 | 5.99
1 | 16 | 5.38
2 | 16 | 5.26
3 | 16 | 5.11
4 | 16 | 5.04
5 | 15 | 4.93
6 | 15 | 4.81
7 | 15 | 4.74
8 | 15 | 4.61
9 | 15 | 4.15
10 | 15 | 4.09
11 | 15 | 3.94
12 | 15 | 3.89
13 | 15 | 3.80
14 | 15 | 3.69
15 | 15 | 3.67
16 | 15 | 3.62
17 | 15 | 3.56
18 | 15 | 3.51
19 | 15 | 3.44
20 | 15 | 3.37
21 | 15 | 3.24
22 | 15 | 3.13
23 | 15 | 3.03
However, I cannot get the data to distribute in the cluster based on node_id:
dbadmin=> create local temp table t2 on commit preserve rows as
dbadmin-> select t.*, n.node_id
dbadmin-> from t, ( select zip_prefix, count(*) cnt, (row_number() over ( order by count(*) desc )-1) % 24 node_id from t group by 1 ) n
dbadmin-> where n.zip_prefix = t.zip_prefix
dbadmin-> order by node_id, zip_prefix, matrix
dbadmin-> segmented by node_id all nodes
dbadmin-> ;
WARNING 5993: Projection is irregularly segmented by column
HINT: Consider using a segmentation expression, such as SEGMENTED BY HASH(column)
CREATE TABLE
dbadmin=>
dbadmin=>
dbadmin=> select node_name, (sum(row_count)/sum(sum(row_count)) over()*100)::numeric(10,2) row_distribution_perc from projection_storage where projection_name = 't2_b0' group by 1 order by 1;
node_name | row_distribution_perc
----------------+-----------------------
v_udb_node0001 | 100.00
v_udb_node0002 | 0.00
v_udb_node0003 | 0.00
v_udb_node0004 | 0.00
v_udb_node0005 | 0.00
v_udb_node0006 | 0.00
v_udb_node0007 | 0.00
v_udb_node0008 | 0.00
v_udb_node0009 | 0.00
v_udb_node0010 | 0.00
v_udb_node0011 | 0.00
v_udb_node0012 | 0.00
v_udb_node0013 | 0.00
v_udb_node0014 | 0.00
v_udb_node0015 | 0.00
v_udb_node0016 | 0.00
v_udb_node0017 | 0.00
v_udb_node0018 | 0.00
v_udb_node0019 | 0.00
v_udb_node0020 | 0.00
v_udb_node0021 | 0.00
v_udb_node0022 | 0.00
v_udb_node0023 | 0.00
v_udb_node0024 | 0.00
Is there a segmentation clause of some other way that I can distribute the data this way?
Controlling data segmentation
Sign up
Already have an account? Login
Welcome to the Rocket Forum!
Please log in or register:
Employee Login | Registration Member Login | RegistrationEnter your E-mail address. We'll send you an e-mail with instructions to reset your password.