Changeset: d545d1b3a8ba for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB/rev/d545d1b3a8ba
Modified Files:
monetdb5/mal/mal.h
monetdb5/mal/mal_client.c
monetdb5/mal/mal_client.h
monetdb5/mal/mal_dataflow.c
monetdb5/mal/mal_instruction.c
monetdb5/mal/mal_interpreter.c
monetdb5/mal/mal_resource.c
monetdb5/mal/mal_runtime.c
monetdb5/modules/mal/sysmon.c
Branch: default
Log Message:
Move maintenance of number of threads working for a query to query context.
Also use atomic variables for this and the maximum number of workers
that worked on a query (which is maintained in the MalBlk record).
diffs (237 lines):
diff --git a/monetdb5/mal/mal.h b/monetdb5/mal/mal.h
--- a/monetdb5/mal/mal.h
+++ b/monetdb5/mal/mal.h
@@ -185,7 +185,7 @@ typedef struct MALBLK {
ptr replica; /* for the replicator tests */
/* During the run we keep track on the maximum number of concurrent
threads and memory claim */
- int workers;
+ ATOMIC_TYPE workers;
lng memory;
lng runtime; /* average execution time of block in
ticks */
int calls; /* number of calls */
@@ -222,7 +222,6 @@ typedef struct MALSTK {
char status; /* srunning 'R' suspended 'S', quiting
'Q' */
int pcup; /* saved pc upon a recursive
all */
oid tag; /* unique invocation call tag */
- int workers; /* Actual number of concurrent
workers */
lng memory; /* Actual memory claims for
highwater mark */
struct MALSTK *up; /* stack trace list */
diff --git a/monetdb5/mal/mal_client.c b/monetdb5/mal/mal_client.c
--- a/monetdb5/mal/mal_client.c
+++ b/monetdb5/mal/mal_client.c
@@ -61,6 +61,7 @@ mal_client_reset(void)
if (mal_clients) {
for (int i = 0; i < MAL_MAXCLIENTS; i++) {
ATOMIC_DESTROY(&mal_clients[i].lastprint);
+ ATOMIC_DESTROY(&mal_clients[i].workers);
ATOMIC_DESTROY(&mal_clients[i].qryctx.datasize);
}
GDKfree(mal_clients);
@@ -93,6 +94,7 @@ MCinit(void)
}
for (int i = 0; i < MAL_MAXCLIENTS; i++){
ATOMIC_INIT(&mal_clients[i].lastprint, 0);
+ ATOMIC_INIT(&mal_clients[i].workers, 1);
ATOMIC_INIT(&mal_clients[i].qryctx.datasize, 0);
mal_clients[i].idx = -1; /* indicate it's available */
}
diff --git a/monetdb5/mal/mal_client.h b/monetdb5/mal/mal_client.h
--- a/monetdb5/mal/mal_client.h
+++ b/monetdb5/mal/mal_client.h
@@ -68,7 +68,7 @@ typedef struct CLIENT {
*/
char optimizer[IDLENGTH];/* The optimizer pipe preferred for this
session */
int workerlimit; /* maximum number of workthreads
processing a query */
- int memorylimit; /* Memory claim highwater mark,
0 = no limit */
+ int memorylimit; /* maximum memory currently
allowed in MB */
lng maxmem; /* maximum memory from
db_user_info table */
lng sessiontimeout; /* session abort after x usec,
0 = no limit */
QryCtx qryctx; /* per query limitations */
@@ -88,6 +88,7 @@ typedef struct CLIENT {
BAT *profevents;
ATOMIC_TYPE lastprint; /* when we last printed the query, to
be deprecated */
+ ATOMIC_TYPE workers; /* number of threads working for this
context */
/*
* Communication channels for the interconnect are stored here.
* It is perfectly legal to have a client without input stream.
diff --git a/monetdb5/mal/mal_dataflow.c b/monetdb5/mal/mal_dataflow.c
--- a/monetdb5/mal/mal_dataflow.c
+++ b/monetdb5/mal/mal_dataflow.c
@@ -50,7 +50,7 @@ typedef struct FLOWEVENT {
sht cost;
lng hotclaim; /* memory foot print of result variables */
lng argclaim; /* memory foot print of arguments */
- lng maxclaim; /* memory foot print of largest argument, counld be
used to indicate result size */
+ lng maxclaim; /* memory foot print of largest argument, could be used
to indicate result size */
} *FlowEvent, FlowEventRec;
typedef struct queue {
@@ -375,15 +375,25 @@ DFLOWworker(void *T)
if( p->fcn != (MALfcn) deblockdataflow){
fe->hotclaim = 0; /* don't assume
priority anymore */
fe->maxclaim = 0;
- if (todo->last == 0)
+ MT_lock_set(&todo->l);
+ int last = todo->last;
+ MT_lock_unset(&todo->l);
+ if (last == 0)
MT_sleep_ms(DELAYUNIT);
q_requeue(todo, fe);
continue;
}
}
+ ATOMIC_BASE_TYPE wrks =
ATOMIC_INC(&flow->cntxt->workers);
+ ATOMIC_BASE_TYPE mwrks = ATOMIC_GET(&flow->mb->workers);
+ while (wrks > mwrks) {
+ if (ATOMIC_CAS(&flow->mb->workers, &mwrks,
wrks))
+ break;
+ }
error = runMALsequence(flow->cntxt, flow->mb, fe->pc,
fe->pc + 1, flow->stk, 0, 0);
+ (void) ATOMIC_DEC(&flow->cntxt->workers);
/* release the memory claim */
- MALadmission_release(flow->cntxt, flow->mb, flow->stk,
p, claim);
+ MALadmission_release(flow->cntxt, flow->mb, flow->stk,
p, claim);
MT_lock_set(&flow->flowlock);
fe->state = DFLOWwrapup;
@@ -724,12 +734,14 @@ DFLOWscheduler(DataFlow flow, struct wor
/* initialize the eligible statements */
fe = flow->status;
+ (void) ATOMIC_DEC(&flow->cntxt->workers);
MT_lock_set(&flow->flowlock);
for (i = 0; i < actions; i++)
if (fe[i].blocks == 0) {
p = getInstrPtr(flow->mb,fe[i].pc);
if (p == NULL) {
MT_lock_unset(&flow->flowlock);
+ (void) ATOMIC_INC(&flow->cntxt->workers);
throw(MAL, "dataflow", "DFLOWscheduler():
getInstrPtr(flow->mb,fe[i].pc) returned NULL");
}
fe[i].argclaim = 0;
@@ -745,8 +757,10 @@ DFLOWscheduler(DataFlow flow, struct wor
f = q_dequeue(flow->done, NULL);
if (ATOMIC_GET(&exiting))
break;
- if (f == NULL)
+ if (f == NULL) {
+ (void) ATOMIC_INC(&flow->cntxt->workers);
throw(MAL, "dataflow", "DFLOWscheduler():
q_dequeue(flow->done) returned NULL");
+ }
/*
* When an instruction is finished we have to reduce the blocked
@@ -772,6 +786,7 @@ DFLOWscheduler(DataFlow flow, struct wor
/* release the worker from its specific task (turn it into a
* generic worker) */
ATOMIC_PTR_SET(&w->cntxt, NULL);
+ (void) ATOMIC_INC(&flow->cntxt->workers);
/* wrap up errors */
assert(flow->done->last == 0);
if ((ret = ATOMIC_PTR_XCG(&flow->error, NULL)) != NULL ) {
diff --git a/monetdb5/mal/mal_instruction.c b/monetdb5/mal/mal_instruction.c
--- a/monetdb5/mal/mal_instruction.c
+++ b/monetdb5/mal/mal_instruction.c
@@ -131,6 +131,7 @@ newMalBlk(int elements)
GDKfree(mb);
return NULL;
}
+ ATOMIC_INIT(&mb->workers, 1);
return mb;
}
@@ -267,7 +268,6 @@ freeMalBlk(MalBlkPtr mb)
mb->binding[0] = 0;
mb->tag = 0;
mb->memory = 0;
- mb->workers = 0;
if (mb->help && mb->statichelp != mb->help)
GDKfree(mb->help);
mb->help = 0;
@@ -275,6 +275,7 @@ freeMalBlk(MalBlkPtr mb)
mb->inlineProp = 0;
mb->unsafeProp = 0;
freeException(mb->errors);
+ ATOMIC_DESTROY(&mb->workers);
GDKfree(mb);
}
diff --git a/monetdb5/mal/mal_interpreter.c b/monetdb5/mal/mal_interpreter.c
--- a/monetdb5/mal/mal_interpreter.c
+++ b/monetdb5/mal/mal_interpreter.c
@@ -285,7 +285,6 @@ prepareMALstack(MalBlkPtr mb, int size)
return NULL;
stk->stktop = mb->vtop;
stk->blk = mb;
- stk->workers = 0;
stk->memory = 0;
initStack(0, res);
if(!res) {
diff --git a/monetdb5/mal/mal_resource.c b/monetdb5/mal/mal_resource.c
--- a/monetdb5/mal/mal_resource.c
+++ b/monetdb5/mal/mal_resource.c
@@ -123,7 +123,7 @@ MALadmission_claim(Client cntxt, MalBlkP
* A way out is to attach the thread count to the MAL stacks, which
just limits the level
* of parallism for a single dataflow graph.
*/
- if(cntxt->workerlimit && cntxt->workerlimit < stk->workers){
+ if (cntxt->workerlimit && (int) ATOMIC_GET(&cntxt->workers) >=
cntxt->workerlimit) {
return -1;
}
if (argclaim == 0)
@@ -137,7 +137,7 @@ MALadmission_claim(Client cntxt, MalBlkP
}
/* the argument claim is based on the input for an instruction */
- if ( memorypool > argclaim || stk->workers == 0 ) {
+ if ( memorypool > argclaim ) {
/* If we are low on memory resources, limit the user if he
exceeds his memory budget
* but make sure there is at least one worker thread active */
if ( cntxt->memorylimit) {
@@ -148,11 +148,8 @@ MALadmission_claim(Client cntxt, MalBlkP
stk->memory += argclaim;
}
memorypool -= argclaim;
- stk->workers++;
stk->memory += argclaim;
MT_lock_set(&mal_delayLock);
- if( mb->workers < stk->workers)
- mb->workers = stk->workers;
if( mb->memory < stk->memory)
mb->memory = stk->memory;
MT_lock_unset(&mal_delayLock);
@@ -181,7 +178,6 @@ MALadmission_release(Client cntxt, MalBl
if ( memorypool > (lng) MEMORY_THRESHOLD ){
memorypool = (lng) MEMORY_THRESHOLD;
}
- stk->workers--;
stk->memory -= argclaim;
MT_lock_unset(&admissionLock);
return;
diff --git a/monetdb5/mal/mal_runtime.c b/monetdb5/mal/mal_runtime.c
--- a/monetdb5/mal/mal_runtime.c
+++ b/monetdb5/mal/mal_runtime.c
@@ -283,7 +283,7 @@ runtimeProfileFinish(Client cntxt, MalBl
if (QRYqueue[i].stk == stk) {
QRYqueue[i].status = "finished";
QRYqueue[i].finished = time(0);
- QRYqueue[i].workers = mb->workers;
+ QRYqueue[i].workers = (int) ATOMIC_GET(&mb->workers);
/* give the MB upperbound by addition of 1 MB */
QRYqueue[i].memory = 1 + (int)(mb->memory /
LL_CONSTANT(1048576));
QRYqueue[i].cntxt = NULL;
diff --git a/monetdb5/modules/mal/sysmon.c b/monetdb5/modules/mal/sysmon.c
--- a/monetdb5/modules/mal/sysmon.c
+++ b/monetdb5/modules/mal/sysmon.c
@@ -241,7 +241,7 @@ SYSMONqueue(Client cntxt, MalBlkPtr mb,
goto bailout;
if( QRYqueue[i].mb)
- wrk = QRYqueue[i].mb->workers;
+ wrk = (int)
ATOMIC_GET(&QRYqueue[i].mb->workers);
else
wrk = QRYqueue[i].workers;
if( QRYqueue[i].mb)
_______________________________________________
checkin-list mailing list -- [email protected]
To unsubscribe send an email to [email protected]