Changeset: 331d0251ab40 for MonetDB
URL: http://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=331d0251ab40
Modified Files:
        MonetDB5/src/optimizer/opt_mapreduce.mx
Branch: default
Log Message:

halfway through rewrite, needto know final signatures before calling 
MRdistributeWork, such that type fixing isn't necessary and won't require bats 
for every column


diffs (218 lines):

diff -r bac195d85dd7 -r 331d0251ab40 MonetDB5/src/optimizer/opt_mapreduce.mx
--- a/MonetDB5/src/optimizer/opt_mapreduce.mx   Sat May 29 17:15:56 2010 +0200
+++ b/MonetDB5/src/optimizer/opt_mapreduce.mx   Mon May 31 08:40:34 2010 +0200
@@ -148,12 +148,13 @@
 
 
 typedef struct _mapcol {
-       int val1;
-       int val1type;
-       int val2;
-       int val2type;
-       int val3;
-       int val3type;
+       int mapid;        /* original column var in map program we
+                            eventually need to replace */
+       int reduceid;     /* var in reduce plan that is in its signature
+                            and return */
+       int type;         /* type of the map plan var */
+       int mapbat;       /* the var that is a BAT containing all values
+                            returned from map nodes (function) */
        struct _mapcol *next;
 } mapcol;
 
@@ -187,19 +188,31 @@
        w = (int *)alloca(retc * sizeof(int));
 
        for (lcol = col, j = 0; lcol != NULL; lcol = lcol->next, j++) {
+               /* define and create the container bat for all results from the
+                * map nodes */
                packs[j] = p = newFcnCall(reduce, batRef, newRef);
-               p = pushType(reduce, p, getHeadType(lcol->val1type));
-               p = pushType(reduce, p, getTailType(lcol->val1type));
-               getArg(p, 0) = lcol->val1;
-               setArgType(reduce, p, 0, lcol->val1type);
+               if (isaBatType(lcol->type)) {
+                       p = pushType(reduce, p, getHeadType(lcol->type));
+                       p = pushType(reduce, p, getTailType(lcol->type));
+               } else {
+                       p = pushNil(reduce, p, TYPE_void);
+                       p = pushType(reduce, p, lcol->type);
+               }
+               setArgType(reduce, p, 0, lcol->type);
+               lcol->mapbat = getArg(p, 0);
 
-               /* same for all sub results that we push into the mat.pack as
-                * arguments at the same time */
+               /* we need to declare the variables that we will use with put,
+                * exec and get */
                for (i = 0; i < n; i++) {
-                       p = newFcnCall(reduce, batRef, newRef);
-                       p = pushType(reduce, p, getHeadType(lcol->val1type));
-                       p = pushType(reduce, p, getTailType(lcol->val1type));
-                       setArgType(reduce, p, 0, lcol->val1type);
+                       if (isaBatType(lcol->type)) {
+                               p = newFcnCall(reduce, batRef, newRef);
+                               p = pushType(reduce, p, 
getHeadType(lcol->type));
+                               p = pushType(reduce, p, 
getTailType(lcol->type));
+                       } else {
+                               p = newAssignment(reduce);
+                               p = pushNil(reduce, p, lcol->type);
+                       }
+                       setArgType(reduce, p, 0, lcol->type);
                        gets[(i * retc) + j] = getArg(p, 0);
                }
        }
@@ -250,14 +263,14 @@
                p = pushArgument(reduce, p, q);
        }
 
-       /* delayed bat.inserts for easily creating an deterministic flow */
+       /* delayed bat.inserts for easily creating a deterministic flow */
        for (lcol = col, j = 0; lcol != NULL; lcol = lcol->next, j++) {
                q = getArg(packs[j], 0);
                /* p := bat.insert(b, y1) */
                for (i = 0; i < n; i++) {
                        p = newStmt(reduce, batRef, insertRef);
                        p = pushArgument(reduce, p, q);
-                       if (!isaBatType(lcol->val1type))
+                       if (!isaBatType(lcol->type))
                                p = pushNil(reduce, p, TYPE_void);
                        p = pushArgument(reduce, p, gets[(i * retc) + j]);
                        q = getArg(p, 0);
@@ -266,7 +279,7 @@
                 * confused by possible duplicate ids */
                p = newFcnCall(reduce, algebraRef, markHRef);
                p = pushArgument(reduce, p, q);
-               getArg(p, 0) = lcol->val1;
+               getArg(p, 0) = lcol->mapbat;
        }
 
 #if defined(_DEBUG_OPT_MAPREDUCE) && _DEBUG_OPT_MAPREDUCE == 0
@@ -577,9 +590,10 @@
                                                                        
alloca(sizeof(mapcol));
                                                                copy = SINGLE;
                                                        }
-                                                       lastcol->val1 = 
getArg(p, 0);
-                                                       lastcol->val1type = 
getArgType(map, p, 0);
-                                                       lastcol->val3 = 
lastcol->val1;
+                                                       lastcol->map = 
getArg(p, 0);
+                                                       lastcol->reduce = 
getarg(oreduce[i], 0);
+                                                       lastcol->type = 
getArgType(map, p, 0);
+                                                       lastcol->mapbat = -1;
                                                        lastcol->next = NULL;
                                                        newComment(map, "= sql 
column bat");
 
@@ -650,8 +664,8 @@
                         * 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->val3);
