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