Changeset: 4c5c629e6493 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=4c5c629e6493
Modified Files:
MonetDB5/src/optimizer/opt_mapreduce.mx
Branch: default
Log Message:
try to delay generating the reduce plan as to avoid having to rewrite types
diffs (274 lines):
diff -r 331d0251ab40 -r 4c5c629e6493 MonetDB5/src/optimizer/opt_mapreduce.mx
--- a/MonetDB5/src/optimizer/opt_mapreduce.mx Mon May 31 08:40:34 2010 +0200
+++ b/MonetDB5/src/optimizer/opt_mapreduce.mx Wed Jun 09 11:18:23 2010 +0200
@@ -389,15 +389,18 @@
* counter cases */
if (trackstack_contains(&avgtrack,
getArg(p, j + 1))) {
mapcol *sum, *count;
- /* got it, time to copy
instructions */
- lastcol->val1 = getArg(p, 1);
- lastcol->val1type =
getArgType(map, p, 1);
+ /* go from a single to two
columns */
+ lastcol->mapid = getArg(p, 1);
+ lastcol->reduceid =
getArg(oreduce[i], 1);
+ lastcol->type = getArgType(map,
p, 1);
sum = lastcol;
lastcol = lastcol->next =
alloca(sizeof(mapcol));
- lastcol->val1 = getArg(p, 2);
- lastcol->val1type =
getArgType(map, p, 2);
+ lastcol->mapid = getArg(p, 2);
+ lastcol->reduceid =
getArg(oreduce[i], 2);
+ lastcol->type = getArgType(map,
p, 2);
lastcol->next = NULL;
count = lastcol;
+ /* got it, time to copy
instructions */
j = i;
for (i = k + 1; i < j; i++) {
p = omap[i];
@@ -408,36 +411,6 @@
}
newComment(map, "= AVG
columns");
- /* BAT for return */
- p = newFcnCall(map, batRef,
newRef);
- p = pushNil(map, p, TYPE_void);
- p = pushType(map, p,
sum->val1type);
- sum->val3 = getArg(p, 0);
- /* bat.insert */
- p = newFcnCall(map, batRef,
insertRef);
- p = pushArgument(map, p,
sum->val3);
- p = pushNil(map, p, TYPE_void);
- p = pushArgument(map, p,
sum->val1);
- /* fix return */
- setArgType(map, p, 0,
-
newBatType(TYPE_void, sum->val1type));
- sum->val3 = getArg(p, 0);
-
- /* BAT for return */
- p = newFcnCall(map, batRef,
newRef);
- p = pushNil(map, p, TYPE_void);
- p = pushType(map, p,
count->val1type);
- count->val3 = getArg(p, 0);
- /* bat.insert */
- p = newFcnCall(map, batRef,
insertRef);
- p = pushArgument(map, p,
count->val3);
- p = pushNil(map, p, TYPE_void);
- p = pushArgument(map, p,
count->val1);
- /* fix return */
- setArgType(map, p, 0,
-
newBatType(TYPE_void, count->val1type));
- count->val3 = getArg(p, 0);
-
i = limit;
break;
}
@@ -562,9 +535,9 @@
}
}
- /* move over statement that depend (indirectly) on the sql.bind
+ /* move over statements that depend (indirectly) on the sql.bind
* calls */
- if (!trackstack_isempty(&tracker)) for (j = p->retc; j <
p->argc; j++) {
+ for (j = p->retc; j < p->argc; j++) {
if (trackstack_contains(&tracker, getArg(p, j))) {
if (getModuleId(p) == algebraRef) {
if (getFunctionId(p) == kunionRef) {
@@ -590,8 +563,8 @@
alloca(sizeof(mapcol));
copy = SINGLE;
}
- lastcol->map =
getArg(p, 0);
- lastcol->reduce =
getarg(oreduce[i], 0);
+ lastcol->mapid =
getArg(p, 0);
+ lastcol->reduceid =
getArg(oreduce[i], 0);
lastcol->type =
getArgType(map, p, 0);
lastcol->mapbat = -1;
lastcol->next = NULL;
@@ -642,7 +615,6 @@
break;
}
}
- trackstack_destroy(&tracker);
trackstack_destroy(&avgtrack);
if (hadBinds == 0) {
@@ -651,30 +623,17 @@
reduce->stop = limit;
GDKfree(omap);
freeSymbol(new);
+ trackstack_destroy(&tracker);
return 0;
}
+ trackstack_clear(&tracker);
copy = STICK;
for (i = 0; i < limit; i++) { /* phase 2 */
p = oreduce[i];
if (p == (InstrPtr)1) {
- newComment(reduce, "{ map");
- /* this is the moment the first sql.bind occurred, set
the
- * calling signature of the map program */
- getArg(sig, 0) = -1; /* get rid of default retval */
- for (lastcol = col; lastcol != NULL; lastcol =
lastcol->next)
- sig = pushReturn(map, sig, lastcol->map);
- map->stmt[0] = sig; /* many args realloc sig, so reset
it */
- MRdistributework(cntxt, reduce, col, sig, mrcluster);
- newComment(reduce, "} map");
- /* bogus instruction to be able to set to NOOP */
- oreduce[i] = newInstruction(reduce, ASSIGNsymbol);
- setModuleId(oreduce[i], ioRef);
- setFunctionId(oreduce[i], printRef);
- oreduce[i] = pushReturn(reduce, oreduce[i], 0);
- oreduce[i] = pushArgument(reduce, oreduce[i], 0);
- oreduce[i]->token = NOOPsymbol;
+ trackstack_push(&tracker, -1);
continue;
}
@@ -689,13 +648,13 @@
{
/* simple ORDER BY [DESC] or LIMIT/OFFSET */
for (lastcol = col; lastcol != NULL; lastcol =
lastcol->next) {
- if (getArg(p, 1) == lastcol->reduce) {
+ if (getArg(p, 1) == lastcol->reduceid) {
copy = SINGLE_DUP;
newComment(map, "ORDER BY [DESC] or
LIMIT/OFFSET");
newComment(reduce, "ORDER BY [DESC] or
LIMIT/OFFSET");
/* fix return */
- lastcol->map = getArg(omap[i], 0);
- lastcol->reduce = getArg(p, 0);
+ lastcol->mapid = getArg(omap[i], 0);
+ lastcol->reduceid = getArg(p, 0);
break;
}
}
@@ -705,14 +664,13 @@
{
/* MAX/MIN aggregation */
for (lastcol = col; lastcol != NULL; lastcol =
lastcol->next) {
- if (getArg(p, 1) == lastcol->reduce) {
- int t;
+ if (getArg(p, 1) == lastcol->reduceid) {
newComment(map, "MAX/MIN");
newComment(reduce, "MAX/MIN");
/* basically perform a MAX over all
MAXes */
pushInstruction(map, omap[i]);
- lastcol->map = getArg(omap[i], 0);
- lastcol->reduce = getArg(oreduce[i], 0);
+ lastcol->mapid = getArg(omap[i], 0);
+ lastcol->reduceid = getArg(oreduce[i],
0);
/* we already DUP_SINGLE'd */
copy = STICK;
p = oreduce[i];
@@ -746,28 +704,13 @@
getArg(p, 1) = getArg(oreduce[i
- 1], 1);
}
for (lastcol = col; lastcol != NULL; lastcol =
lastcol->next) {
- if (getArg(p, 1) == lastcol->reduce) {
- int t, tt;
+ if (getArg(p, 1) == lastcol->reduceid) {
newComment(map, "COUNT/SUM");
newComment(reduce, "COUNT/SUM");
pushInstruction(map, omap[i]);
- lastcol->map = getArg(omap[i], 0);
- /* BAT for return */
- p = newFcnCall(map, batRef, newRef);
- p = pushType(map, p,
getHeadType(lastcol->val1type));
- p = pushType(map, p, getArgType(map,
omap[i], 0));
- setArgType(map, p, 0,
-
newBatType(getHeadType(lastcol->val1type),
- getArgType(map,
omap[i], 0)));
- t = getArg(p, 0);
- tt = getArgType(map, p, 0);
- /* bat.insert */
- p = newFcnCall(map, batRef, insertRef);
- p = pushArgument(map, p, t);
- p = pushNil(map, p,
getHeadType(lastcol->val1type));
- p = pushArgument(map, p, lastcol->map);
- lastcol->map = t = getArg(p, 0);
+ lastcol->mapid = getArg(omap[i], 0);
+#if 0
/* propagate type change */
for (j = 0; j < sig->retc; j++) {
if (getArg(sig, j) ==
lastcol->val1) {
@@ -828,16 +771,17 @@
/* we injected an alternative
instruction */
copy = cNONE;
break;
+#endif
}
}
} else if (getModuleId(p) == calcRef && getFunctionId(p) ==
divRef) {
for (lastcol = col; lastcol != NULL; lastcol =
lastcol->next) {
- if (getArg(p, 1) == lastcol->reduce ||
- getArg(p, 2) == lastcol->reduce)
+ if (getArg(p, 1) == lastcol->reduceid ||
+ getArg(p, 2) ==
lastcol->reduceid)
{
p = newFcnCall(reduce, aggrRef, sumRef);
p = pushArgument(reduce, p,
lastcol->mapbat);
- getArg(p, 0) = lastcol->reduce;
+ getArg(p, 0) = lastcol->reduceid;
copy = STICK;
}
}
@@ -845,6 +789,7 @@
/* terminate both map and reduce functions properly */
if (p->token == ENDsymbol) {
+ size_t l;
/* make sure the return comes at the end, as we may have
* added some stuff to the MAP program in this phase,
* changing the actual return variable */
@@ -852,18 +797,35 @@
ret->barrier = RETURNsymbol;
/* TODO: this is the moment that we should call
MRdistribute
* work as well */
- if (col == NULL) {
- getArg(ret, 0) = newTmpVariable(map, TYPE_void);
- } else for (lastcol = col; lastcol != NULL; lastcol =
lastcol->next) {
- ret = pushReturn(map, ret, lastcol->val3);
+
+ /* nothing can change now any more, so finally set the
+ * calling signature of the map program */
+ getArg(sig, 0) = -1; /* get rid of default retval */
+ for (lastcol = col; lastcol != NULL; lastcol =
lastcol->next) {
+ sig = pushReturn(map, sig, lastcol->mapid);
+ ret = pushReturn(map, ret, lastcol->mapid);
}
+ map->stmt[0] = sig; /* many args realloc sig, so reset
it */
pushInstruction(map, ret);
+
+ for (l = 0; l < tracker.cur; l++) {
+ if (tracker.stack[l] == -1) {
+ newComment(reduce, "{ call-map");
+ MRdistributework(cntxt, reduce, col,
sig, mrcluster);
+ newComment(reduce, "} call-map");
+ oreduce[l] = NULL;
+ } else {
+ pushInstruction(reduce,
oreduce[tracker.stack[l]]);
+ }
+ }
+
copy = DUP;
}
switch (copy) {
case STICK:
- pushInstruction(reduce, p);
+ trackstack_push(&tracker, i);
+// pushInstruction(reduce, p);
break;
case SINGLE_DUP:
copy = STICK;
@@ -894,6 +856,7 @@
GDKfree(omap);
GDKfree(oreduce);
+ trackstack_destroy(&tracker);
insertSymbol(findModule(cntxt->nspace, userRef), new);
return 1;
_______________________________________________
Checkin-list mailing list
[email protected]
http://mail.monetdb.org/mailman/listinfo/checkin-list