Changeset: 7b3c414bad37 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/7b3c414bad37
Modified Files:
        sql/backends/monet5/rel_physical.c
        sql/backends/monet5/sql_execute.c
        sql/server/rel_dump.c
        sql/server/rel_rel.c
Branch: pp_hashjoin
Log Message:

allow for initial fetchjoins (using classic joins)
improved dump/read for pipeline plans (ie explain after logical physical)


diffs (truncated from 350 to 300 lines):

diff --git a/sql/backends/monet5/rel_physical.c 
b/sql/backends/monet5/rel_physical.c
--- a/sql/backends/monet5/rel_physical.c
+++ b/sql/backends/monet5/rel_physical.c
@@ -23,8 +23,6 @@
                                (argc == 2 && (strcmp((fname), "quantile") == 0 
|| strcmp((fname), "quantile_avg") == 0)) || \
                                (argc == 1 && (strcmp((fname), "median") == 0 
|| strcmp((fname), "median_avg") == 0)))
 
-static int do_oahash_join(sql_rel *rel);
-
 /* Returns the row count of a base table or any count info we can get fom the
  * PROP_COUNT of this 'rel' (i.e.  get_rel_count()). */
 static lng
@@ -426,12 +424,35 @@ rel_groupby_partition_safe(sql_rel *rel)
 }
 
 static int