-                       map->stmt[0] = sig;
+                               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 */
@@ -675,13 +689,13 @@
                {
                        /* simple ORDER BY [DESC] or LIMIT/OFFSET */
                        for (lastcol = col; lastcol != NULL; lastcol = 
lastcol->next) { 
-                               if (getArg(p, 1) == lastcol->val1) {
+                               if (getArg(p, 1) == lastcol->reduce) {
                                        copy = SINGLE_DUP;
                                        newComment(map, "ORDER BY [DESC] or 
LIMIT/OFFSET");
                                        newComment(reduce, "ORDER BY [DESC] or 
LIMIT/OFFSET");
                                        /* fix return */
-                                       lastcol->val3 = getArg(omap[i], 0);
-                                       lastcol->val1 = getArg(p, 0);
+                                       lastcol->map = getArg(omap[i], 0);
+                                       lastcol->reduce = getArg(p, 0);
                                        break;
                                }
                        }
@@ -689,28 +703,16 @@
                                        getFunctionId(p) == maxRef ||
                                        getFunctionId(p) == minRef))
                {
-                       /* MAX/MIN aggregation, we cannot return a single val, 
since
-                        * we already fixed the signature, so create a 
container BAT */
+                       /* MAX/MIN aggregation */
                        for (lastcol = col; lastcol != NULL; lastcol = 
lastcol->next) { 
-                               if (getArg(p, 1) == lastcol->val1) {
+                               if (getArg(p, 1) == lastcol->reduce) {
                                        int t;
                                        newComment(map, "MAX/MIN");
                                        newComment(reduce, "MAX/MIN");
                                        /* basically perform a MAX over all 
MAXes */
                                        pushInstruction(map, omap[i]);
-                                       lastcol->val2 = getArg(omap[i], 0);
-                                       /* BAT for return */
-                                       p = newFcnCall(map, batRef, newRef);
-                                       p = pushType(map, p, 
getHeadType(lastcol->val1type));
-                                       p = pushType(map, p, 
getTailType(lastcol->val1type));
-                                       t = getArg(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->val2);
-                                       /* fix return */
-                                       lastcol->val1 = getArg(p, 0);
+                                       lastcol->map = getArg(omap[i], 0);
+                                       lastcol->reduce = getArg(oreduce[i], 0);
                                        /* we already DUP_SINGLE'd */
                                        copy = STICK;
                                        p = oreduce[i];
@@ -744,12 +746,12 @@
                                                getArg(p, 1) = getArg(oreduce[i 
- 1], 1);
                        }
                        for (lastcol = col; lastcol != NULL; lastcol = 
lastcol->next) { 
-                               if (getArg(p, 1) == lastcol->val1) {
+                               if (getArg(p, 1) == lastcol->reduce) {
                                        int t, tt;
                                        newComment(map, "COUNT/SUM");
                                        newComment(reduce, "COUNT/SUM");
                                        pushInstruction(map, omap[i]);
-                                       lastcol->val2 = getArg(omap[i], 0);
+                                       lastcol->map = getArg(omap[i], 0);
                                        /* BAT for return */
                                        p = newFcnCall(map, batRef, newRef);
                                        p = pushType(map, p, 
getHeadType(lastcol->val1type));
@@ -763,8 +765,8 @@
                                        p = newFcnCall(map, batRef, insertRef);
                                        p = pushArgument(map, p, t);
                                        p = pushNil(map, p, 
getHeadType(lastcol->val1type));
-                                       p = pushArgument(map, p, lastcol->val2);
-                                       t = getArg(p, 0);
+                                       p = pushArgument(map, p, lastcol->map);
+                                       lastcol->map = t = getArg(p, 0);
 
                                        /* propagate type change */
                                        for (j = 0; j < sig->retc; j++) {
@@ -830,10 +832,13 @@
                        }
                } else if (getModuleId(p) == calcRef && getFunctionId(p) == 
divRef) {
                        for (lastcol = col; lastcol != NULL; lastcol = 
lastcol->next) {
-                               if (getArg(p, 1) == lastcol->val1 ||
-                                               getArg(p, 2) == lastcol->val1)
+                               if (getArg(p, 1) == lastcol->reduce ||
+                                               getArg(p, 2) == lastcol->reduce)
                                {
-                                       /* TODO */
+                                       p = newFcnCall(reduce, aggrRef, sumRef);
+                                       p = pushArgument(reduce, p, 
lastcol->mapbat);
+                                       getArg(p, 0) = lastcol->reduce;
+                                       copy = STICK;
                                }
                        }
                }
@@ -845,6 +850,8 @@
                         * changing the actual return variable */
                        ret = newInstruction(map, ASSIGNsymbol);
                        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) {  
_______________________________________________
Checkin-list mailing list
[email protected]
http://mail.monetdb.org/mailman/listinfo/checkin-list

Reply via email to