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

Reply via email to