-do_oahash_join(sql_rel *rel)
+do_oahash_join(visitor *v, sql_rel *rel, int *side)
 {
        ATOMIC_TYPE oahash_enabled = (1U<<19);
        if (!(GDKdebug & oahash_enabled))
                return 0;
 
+       /* fetch join */
+       if (is_innerjoin(rel->op) && list_length(rel->exps) == 1 /* single 
JOINIDX */) {
+               sql_rel *l = rel->l;
+               sql_rel *r = rel->r;
+               list *exps = rel->exps;
+               sql_exp *je = exps->h->data;
+               prop *p = find_prop(je->p, PROP_JOINIDX);
+               if (p && (is_basetable(l->op) || is_basetable(r->op))) { /* 
only use join idx on direct basetable access */
+                       sql_idx *idx = p->value.pval;
+                       sql_table *lt = l->l, *rt = r->l;
+                       sql_trans *tr = v->sql->session->tr;
+                       sql_key *rk = (sql_key*)os_find_id(tr->cat->objects, 
tr, ((sql_fkey*)idx->key)->rkey);
+                       if (rk && is_basetable(l->op) && lt == rk->t) {
+                               *side = 2; /* primary left, ie continue with 
right hand side */
+                               return 0;
+                       }
+
+                       if (rk && is_basetable(r->op) && rt == rk->t) {
+                               *side = 1; /* primary right, ie continue with 
left hand side */
+                               return 0;
+                       }
+               }
+       }
        // TODO groupjoin other then mark/exist
     if (list_length(rel->attr) == 1) {
         sql_exp *e = rel->attr->h->data;
@@ -858,7 +879,8 @@ rel_pipeline(visitor *v, sql_rel *rel, b
                if (is_delete(rel->op) && !rel->r && pb)
                        rel_dup(rel);
        } else if (is_join(rel->op)) {
-               if (do_oahash_join(rel)) {
+               int side = 0;
+               if (do_oahash_join(v, rel, &side)) {
                        list *eq_exps = sa_list(v->sql->sa);
                        list *other = sa_list(v->sql->sa);
                        if (!list_empty(rel->attr))
@@ -984,29 +1006,11 @@ rel_pipeline(visitor *v, sql_rel *rel, b
                        if (pb)
                                rel->spb = 1;
                        res = SPB;
-               } else {
-                       if (pb && is_outerjoin(rel->op))
-                               res = 0;
-                       /* For now we only try to partition in case of a 
equi-join.
-                        * The other joins are too complex to handle. */
-                       else if (pb) { /* and rel->op == op_join */
-                               if (!rel->partition)
-                                       res = _rel_partition(v->sql, rel);
-                               if (res) {
-                                       int lres = rel_pipeline(v, rel->l, 
false, (rel->partition==1 && rel->spb)?pb:0);
-                                       if (lres == EPB && pb)
-                                               rel_dup(rel->l);
-                                       int rres = rel_pipeline(v, rel->r, 
false, (rel->partition==2 && rel->spb)?pb:0);
-                                       if (rres == EPB && pb)
-                                               rel_dup(rel->r);
-                                       if (pb)
-                                               res = 0;
-                               }
-                               if (!res) {
-                                       rel->spb = 1;
-                                       res = SPB;
-                               }
-                       }
+               } else if (pb && side) { /* handle fetch join */
+                       if (side == 1)
+                               res = rel_pipeline(v, rel->l, false, pb);
+                       else
+                               res = rel_pipeline(v, rel->r, false, pb);
                }
        } else if (is_ddl(rel->op)) {
                if (rel->flag == ddl_output || rel->flag == ddl_create_seq || 
rel->flag == ddl_alter_seq || rel->flag == ddl_alter_table || rel->flag == 
ddl_create_table || rel->flag == ddl_create_view) {
diff --git a/sql/backends/monet5/sql_execute.c 
b/sql/backends/monet5/sql_execute.c
--- a/sql/backends/monet5/sql_execute.c
+++ b/sql/backends/monet5/sql_execute.c
@@ -438,22 +438,28 @@ RAstatement(Client c, MalBlkPtr mb, MalS
                else
                        msg = createException(SQL, "RAstatement", 
SQLSTATE(42000) "%s", m->errstr);
        } else {
+               Symbol backup = NULL;
+               if (c->curprg) {
+                       backup = c->curprg;
+                       c->curprg = NULL;
+               }
                if ((msg = MSinitClientPrg(c, sql_private_module_name, "test")) 
!= MAL_SUCCEED)
                        return RAcommit_statement(be, msg);
 
                /* generate MAL code, ignoring any code generation error */
+               m->type = Q_TABLE;
                setVarType(c->curprg->def, 0, 0);
-               if (backend_dumpstmt(be, c->curprg->def, rel, 0, 1, NULL) < 0) {
+               if (backend_dumpstmt(be, c->curprg->def, rel, 1, 1, NULL) < 0) {
                        msg = createException(SQL,"RAstatement","Program 
contains errors"); // TODO: use macro definition.
                } else {
                        msg = SQLoptimizeFunction(c, c->curprg->def);
                        if (msg == MAL_SUCCEED)
                                msg = SQLrun(c, be);
-                       if (msg == MAL_SUCCEED)
-                               msg = resetMalBlk(&c->curprg->def);
                }
+               c->curprg = backup;
                rel_destroy(m, rel);
        }
+       sqlcleanup(be, 0);
        return RAcommit_statement(be, msg);
 }
 
diff --git a/sql/server/rel_dump.c b/sql/server/rel_dump.c
--- a/sql/server/rel_dump.c
+++ b/sql/server/rel_dump.c
@@ -533,7 +533,7 @@ rel_print_rel(mvc *sql, stream  *fout, s
        if ((ATOMIC_GET(&GDKdebug) & TESTINGMASK) == 0 && rel->spb)
                        mnstr_printf(fout, " start ");
        if ((ATOMIC_GET(&GDKdebug) & TESTINGMASK) == 0 && rel->parallel)
-                       mnstr_printf(fout, " || ");
+                       mnstr_printf(fout, " parallel ");
 
        if (is_single(rel))
                mnstr_printf(fout, "single ");
@@ -1973,17 +1973,20 @@ sql_rel*
 rel_read(mvc *sql, char *r, int *pos, list *refs)
 {
        sql_rel *rel = NULL, *nrel, *lrel, *rrel = NULL;
-       list *exps, *gexps, *rels = NULL;
+       list *exps, *attr, *gexps, *rels = NULL;
        int distinct = 0, dependent = 0, single = 0, recursive = 0;
        operator_type j = op_basetable;
        bool groupjoin = false;
+       bool start = 0, parallel = 0;
 
        skipWS(r,pos);
-       if (r[*pos] == 'R') {
+       while (r[*pos] == 'R') {
                *pos += (int) strlen("REF");
 
                skipWS(r, pos);
-               (void)readInt(r,pos);
+               int nr = readInt(r,pos);
+               if (list_length(refs) != (nr-1))
+                       return NULL;
                skipWS(r, pos);
                (*pos)++; /* ( */
                int cnt = readInt(r,pos);
@@ -2004,6 +2007,20 @@ rel_read(mvc *sql, char *r, int *pos, li
                return list_fetch(refs, nr-1);
        }
 
+       /* start */
+       if (r[*pos] == 's' && r[*pos+1] == 't' && r[*pos+2] == 'a' && r[*pos+3] 
== 'r' && r[*pos+4] == 't') {
+               start = true;
+               *pos += 5;
+               skipWS(r, pos);
+       }
+       /* parallel */
+       if (r[*pos] == 'p' && r[*pos+1] == 'a' && r[*pos+2] == 'r' && r[*pos+3] 
== 'a' &&
+               r[*pos+4] == 'l' && r[*pos+5] == 'l' && r[*pos+6] == 'e' && 
r[*pos+7] == 'l') {
+               parallel = true;
+               *pos += 8;
+               skipWS(r, pos);
+       }
+
        if (r[*pos] == 'i' && r[*pos+1] == 'n' && r[*pos+2] == 's') {
                sql_table *t;
 
@@ -2024,6 +2041,8 @@ rel_read(mvc *sql, char *r, int *pos, li
 
                if (!(rel = rel_insert(sql, lrel, rrel)) || !(rel = 
read_rel_properties(sql, rel, r, pos)))
                        return NULL;
+               rel->spb = start;
+               rel->parallel = parallel;
                return rel;
        }
 
@@ -2048,6 +2067,8 @@ rel_read(mvc *sql, char *r, int *pos, li
                if (!(rel = rel_delete(sql->sa, lrel, rrel)) || !(rel = 
read_rel_properties(sql, rel, r, pos)))
                        return NULL;
 
+               rel->spb = start;
+               rel->parallel = parallel;
                return rel;
        }
 
@@ -2124,6 +2145,8 @@ rel_read(mvc *sql, char *r, int *pos, li
                if (!(rel = rel_update(sql, lrel, rrel, NULL, nexps)) || !(rel 
= read_rel_properties(sql, rel, r, pos)))
                        return NULL;
 
+               rel->spb = start;
+               rel->parallel = parallel;
                return rel;
        }
 
@@ -2155,6 +2178,7 @@ rel_read(mvc *sql, char *r, int *pos, li
                recursive = 1;
        }
 
+       bool probe = false;
        switch(r[*pos]) {
        case 't':
                if (r[*pos+1] == 'a') {
@@ -2331,6 +2355,8 @@ rel_read(mvc *sql, char *r, int *pos, li
                                        rel_base_use_all(sql, rel);
                                        rel = rewrite_basetable(sql, rel, true);
                                }
+                               rel->spb = start;
+                               rel->parallel = parallel;
 
                                if (!r[*pos])
                                        return rel;
@@ -2355,12 +2381,17 @@ rel_read(mvc *sql, char *r, int *pos, li
                        skipWS(r, pos);
                        if (!(exps = read_exps(sql, nrel, NULL, NULL, r, pos, 
'[', 0, 1)))
                                return NULL;
-                       rel = rel_topn(sql->sa, nrel, exps);
-                       set_processed(rel);
+                       rel = rel_topn(sql->sa, nrel, exps); set_processed(rel);
                }
                break;
        case 'p':
-               *pos += (int) strlen("project");
+               /* partition missing */
+               if (strncmp(r+*pos, "project", 7) == 0) {
+                       *pos += (int) strlen("project");
+               } else if (strncmp(r+*pos, "probe", 5) == 0) {
+                       *pos += (int) strlen("probe");
+                       probe = true;
+               }
                skipWS(r, pos);
 
                if (r[*pos] != '(')
@@ -2379,11 +2410,44 @@ rel_read(mvc *sql, char *r, int *pos, li
                if (!(exps = read_exps(sql, is_modify?nrel->l : nrel, NULL, 
NULL, r, pos, '[', 0, 1)))
                        return NULL;
                rel = rel_project(sql->sa, nrel, exps);
+               if (probe)
+                       rel->op = op_probehash;
                set_processed(rel);
                /* order by ? */
                /* first projected expressions, then left relation projections 
*/
-               if (r[*pos] == '[' && !(rel->r = read_exps(sql, rel, nrel, 
NULL, r, pos, '[', 0, 1)))
+               if (r[*pos] == '[' && !(rel->r = read_exps(sql, 
probe?rel->l:rel, probe?NULL:nrel, NULL, r, pos, '[', 0, 1)))
                        return NULL;
+               if (probe) {
+                       rel->attr = rel->exps;
+                       rel->exps = rel->r;
+                       rel->r = NULL;
+               }
+               break;
+       case 'b':
+               if (strncmp(r+*pos, "buildhash", 9) == 0)
+                       *pos += (int) strlen("buildhash");
+               skipWS(r, pos);
+
+               if (r[*pos] != '(')
+                       return sql_error(sql, -1, SQLSTATE(42000) "Project: 
missing '('\n");
+               (*pos)++;
+               skipWS(r, pos);
+               if (!(nrel = rel_read(sql, r, pos, refs)))
+                       return NULL;
+               skipWS(r, pos);
+               if (r[*pos] != ')')
+                       return sql_error(sql, -1, SQLSTATE(42000) "Project: 
missing ')'\n");
+               (*pos)++;
+               skipWS(r, pos);
+
+               if (!(attr = read_exps(sql, nrel, NULL, NULL, r, pos, '[', 0, 
1)))
+                       return NULL;
+               if (!(exps = read_exps(sql, nrel, NULL, NULL, r, pos, '[', 0, 
1)))
+                       return NULL;
+               rel = rel_project(sql->sa, nrel, exps);
+               rel->op = op_buildhash;
+               rel->attr = attr;
+               set_processed(rel);
                break;
        case 's':
        case 'a':
@@ -2558,6 +2622,8 @@ rel_read(mvc *sql, char *r, int *pos, li
                if (!(exps = read_exps(sql, lrel, rrel, NULL, r, pos, '[', 0, 
1)))
                        return NULL;
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to