diff --git a/Dockerfile b/Dockerfile index a938575..bb3835b 100644 --- a/Dockerfile +++ b/Dockerfile @@ -7,4 +7,8 @@ COPY ./target/scala-2.13/WeatherTool-assembly-0.1.1-SNAPSHOT.jar app.jar VOLUME /app/data -CMD ["java", "-jar", "app.jar"] \ No newline at end of file + +# flags setting heap to 256mb and max 384mb, plus g1 garbage collector +# CMD ["java", "-Xms256m", "-Xmx384m", "-XX:+UseG1GC", "-jar", "app.jar"] +# CMD ["java", "-jar", "app.jar"] +CMD ["java", "-Xms512m", "-Xmx600m", "-XX:+UseG1GC", "-XX:MaxGCPauseMillis=200", "-jar", "app.jar"] \ No newline at end of file diff --git a/fly.toml b/fly.toml index 3beb42e..d21e977 100644 --- a/fly.toml +++ b/fly.toml @@ -18,7 +18,7 @@ primary_region = "waw" size_gb = 10 [[vm]] - memory = 512 + memory = 1024 cpu_kind = 'shared' cpus = 1 diff --git a/src/main/scala/data/DataService.scala b/src/main/scala/data/DataService.scala index 0738093..e808163 100644 --- a/src/main/scala/data/DataService.scala +++ b/src/main/scala/data/DataService.scala @@ -64,29 +64,40 @@ class DataService(log: Logger[IO]) { def getAllFileStructure(): IO[List[Grib]] = { for { fileList <- getFileList() - allStructure <- fileList.parTraverseN(4)(getGribStucture) // limit concurrency to 4 + allStructure <- fileList.parTraverseN(4)(getGribStucture) // limit concurrency } yield allStructure.flatten } -// private val blockingEC = ExecutionContext.fromExecutorService( -// Executors.newFixedThreadPool(4) // limit concurrency to 4 -// ) -// -// private val semaphore = Semaphore[IO](4).unsafeRunSync() + private val semaphore = Semaphore[IO](4).unsafeRunSync() + + private val blockingEC = ExecutionContext.fromExecutorService( + Executors.newFixedThreadPool(4) // limit concurrency + ) def getBinaryChunk(offset: Int, length: Int, fileName: String): IO[Array[Byte]] = { val fileResource = Resource.make( IO.blocking(new RandomAccessFile(s"$GRIB_FOLDER/$fileName", "r")) )(file => IO.blocking(file.close())) - fileResource.use { file => - IO.blocking { - val buffer = new Array[Byte](length) - file.seek(offset) - file.readFully(buffer) - buffer - } + val logMemory = IO { + val runtime = Runtime.getRuntime + val usedMemoryMB = (runtime.totalMemory - runtime.freeMemory) / 1024 / 1024 + val maxMemoryMB = runtime.maxMemory / 1024 / 1024 + println(s"Memory usage before reading $fileName: $usedMemoryMB MB / $maxMemoryMB MB max") } + + for { + _ <- logMemory + result <- fileResource.use { file => + IO.blocking { + val buffer = new Array[Byte](length) + file.seek(offset) + file.readFully(buffer) + buffer + }.evalOn(blockingEC) + } + _ <- IO { System.gc() } + } yield result } def getForecasts(): IO[List[(ZonedDateTime, ZonedDateTime)]] = { diff --git a/web/src/harmonie/ReferenceTimes.tsx b/web/src/harmonie/ReferenceTimes.tsx index f06f4a9..b897b5e 100644 --- a/web/src/harmonie/ReferenceTimes.tsx +++ b/web/src/harmonie/ReferenceTimes.tsx @@ -6,7 +6,7 @@ import { drawGrib } from './draw/drawGrib' import { CROP_BOUNDS } from './DrawView' import styles from './harmonie.module.css' -import { handleProgressivePromises } from '../helpers/progressivePromises' +import { processPromisesInBatches } from '../helpers/progressivePromises' const METEO_PARAMS: [string, MeteoParam][] = [ ['temperature', { discipline: 0, category: 0, product: 0, levelType: -1, levelValue: -1, subType: 'now' }], @@ -62,7 +62,7 @@ export const ReferenceTimes: Component<{ .map((grib): [string, undefined] => [grib.time.forecastTime, undefined]) .sort((a, b) => a[0] > b[0] ? 1 : -1) setImgList(emptyImgList) - const promiseList = forecastList.map(async (grib): Promise<[string, ImageBitmap]> => { + const promiseFnsList = forecastList.map((grib): () => Promise<[string, ImageBitmap]> => async () => { const canvas = document.createElement('canvas') const [messages, buffers, bitmasks] = await fetchGribBinaries(grib, getGribList()) drawGrib(canvas, messages, buffers, bitmasks, cropBounds, contour, isInterpolated) @@ -72,8 +72,8 @@ export const ReferenceTimes: Component<{ return [grib.time.forecastTime, img] }) - handleProgressivePromises( - promiseList, + processPromisesInBatches( + promiseFnsList, ([forecastDate, img]) => { const udpdatedImgList = [...getImgList()] const idx = udpdatedImgList.findIndex(([d]) => forecastDate === d) diff --git a/web/src/helpers/progressivePromises.ts b/web/src/helpers/progressivePromises.ts index 3aee7cc..1572928 100644 --- a/web/src/helpers/progressivePromises.ts +++ b/web/src/helpers/progressivePromises.ts @@ -1,10 +1,10 @@ -export async function handleProgressivePromises( - promises: Promise[], +async function processPromises( + promiseFns: (() => Promise)[], onProgress: (result: T) => void, ): Promise<(T | undefined)[]> { - const allPromises = promises.map(async promise => { + const allPromises = promiseFns.map(async promise => { try { - const result = await promise; + const result = await promise(); onProgress(result) return result } catch (error) { @@ -17,3 +17,20 @@ export async function handleProgressivePromises( result.status === 'fulfilled' ? result.value : undefined ) } + +export async function processPromisesInBatches( + promiseFns: (() => Promise)[], + onProgress: (result: T) => void, + batchSize = 3, +): Promise<(T | undefined)[]> { + const results = [] + while(await promiseFns.length > 0) { + const size = Math.min(batchSize, promiseFns.length) + const batchPromises = promiseFns.splice(0, size) + console.log('process:', batchPromises.length) + const batchResults = await processPromises(batchPromises, onProgress) + results.push(...batchResults) + } + + return results +}