Version 2 - per-producer queues:
The original design used one shared array queue that every producer kernel wrote into and the trainer read from, which meant every kernel touched the same underlying data structure and needed offset bookkeeping (trimmed, takenSoFar) to avoid stopping on each other — the source of the occasional duplicate-delivery race.
​
Version 2 solves that at the structural level:
each producer gets its own named queue variable and a single boolean ready flag instead of an offset-tracked array slice. There’s no more shared array to trim or index into — the trainer just checks each producer’s flag and swaps it off atomically once it’s consumed the batch. This removes the trim-race entirely and makes the flow easier to reason about, but it still runs through WL’s parallel-kernel shared-variable mechanism, and that mechanism has its own coordination floor: every check of ready[[cI]] or activeProducers is a round-trip through the kernel link’s shared-variable protocol. That overhead shows up directly in the numbers — batch retrieval took ~0.133s on average even though the actual generation cost per chunk was ~0.57s, so a meaningful fraction of the trainer’s time was spent waiting on shared-variable synchronization rather than on real work.
​
Version 3 — raw WSTP links:
This version drops ParallelSubmit/SetSharedVariable coordination altogether and talks to producers as plain WSTP subprocesses launched with LinkLaunch. Each producer is a dedicated point-to-point link: the trainer writes produceChunk[batchSize] directly onto a producer’s link and polls that link with LinkReadyQ/LinkRead for its result, with no shared kernel state involved in the handoff.
​
That eliminates the synchronization layer that Version 2 was still paying for. The result is visible directly in the batch-wait measurements — average wait dropped from ~0.133s (V2) to ~0.0025s (V3), roughly a 50x reduction — while the actual per-chunk generation cost stayed effectively the same (~0.53–0.57s in both). For the full 128-round run this took total wall time from ~45.2s down to ~31.1s. Scaled to a real ~12-hour training run, closing that per-batch gap.
For a typical 400 round training with around 128 batch retreivals that normally runs for 12 hours this shaves of around 2 hours of time.

Version 2: separate queues for each kernel

