Changeset: 85eadf47444d for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=85eadf47444d
Modified Files:
gdk/gdk_analytic.c
gdk/gdk_analytic.h
sql/backends/monet5/sql_rank.c
sql/backends/monet5/sql_rank.h
sql/backends/monet5/sql_rank.mal
sql/backends/monet5/sql_rank.mal.sh
sql/common/sql_types.c
Branch: analytics
Log Message:
Implemented lag window function in MAL and GDK layers.
Now I have to check how I'm going to extend the SQL compiler to support
aggregates with 3 parameters.
diffs (truncated from 1134 to 300 lines):
diff --git a/gdk/gdk_analytic.c b/gdk/gdk_analytic.c
--- a/gdk/gdk_analytic.c
+++ b/gdk/gdk_analytic.c
@@ -162,13 +162,15 @@ GDKanalyticaldiff(BAT *r, BAT *b, BAT *c
} while(0);
gdk_return
-GDKanalyticalntile(BAT *r, BAT *b, BAT *p, BAT *o, int tpe, ptr ntile)
+GDKanalyticalntile(BAT *r, BAT *b, BAT *p, BAT *o, int tpe, const void*
restrict ntile)
{
BUN cnt = BATcount(b);
bit *np, *pnp;
bool has_nils = false;
gdk_return gdk_res = GDK_SUCCEED;
+ assert(ntile);
+
switch (tpe) {
case TYPE_bte:
ANALYTICAL_NTILE_IMP(bte)
@@ -496,6 +498,7 @@ GDKanalyticalnthvalue(BAT *r, BAT *b, BA
gdk_return gdk_res = GDK_SUCCEED;
bool has_nils = false;
+ assert(is_lng_nil(nth) || nth >= 0);
(void) o;
switch (tpe) {
case TYPE_bte:
@@ -584,6 +587,142 @@ finish:
#undef ANALYTICAL_NTHVALUE_IMP
+#define ANALYTICAL_LAG_IMP(TPE) \
+ do { \
+ TPE *rp, *rb, *bp, *end, def = *((TPE *) default_value); \
+ bp = (TPE*)Tloc(b, 0); \
+ rb = rp = (TPE*)Tloc(r, 0); \
+ if(is_lng_nil(lag)) { \
+ has_nils = true; \
+ end = rb + cnt; \
+ for(; rb<end; rb++) \
+ *rb = TPE##_nil;
\
+ } else if(p) { \
+ end = rp + cnt; \
+ np = (bit*)Tloc(p, 0); \
+ for(; rp<end; np++, rp++) { \
+ if (*np) {
\
+ bp += (rp - rb);
\
+ if(lag > 0) {
\
+ for(; i<lag && rb<rp; i++,
rb++) \
+ *rb = def;
\
+ if(is_##TPE##_nil(def))
\
+ has_nils = true;
\
+ }
\
+ bp += i;
\
+ for(;rb<rp; rb++, bp++) {
\
+ *rb = *bp;
\
+ if(is_##TPE##_nil(*rb))
\
+ has_nils = true;
\
+ }
\
+ i = 0;
\
+ }
\
+ } \
+ bp += (rp - rb); \
+ if(lag > 0) { \
+ for(; i<lag && rb<end; i++, rb++)
\
+ *rb = def;
\
+ if(is_##TPE##_nil(def))
\
+ has_nils = true;
\
+ } \
+ bp += i; \
+ for(;rb<end; rb++, bp++) { \
+ *rb = *bp;
\
+ if(is_##TPE##_nil(*rb))
\
+ has_nils = true;
\
+ } \
+ } else { \
+ end = rb + cnt; \
+ if(lag > 0) { \
+ for(; i<lag && rb<end; i++, rb++)
\
+ *rb = def;
\
+ if(is_##TPE##_nil(def))
\
+ has_nils = true;
\
+ } \
+ bp += i; \
+ for(;rb<end; rb++, bp++) { \
+ *rb = *bp;
\
+ if(is_##TPE##_nil(*rb))
\
+ has_nils = true;
\
+ } \
+ } \
+ goto finish; \
+ } while(0);
+
+gdk_return
+GDKanalyticallag(BAT *r, BAT *b, BAT *p, BAT *o, lng lag, const void* restrict
default_value, int tpe)
+{
+ /*int (*atomcmp)(const void *, const void *);
+ const void *nil;*/
+ BUN cnt = BATcount(b);
+ lng i = 0;
+ bit *np;
+ gdk_return gdk_res = GDK_SUCCEED;
+ bool has_nils = false;
+
+ assert(default_value);
+ assert(is_lng_nil(lag) || lag >= 0);
+
+ (void) o;
+ switch (tpe) {
+ case TYPE_bte:
+ ANALYTICAL_LAG_IMP(bte)
+ break;
+ case TYPE_sht:
+ ANALYTICAL_LAG_IMP(sht)
+ break;
+ case TYPE_int:
+ ANALYTICAL_LAG_IMP(int)
+ break;
+ case TYPE_lng:
+ ANALYTICAL_LAG_IMP(lng)
+ break;
+#ifdef HAVE_HGE
+ case TYPE_hge:
+ ANALYTICAL_LAG_IMP(hge)
+ break;
+#endif
+ case TYPE_flt:
+ ANALYTICAL_LAG_IMP(flt)
+ break;
+ case TYPE_dbl:
+ ANALYTICAL_LAG_IMP(dbl)
+ break;
+ default: {
+ }
+ }
+finish:
+ BATsetcount(r, cnt);
+ r->tnonil = !has_nils;
+ r->tnil = has_nils;
+ return gdk_res;
+}
+
+#undef ANALYTICAL_LAG_IMP
+
+gdk_return
+GDKanalyticallead(BAT *r, BAT *b, BAT *p, BAT *o, lng lead, const void*
restrict default_value, int tpe)
+{
+ //int (*atomcmp)(const void *, const void *);
+ //const void *nil;
+ BUN /*i, j,*/ cnt = BATcount(b);
+ //bit *np;
+ gdk_return gdk_res = GDK_SUCCEED;
+ bool has_nils = false;
+
+ assert(default_value);
+ assert(is_lng_nil(lead) || lead <= 0);
+
+ (void) o;
+ (void) p;
+ (void) tpe;
+//finish:
+ BATsetcount(r, cnt);
+ r->tnonil = !has_nils;
+ r->tnil = has_nils;
+ return gdk_res;
+}
+
#define ANALYTICAL_LIMIT_IMP(TPE, OP) \
do { \
TPE *rp, *rb, *restrict bp, *end, curval; \
@@ -807,6 +946,7 @@ GDKanalyticalcount(BAT *r, BAT *b, BAT *
BUN i, cnt = BATcount(b);
gdk_return gdk_res = GDK_SUCCEED;
+ assert(ignore_nils);
(void) o;
if(!*ignore_nils || b->T.nonil) {
bit *np, *pnp;
diff --git a/gdk/gdk_analytic.h b/gdk/gdk_analytic.h
--- a/gdk/gdk_analytic.h
+++ b/gdk/gdk_analytic.h
@@ -17,13 +17,15 @@
#include "gdk.h"
gdk_export gdk_return GDKanalyticaldiff(BAT *r, BAT *b, BAT *c, int tpe);
-gdk_export gdk_return GDKanalyticalntile(BAT *r, BAT *b, BAT *p, BAT *o, int
tpe, ptr ntile);
+gdk_export gdk_return GDKanalyticalntile(BAT *r, BAT *b, BAT *p, BAT *o, int
tpe, const void* restrict ntile);
gdk_export gdk_return GDKanalyticalfirst(BAT *r, BAT *b, BAT *p, BAT *o, int
tpe);
gdk_export gdk_return GDKanalyticallast(BAT *r, BAT *b, BAT *p, BAT *o, int
tpe);
gdk_export gdk_return GDKanalyticalnthvalue(BAT *r, BAT *b, BAT *p, BAT *o,
lng nth, int tpe);
+gdk_export gdk_return GDKanalyticallag(BAT *r, BAT *b, BAT *p, BAT *o, lng
lag, const void* restrict default_value, int tpe);
+gdk_export gdk_return GDKanalyticallead(BAT *r, BAT *b, BAT *p, BAT *o, lng
lead, const void* restrict default_value, int tpe);
gdk_export gdk_return GDKanalyticalmin(BAT *r, BAT *b, BAT *p, BAT *o, int
tpe);
gdk_export gdk_return GDKanalyticalmax(BAT *r, BAT *b, BAT *p, BAT *o, int
tpe);
-gdk_export gdk_return GDKanalyticalcount(BAT *r, BAT *b, BAT *p, BAT *o, const
bit *ignore_nils, int tpe);
+gdk_export gdk_return GDKanalyticalcount(BAT *r, BAT *b, BAT *p, BAT *o, const
bit* restrict ignore_nils, int tpe);
gdk_export gdk_return GDKanalyticalsum(BAT *r, BAT *b, BAT *p, BAT *o, int
tp1, int tp2);
gdk_export gdk_return GDKanalyticalprod(BAT *r, BAT *b, BAT *p, BAT *o, int
tp1, int tp2);
gdk_export gdk_return GDKanalyticalavg(BAT *r, BAT *b, BAT *p, BAT *o, int
tpe);
diff --git a/sql/backends/monet5/sql_rank.c b/sql/backends/monet5/sql_rank.c
--- a/sql/backends/monet5/sql_rank.c
+++ b/sql/backends/monet5/sql_rank.c
@@ -52,7 +52,7 @@ SQLdiff(Client cntxt, MalBlkPtr mb, MalS
if(gdk_code == GDK_SUCCEED)
BBPkeepref(*res = r->batCacheid);
else
- throw(SQL, "sql.diff", SQLSTATE(HY001) "Unknown GDK
error");
+ throw(SQL, "sql.diff", SQLSTATE(HY001) MAL_MALLOC_FAIL);
} else {
bit *res = getArgReference_bit(stk, pci, 0);
@@ -536,7 +536,7 @@ SQLntile(Client cntxt, MalBlkPtr mb, Mal
if(gdk_code == GDK_SUCCEED)
BBPkeepref(*res = r->batCacheid);
else
- throw(SQL, "sql.ntile", SQLSTATE(HY001) "Unknown GDK
error");
+ throw(SQL, "sql.ntile", SQLSTATE(HY001)
MAL_MALLOC_FAIL);
} else {
ptr res = getArgReference_ptr(stk, pci, 0);
ptr in = getArgReference_ptr(stk, pci, 1);
@@ -757,7 +757,7 @@ SQLnth_value(Client cntxt, MalBlkPtr mb,
if(gdk_code == GDK_SUCCEED)
BBPkeepref(*res = r->batCacheid);
else
- throw(SQL, "sql.nth_value", SQLSTATE(HY001) "Unknown
GDK error");
+ throw(SQL, "sql.nth_value", SQLSTATE(HY001)
MAL_MALLOC_FAIL);
} else {
ptr res = getArgReference_ptr(stk, pci, 0);
ptr in = getArgReference_ptr(stk, pci, 1);
@@ -791,6 +791,134 @@ SQLnth_value(Client cntxt, MalBlkPtr mb,
#undef NTH_VALUE_IMP
#undef NTH_VALUE_SINGLE_IMP
+#define CHECK_L_VALUE(TPE)
\
+ do {
\
+ TPE rval = *getArgReference_##TPE(stk, pci, 2);
\
+ l_value = is_##TPE##_nil(rval) ? lng_nil : (rval > 0 ?
default_l * (TPE)rval : m * (TPE)rval); \
+ } while(0);
+
+static str /* the variable m is used to fix the multiplier */
+do_lead_lag(Client cntxt, MalBlkPtr mb, MalStkPtr stk, InstrPtr pci, const str
op, const str desc,
+ gdk_return (*func)(BAT *, BAT *, BAT *, BAT *, lng,
const void* restrict, int), lng default_l, lng m)
+{
+ int tp1, tp2, tp3, base = 2;
+ lng l_value = default_l;
+ const void *restrict default_value;
+ size_t default_value_size = 0;
+
+ (void)cntxt;
+ if (pci->argc < 8 || pci->argc > 10)
+ throw(SQL, op, SQLSTATE(42000) "%s called with invalid number
of arguments", desc);
+
+ tp1 = getArgType(mb, pci, 1);
+
+ if (pci->argc > 8) { //contains (lag or lead) value;
+ tp2 = getArgType(mb, pci, 2);
+ if (isaBatType(tp2))
+ throw(SQL, op, SQLSTATE(42000) "%s second argument must
a single atom", desc);
+ switch (tp2) {
+ case TYPE_bte:
+ CHECK_L_VALUE(bte)
+ break;
+ case TYPE_sht:
+ CHECK_L_VALUE(sht)
+ break;
+ case TYPE_int:
+ CHECK_L_VALUE(int)
+ break;
+ case TYPE_lng:
+ CHECK_L_VALUE(lng)
+ break;
+#ifdef HAVE_HGE
+ case TYPE_hge:
+ CHECK_L_VALUE(hge)
+ break;
+#endif
+ default:
+ throw(SQL, "sql.lag", SQLSTATE(42000) "%s value
not available for %s", desc, ATOMname(tp2));
+ }
+ base = 3;
+ }
+
+ if (pci->argc > 9) { //contains default value;
+ ValRecord *vin = &(stk)->stk[(pci)->argv[3]];
+ tp3 = getArgType(mb, pci, 3);
+ if (isaBatType(tp3))
+ throw(SQL, op, SQLSTATE(42000) "%s third argument must
a single atom", desc);
+ default_value = vin->val.pval;
+ default_value_size = vin->len;
+ base = 4;
+ } else {
+ int tpe = tp1;
+ if (isaBatType(tpe))
+ tpe = getBatType(tp1);
+ default_value = ATOMnilptr(tpe);
+ default_value_size = ATOMlen(tpe, default_value);
+ }
+
+ assert(default_value); //default value must be set
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list