Changeset: e04248dbe275 for MonetDB URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=e04248dbe275 Modified Files: gdk/gdk_analytic.c sql/backends/monet5/sql_rank.c sql/backends/monet5/sql_rank.mal sql/backends/monet5/sql_rank.mal.sh sql/server/sql_scan.c Branch: analytics Log Message:
Fixed lag implementation.
diffs (truncated from 620 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
@@ -247,7 +247,7 @@ gdk_return
GDKanalyticalfirst(BAT *r, BAT *b, BAT *p, BAT *o, int tpe)
{
int (*atomcmp)(const void *, const void *);
- const void *nil;
+ const void* restrict nil;
bool has_nils = false;
BUN i, j, cnt = BATcount(b);
bit *restrict np;
@@ -359,7 +359,7 @@ gdk_return
GDKanalyticallast(BAT *r, BAT *b, BAT *p, BAT *o, int tpe)
{
int (*atomcmp)(const void *, const void *);
- const void *nil;
+ const void* restrict nil;
bool has_nils = false;
BUN i, j, cnt = BATcount(b);
bit *restrict np;
@@ -492,7 +492,7 @@ gdk_return
GDKanalyticalnthvalue(BAT *r, BAT *b, BAT *p, BAT *o, lng nth, int tpe)
{
int (*atomcmp)(const void *, const void *);
- const void *nil;
+ const void* restrict nil;
BUN i, j, cnt = BATcount(b);
bit *np;
gdk_return gdk_res = GDK_SUCCEED;
@@ -587,76 +587,69 @@ 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; \
+#define ANALYTICAL_LAG_IMP(TPE) \
+ do { \
+ TPE *rp, *rb, *bp, *end, \
+ def = *((TPE *) default_value), next; \
+ 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) { \
+ for(i=0; i<lag && rb<rp; i++, rb++) \
+ *rb = def;
\
+ if(lag > 0 && is_##TPE##_nil(def)) \
+ has_nils = true;
\
+ for(;rb<rp; rb++, bp++) { \
+ next = *bp;
\
+ *rb = next;
\
+ if(is_##TPE##_nil(next))
\
+ has_nils = true;
\
+ } \
+ bp += lag; \
+ } \
+ } \
+ for(i=0; i<lag && rb<end; i++, rb++) \
+ *rb = def; \
+ if(lag > 0 && is_##TPE##_nil(def)) \
+ has_nils = true; \
+ for(;rb<end; rb++, bp++) { \
+ next = *bp; \
+ *rb = next; \
+ if(is_##TPE##_nil(next)) \
+ has_nils = true; \
+ } \
+ } else { \
+ end = rb + cnt; \
+ for(i=0; i<lag && rb<end; i++, rb++) \
+ *rb = def; \
+ if(lag > 0 && is_##TPE##_nil(def)) \
+ has_nils = true; \
+ for(;rb<end; rb++, bp++) { \
+ next = *bp; \
+ *rb = next; \
+ if(is_##TPE##_nil(next)) \
+ 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);
+ int (*atomcmp)(const void *, const void *);
+ const void *restrict nil;
lng i = 0;
- bit *np;
+ BUN j = 0, k = 0, l, cnt = BATcount(b);
+ bit *restrict np;
gdk_return gdk_res = GDK_SUCCEED;
bool has_nils = false;
@@ -689,6 +682,64 @@ GDKanalyticallag(BAT *r, BAT *b, BAT *p,
ANALYTICAL_LAG_IMP(dbl)
break;
default: {
+ BATiter bpi = bat_iterator(b);
+ const void *restrict curval;
+ nil = ATOMnilptr(tpe);
+ atomcmp = ATOMcompare(tpe);
+ if(is_lng_nil(lag)) {
+ has_nils = true;
+ for (j=0;j < cnt; j++) {
+ if ((gdk_res = BUNappend(r, nil,
false)) != GDK_SUCCEED)
+ goto finish;
+ }
+ } else if(p) {
+ np = (bit*)Tloc(p, 0);
+ for(j=0,k=0; j<cnt; j++, np++) {
+ if (*np) {
+ for(i=0; i<lag && k<j; i++,
k++) {
+ if ((gdk_res =
BUNappend(r, default_value, false)) != GDK_SUCCEED)
+ goto finish;
+ }
+ if(lag > 0 &&
(*atomcmp)(default_value, nil) == 0)
+ has_nils = true;
+ for(l=0; k<j; k++, l++) {
+ curval = BUNtail(bpi,
l);
+ if ((gdk_res =
BUNappend(r, curval, false)) != GDK_SUCCEED)
+ goto finish;
+ if((*atomcmp)(curval,
nil) == 0)
+ has_nils = true;
+ }
+ }
+ }
+ for(i=0; i<lag && k<cnt; i++, k++) {
+ if ((gdk_res = BUNappend(r,
default_value, false)) != GDK_SUCCEED)
+ goto finish;
+ }
+ if(lag > 0 && (*atomcmp)(default_value, nil) ==
0)
+ has_nils = true;
+ for(l=0; k<cnt; k++, l++) {
+ curval = BUNtail(bpi, l);
+ if ((gdk_res = BUNappend(r, curval,
false)) != GDK_SUCCEED)
+ goto finish;
+ if((*atomcmp)(curval, nil) == 0)
+ has_nils = true;
+ }
+ } else {
+ lng lcnt = (lng) cnt;
+ for(i=0; i<lag && i<lcnt; i++) {
+ if ((gdk_res = BUNappend(r,
default_value, false)) != GDK_SUCCEED)
+ goto finish;
+ }
+ if(lag > 0 && (*atomcmp)(default_value, nil) ==
0)
+ has_nils = true;
+ for(l=0, k=(BUN)lag; k<cnt; k++, l++) {
+ curval = BUNtail(bpi, l);
+ if ((gdk_res = BUNappend(r, curval,
false)) != GDK_SUCCEED)
+ goto finish;
+ if((*atomcmp)(curval, nil) == 0)
+ has_nils = true;
+ }
+ }
}
}
finish:
@@ -704,9 +755,9 @@ 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;
+ //const void* restrict nil;
BUN /*i, j,*/ cnt = BATcount(b);
- //bit *np;
+ //bit *restrict np;
gdk_return gdk_res = GDK_SUCCEED;
bool has_nils = false;
@@ -787,7 +838,7 @@ gdk_return
GDKanalytical##OP(BAT *r, BAT *b, BAT *p, BAT *o, int tpe)
\
{
\
int (*atomcmp)(const void *, const void *);
\
- const void *nil;
\
+ const void* restrict nil;
\
bool has_nils = false;
\
BUN i, j, cnt = BATcount(b);
\
bit *restrict np;
\
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
@@ -835,7 +835,7 @@ do_lead_lag(Client cntxt, MalBlkPtr mb,
break;
#endif
default:
- throw(SQL, "sql.lag", SQLSTATE(42000) "%s value
not available for %s", desc, ATOMname(tp2));
+ throw(SQL, op, SQLSTATE(42000) "%s value not
available for %s", desc, ATOMname(tp2));
}
base = 3;
}
@@ -868,16 +868,16 @@ do_lead_lag(Client cntxt, MalBlkPtr mb,
gdk_return gdk_code;
tp1 = getBatType(tp1);
- voidresultBAT(r, tp1, cnt, b, "sql.lag");
- if (isaBatType(getArgType(mb, pci, base + 1))) {
- p = BATdescriptor(*getArgReference_bat(stk, pci, base +
1));
+ voidresultBAT(r, tp1, cnt, b, op);
+ if (isaBatType(getArgType(mb, pci, base))) {
+ p = BATdescriptor(*getArgReference_bat(stk, pci, base));
if (!p) {
BBPunfix(b->batCacheid);
throw(SQL, op, SQLSTATE(HY005) "Cannot access
column descriptor");
}
}
- if (isaBatType(getArgType(mb, pci, base + 2))) {
- o = BATdescriptor(*getArgReference_bat(stk, pci, base +
2));
+ if (isaBatType(getArgType(mb, pci, base + 1))) {
+ o = BATdescriptor(*getArgReference_bat(stk, pci, base +
1));
if (!o) {
BBPunfix(b->batCacheid);
BBPunfix(p->batCacheid);
diff --git a/sql/backends/monet5/sql_rank.mal b/sql/backends/monet5/sql_rank.mal
--- a/sql/backends/monet5/sql_rank.mal
+++ b/sql/backends/monet5/sql_rank.mal
@@ -124,54 +124,54 @@ address SQLlag
comment "return the value in the previous row in the partition or NULL if non
existent";
-pattern sql.lag(b:any_1, lag:any_2, p:any_3, o:any_4, unit:int, s:int, e:int,
exl:int) :any_1
+pattern sql.lag(b:any_1, l:any_2, p:any_3, o:any_4, unit:int, s:int, e:int,
exl:int) :any_1
address SQLlag
-comment "return the value in the previous 'lag' row in the partition or NULL
if non existent";
+comment "return the value in the previous 'l' row in the partition or NULL if
non existent";
-pattern batsql.lag(b:bat[:any_1], lag:any_2, p:any_3, o:any_4, unit:int,
s:int, e:int, exl:int) :bat[:any_1]
+pattern batsql.lag(b:bat[:any_1], l:any_2, p:any_3, o:any_4, unit:int, s:int,
e:int, exl:int) :bat[:any_1]
address SQLlag
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list
