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]