In[]:=
batchSize=2;​​rounds=128;​​stallTimeout=4;​​​​producerStatusQ[]:=(​​Dynamic[Grid[{​​{Style["Active producers: "<>ToString[activeProducers],Bold],SpanFromLeft},​​Join[{Style["Produced:",Bold],Total[produced]," - "},produced],​​Join[{Style["Used:",Bold],Total[used]," - "},used],​​Join[{Style["Ready:",Bold],Total[ready]," - "},ready]/.{False->0,True->1}​​},Alignment->Center]]​​);​​​​CloseKernels[];​​LaunchKernels[];​​nProducers=activeProducers=Length[Kernels[]]-1;​​​​queue="queue"<>ToString[#]&;​​trainingDone=False;​​queueVars=Table[queue[i],{i,nProducers}];​​Clear/@queueVars;​​Table[ToExpression[queue[i]<>" = {}"],{i,nProducers}];​​ready=Table[False,nProducers];​​produced=used=Table[0,nProducers];​​​​produceChunk[n_]:={​​$ProcessID,​​AbsoluteTiming[Table[​​Pause[RandomReal[{0.2,.3}]];​​NumericArray[RandomReal[{0,1},{32,96,96}],"Real32"]->​​NumericArray[RandomInteger[{0,1},{32,21,96,96}],"Byte"]​​,{n}]]​​};​​​​cI=0;​​getBatchQ[]:=Block[{q,bs,startTime,offset,attempt},​​startTime=AbsoluteTime[];​​Catch[While[True,Do[​​cI=Mod[cI,nProducers]+1;​​If[ready[[cI]],​​q=ToExpression[queue[cI]];​​used[[cI]]++;​​ready[[cI]]=False;​​Throw[q]];​​,{nProducers}];​​If[activeProducers<=0,Echo["No more generators"];Abort[]];​​If[AbsoluteTime[]-startTime>stallTimeout,Echo["Could not get enough samples in 4s"];Abort[]];​​Pause[0.1]​​]]​​];​​​​SetSharedVariable[trainingDone,activeProducers,used,produced,ready];​​ToExpression["SetSharedVariable["<>StringRiffle[queueVars,","]<>"]"];​​DistributeDefinitions[produceChunk,getBatchQ,batchSize,rounds,queue,nProducers,stallTimeout];​​​​producerJob=Table[With[{qi=i},Block[{},​​ParallelSubmit[​​While[!trainingDone,​​If[ready[[qi]],Pause[0.1],​​ToExpression[queue[qi]<>" = produceChunk[batchSize]"];​​produced[[qi]]++;​​ready[[qi]]=True​​]];​​activeProducers--;​​]]],{i,nProducers}];​​​​trainerJob=ParallelSubmit[CheckAbort[Block[{t0,times,out,bt},​​times={};​​out=Table[​​t0=AbsoluteTime[];​​bt=getBatchQ[];​​AppendTo[times,AbsoluteTime[]-t0];​​Pause[.2];​​bt,{rounds}];​​trainingDone=True;​​{times,out}​​],​​(trainingDone=True;$Aborted)]];​​​​Dynamic[producerStatusQ[]]​​AbsoluteTiming[​​out=First@WaitAll[{trainerJob,producerJob},ProgressReporting->False];​​]​​​​SmoothHistogram[First@out,PlotRange->{{0.,0.3},Full}]​​Through[{Mean,StandardDeviation}[First@out]]​​Through[{Mean,StandardDeviation}[out[[2,All,2,1]]]]​​DeleteDuplicates[out[[2,All,1]]]
Out[]=
producerStatusQ[]
Out[]=
{45.2216,Null}
Out[]=
Out[]=
{0.1331397,0.0736275}
Out[]=
{0.566174,0.0424185}
Out[]=
{46956,42420,64220,42396,21616,54100,48672,44292,42512,32604,38208,60644,58596,57532,54012}

Version3: Using Links

In[]:=
batchSize=2;​​rounds=128;​​stallTimeout=4;​​nProducers=activeProducers=15;​​​​producerStatusL[]:=(​​produced=used=Table[0,nProducers];​​Dynamic[Grid[{​​{Style["Active producers: "<>ToString[activeProducers],Bold],SpanFromLeft},​​Join[{Style["Produced:",Bold],Total[produced]," - "},produced],​​Join[{Style["Used:",Bold],Total[used]," - "},used]​​},Alignment->Center]]​​);​​​​readLink[link_,timeout_:4]:=Module[{pkt,t0=AbsoluteTime[]},​​While[True,​​If[AbsoluteTime[]-t0>timeout,Return[$Failed]];​​If[LinkReadyQ[link],​​pkt=Check[LinkRead[link],$Failed];​​Which[pkt===$Failed,​​Return[$Failed],​​Head[pkt]===ReturnPacket&&First[pkt]=!=Null,​​Return[First[pkt]],​​True,​​Null​​],Pause[0.02]]​​]];​​​​writeProducer[link_]:=LinkWrite[link,Unevaluated[produceChunk[n_]:={​​$ProcessID,​​AbsoluteTiming[Table[​​Pause[RandomReal[{0.2,.3}]];​​NumericArray[RandomReal[{0,1},{32,96,96}],"Real32"]->​​NumericArray[RandomInteger[{0,1},{32,1,96,96}],"Byte"]​​,{n}]]​​}]]​​​​activateProducer[link_,batch_]:=With[{b=batch},​​LinkWrite[link,Unevaluated[produceChunk[b]]];​​];​​​​getBatchL[]:=Block[{startTime=AbsoluteTime[],result},Catch[While[True,Do[​​cI=Mod[cI,nProducers]+1;​​link=producerLinks[[cI]];​​If[LinkReadyQ[link],​​result=readLink[link];​​If[result=!=$Failed,​​used[[cI]]++;​​activateProducer[link,batchSize];​​produced〚cI〛++;​​Throw[{cI,result}],​​activeProducers--;​​Echo["producer "<>ToString[cI]<>" failed"];​​]];​​,{nProducers}];​​If[activeProducers<=0,Echo["No more generators"];Abort[]];​​If[AbsoluteTime[]-startTime>stallTimeout,Echo["stall"];Abort[]];​​Pause[0.05];​​]]];​​​​Dynamic[producerStatusL[]]​​​​AbsoluteTiming[​​producerLinks=Table[​​link=LinkLaunch[First[$CommandLine]<>" -wstp"];​​writeProducer[link];​​activateProducer[link,batchSize];​​produced〚i〛++;​​link​​,{i,nProducers}];​​​​times={};​​cI=0;​​out=Table[t0=AbsoluteTime[];​​bt=getBatchL[];​​AppendTo[times,AbsoluteTime[]-t0];​​Pause[.2];​​bt,{r,rounds}];​​​​LinkClose/@producerLinks;​​]​​​​SmoothHistogram[times,PlotRange->{{0,0.3},Full}]​​Through[{Mean,StandardDeviation}[times]]​​Through[{Mean,StandardDeviation}[out[[All,2,2,1]]]]​​DeleteDuplicates[out[[All,2,1]]]