11/*
22[dependencies]
3+ dsrs = { git = "https://github.com/ProjectASAP/datasketches-rs", rev = "d748ec75c80fff21f7b24897244dd1c895df2e9a" }
34asap_sketchlib = { git = "https://github.com/ProjectASAP/asap_sketchlib" }
45arroyo-udf-plugin = "0.1"
56rmp-serde = "1.1"
@@ -8,11 +9,24 @@ xxhash-rust = { version = "0.8", features = ["xxh32"] }
89*/
910
1011use arroyo_udf_plugin::udf;
12+ use dsrs::KllDoubleSketch;
1113use rmp_serde::Serializer;
1214use serde::{Deserialize, Serialize};
1315use asap_sketchlib::KLL;
1416use xxhash_rust::xxh32::xxh32;
1517
18+ enum ImplMode {
19+ Legacy,
20+ Sketchlib,
21+ }
22+
23+ {% set _impl_mode = impl_mode | default ("Sketchlib" ) %}
24+ const IMPL_MODE: ImplMode = ImplMode::{% if _impl_mode == "Legacy" or _impl_mode == "Sketchlib" %} {{ _impl_mode }}{% else %} Sketchlib{% endif %} ;
25+
26+ fn use_sketchlib_for_kll() -> bool {
27+ matches!(IMPL_MODE, ImplMode::Sketchlib)
28+ }
29+
1630const ROW_NUM: usize = {{ row_num }};
1731const COL_NUM: usize = {{ col_num }};
1832const DEFAULT_K: u16 = {{ k }};
@@ -32,40 +46,76 @@ struct HydraKllSketchData {
3246
3347#[udf]
3448fn hydrakll_(keys: Vec<&str>, values: Vec<f64 >) -> Option<Vec <u8 >> {
35- let mut sketches: Vec<Vec <KLL >> = (0..ROW_NUM)
36- .map(|_| {
37- (0..COL_NUM)
38- .map(|_| KLL::init_kll(DEFAULT_K as i32))
39- .collect()
40- })
41- .collect();
49+ let sketch_data: Vec<Vec <KllSketchData >> = if use_sketchlib_for_kll() {
50+ let mut sketches: Vec<Vec <KLL >> = (0..ROW_NUM)
51+ .map(|_| {
52+ (0..COL_NUM)
53+ .map(|_| KLL::init_kll(DEFAULT_K as i32))
54+ .collect()
55+ })
56+ .collect();
4257
43- for (i, &key) in keys.iter().enumerate() {
44- if i >= values.len() {
45- break;
58+ for (i, &key) in keys.iter().enumerate() {
59+ if i >= values.len() {
60+ break;
61+ }
62+ let key_bytes = key.as_bytes();
63+ for row in 0..ROW_NUM {
64+ let hash_value = xxh32(key_bytes, row as u32);
65+ let col_index = (hash_value as usize) % COL_NUM;
66+ sketches[row][col_index].update(&values[i]);
67+ }
4668 }
47- let key_bytes = key.as_bytes();
48- for row in 0..ROW_NUM {
49- let hash_value = xxh32(key_bytes, row as u32);
50- let col_index = (hash_value as usize) % COL_NUM;
51- sketches[row][col_index].update(&values[i]);
69+
70+ sketches
71+ .iter()
72+ .map(|row| {
73+ row.iter()
74+ .map(|sketch| {
75+ let sketch_bytes = sketch.serialize_to_bytes().ok()?;
76+ Some(KllSketchData {
77+ k: DEFAULT_K,
78+ sketch_bytes,
79+ })
80+ })
81+ .collect::<Option <Vec <_ >>>()
82+ })
83+ .collect::<Option <Vec <_ >>>()?
84+ } else {
85+ let mut sketches: Vec<Vec <KllDoubleSketch >> = (0..ROW_NUM)
86+ .map(|_| {
87+ (0..COL_NUM)
88+ .map(|_| KllDoubleSketch::with_k(DEFAULT_K))
89+ .collect()
90+ })
91+ .collect();
92+
93+ for (i, &key) in keys.iter().enumerate() {
94+ if i >= values.len() {
95+ break;
96+ }
97+ let key_bytes = key.as_bytes();
98+ for row in 0..ROW_NUM {
99+ let hash_value = xxh32(key_bytes, row as u32);
100+ let col_index = (hash_value as usize) % COL_NUM;
101+ sketches[row][col_index].update(values[i]);
102+ }
52103 }
53- }
54104
55- let sketch_data: Vec< Vec < KllSketchData >> = sketches
56- .iter()
57- .map(|row| {
58- row.iter()
59- .map(|sketch| {
60- let sketch_bytes = sketch.serialize_to_bytes().ok()?;
61- Some(KllSketchData {
62- k: DEFAULT_K ,
63- sketch_bytes,
105+ sketches
106+ .iter()
107+ .map(|row| {
108+ row.iter()
109+ .map(|sketch| {
110+ Some(KllSketchData {
111+ k: DEFAULT_K,
112+ sketch_bytes: sketch.serialize().as_ref().to_vec() ,
113+ })
64114 })
65- } )
66- .collect::< Option < Vec < _ >>>( )
67- })
68- .collect::< Option < Vec < _ >>>()? ;
115+ .collect::< Option < Vec < _ >>>( )
116+ } )
117+ .collect::< Option < Vec < _ >>>()?
118+ } ;
69119
70120 let hydra_data = HydraKllSketchData {
71121 row_num: ROW_NUM,
0 commit comments