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

Reply via email to