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]

Reply via email to