astrodatagui: CalcAllFlaresThread: split list of dictionaries

previously we had a lot of overhead spawning new processes for
every single list entry.
Splitting it this way only spawns <core count> processes once,
heaving no continuous overhead.
This commit is contained in:
2024-05-15 10:59:58 +02:00
parent 724ce3e13f
commit 38aa08fcd0
+24 -25
View File
@@ -6,24 +6,27 @@ from ..flaredetector.flaredetector import calculateFlareFitsForLightcurve
import time
def getFlareCount(filesDict):
lc = lk.read(filesDict["FilePath"])
lc.flux = lc["sap_flux"]
lc_flattenend = lc.flatten()
sapPeaks, sapFits = calculateFlareFitsForLightcurve(lc_flattenend)
lc.flux = lc["pdcsap_flux"]
lc_flattenend = lc.flatten()
pdcsapPeaks, pdcsapFits = calculateFlareFitsForLightcurve(lc_flattenend)
def getFlareCount(filesDictList):
retDictList = []
for filesDict in filesDictList:
lc = lk.read(filesDict["FilePath"])
lc.flux = lc["sap_flux"]
lc_flattenend = lc.flatten()
sapPeaks, sapFits = calculateFlareFitsForLightcurve(lc_flattenend)
lc.flux = lc["pdcsap_flux"]
lc_flattenend = lc.flatten()
pdcsapPeaks, pdcsapFits = calculateFlareFitsForLightcurve(lc_flattenend)
retDict = filesDict
retDict["sapPeaks"] = sapPeaks
retDict["sapPeaksCount"] = len(sapPeaks)
retDict["sapFits"] = sapFits
retDict["pdcsapPeaks"] = pdcsapPeaks
retDict["pdcsapPeaksCount"] = len(pdcsapPeaks)
retDict["pdcsapFits"] = pdcsapFits
retDict = filesDict
retDict["sapPeaks"] = sapPeaks
retDict["sapPeaksCount"] = len(sapPeaks)
retDict["sapFits"] = sapFits
retDict["pdcsapPeaks"] = pdcsapPeaks
retDict["pdcsapPeaksCount"] = len(pdcsapPeaks)
retDict["pdcsapFits"] = pdcsapFits
retDictList.append(retDict)
return retDict
return retDictList
class CalcAllFlaresThread(QThread):
progress = pyqtSignal(int)
@@ -36,18 +39,14 @@ class CalcAllFlaresThread(QThread):
#@pyqtSlot()
def run(self):
cpuCount = multiprocessing.cpu_count()
print(len(self.allFlaresDictList))
start = time.time()
splitList = [self.allFlaresDictList[i:i + cpuCount] for i in range(0, len(self.allFlaresDictList), cpuCount)]
executor = concurrent.futures.ProcessPoolExecutor(cpuCount)
futures = [executor.submit(getFlareCount, starDict) for starDict in self.allFlaresDictList]
futures = [executor.submit(getFlareCount, starDictPartList) for starDictPartList in splitList]
concurrent.futures.wait(futures)
end = time.time()
print("Multiprocess time needed: ", end-start)
ret = []
#for dic in executor.map(getFlareCount, self.allFlaresDictList):
# print(dic)
# ret.append(dic)
for future in futures:
ret.append(future.result())
retList = future.result()
for dic in retList:
ret.append(dic)
self.finished.emit(ret)