diff --git a/src/main/java/com/dedicatedcode/paikka/service/importer/BoundaryImportStatistics.java b/src/main/java/com/dedicatedcode/paikka/service/importer/BoundaryImportStatistics.java index 537c1e1..52c872f 100644 --- a/src/main/java/com/dedicatedcode/paikka/service/importer/BoundaryImportStatistics.java +++ b/src/main/java/com/dedicatedcode/paikka/service/importer/BoundaryImportStatistics.java @@ -73,6 +73,10 @@ public String toString() { private final AtomicLong relationsProcessed = new AtomicLong(0); private final AtomicLong h3CellsGenerated = new AtomicLong(0); + private volatile int mergedThreadDatabasesTotal; + private final AtomicLong mergedThreadDatabasesProcessed = new AtomicLong(0); + private final AtomicLong mergedEntriesProcessed = new AtomicLong(0); + private volatile String currentPhase = "Initializing"; private volatile boolean running = true; private final long startTime = System.currentTimeMillis(); @@ -129,6 +133,30 @@ public void addH3CellsGenerated(long count) { h3CellsGenerated.addAndGet(count); } + public void setMergedThreadCount(int count) { + this.mergedThreadDatabasesTotal = count; + } + + public int getMergedThreadDatabasesTotal() { + return mergedThreadDatabasesTotal; + } + + public long getMergedThreadDatabasesProcessed() { + return mergedThreadDatabasesProcessed.get(); + } + + public void incrementMergedThreadDatabasesProcessed() { + mergedThreadDatabasesProcessed.incrementAndGet(); + } + + public long getMergedEntriesProcessed() { + return mergedEntriesProcessed.get(); + } + + public void incrementMergedEntriesProcessed() { + mergedEntriesProcessed.incrementAndGet(); + } + public String getCurrentPhase() { return currentPhase; } @@ -247,6 +275,15 @@ public void startProgressReporter() { if (getErrorsTotal() > 0) { sb.append(String.format(" │ \033[31mErrors:\033[0m %d", getErrorsTotal())); } + } else if (phase.contains("Merging")) { + long dbsProcessed = mergedThreadDatabasesProcessed.get(); + int dbsTotal = mergedThreadDatabasesTotal; + long entriesProcessed = mergedEntriesProcessed.get(); + long entriesPerSec = phaseSeconds > 0 ? (long) (entriesProcessed / phaseSeconds) : 0; + sb.append(String.format("\033[1;36m[%s]\033[0m \033[1mMerging H3 thread DBs\033[0m", formatTime(elapsed))); + sb.append(String.format(" │ \033[32mDBs:\033[0m %d/%d", dbsProcessed, dbsTotal)); + sb.append(String.format(" │ \033[36mEntries:\033[0m %s \033[33m(%s/s)\033[0m", + formatCompactNumber(entriesProcessed), formatCompactRate(entriesPerSec))); } else { sb.append(String.format("\033[1;36m[%s]\033[0m %s", formatTime(elapsed), phase)); } @@ -277,6 +314,8 @@ public void printFinalStatistics() { double totalSeconds = totalTime / 1000.0; double phase1Seconds = Math.max(0.001, phase1Duration / 1000.0); double phase2Seconds = Math.max(0.001, phase2Duration / 1000.0); + long phase3Duration = totalTime - phase1Duration - phase2Duration; + double phase3Seconds = Math.max(0.001, phase3Duration / 1000.0); System.out.printf("\n\033[1;37mTotal Import Time:\033[0m \033[1;33m%s\033[0m%n%n", formatTime(getTotalTime())); @@ -299,6 +338,9 @@ public void printFinalStatistics() { System.out.printf("│ \033[33mH3 Cells Generated\033[0m │ %15s │ %13s/s │%n", formatCompactNumber(getH3CellsGenerated()), formatCompactNumber((long) (getH3CellsGenerated() / phase2Seconds))); + System.out.printf("│ \033[36mH3 Entries Merged\033[0m │ %15s │ %13s/s │%n", + formatCompactNumber(getMergedEntriesProcessed()), + formatCompactNumber((long) (getMergedEntriesProcessed() / phase3Seconds))); System.out.println("└──────────────────────┴─────────────────┴─────────────────┘"); if (h3OsmSizeBytes > 0) { diff --git a/src/main/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporter.java b/src/main/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporter.java index e684ed9..60debee 100644 --- a/src/main/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporter.java +++ b/src/main/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporter.java @@ -45,7 +45,7 @@ * Reads a pre-filtered boundaries_only.pbf (Nodes -> Ways -> Relations ordered) * and produces three RocksDB databases for offline mobile lookup: * - h3_to_osm: H3_CELL_ID (uint64) -> List[OSM_ID] (raw byte array) - * - region_metadata: OSM_ID -> [total cell count (long), h3 resolution (int)] (12 bytes) + * - region_metadata: OSM_ID -> [total cell count (long), h3 resolution (int), admin_level (int)] (16 bytes) * - region_geometry: OSM_ID -> simplified WKB (bytes) */ @Service @@ -275,7 +275,7 @@ public void importBoundaries(List pbfPaths, String outputDir) throws Exc stats.addH3CellsGenerated(totalCells); startTime = System.currentTimeMillis(); - tmpRegionMeta.put(wo, longToBytes(stub.osmId()), cellMetaToBytes(totalCells, resolution)); + tmpRegionMeta.put(wo, longToBytes(stub.osmId()), cellMetaToBytes(totalCells, resolution, stub.adminLevel())); if (stub.adminLevel() <= 3) { logger.debug("H3 Cells written to tmpRegionMeta in {}ms for OSM ID: {}", System.currentTimeMillis() - startTime, stub.osmId()); } @@ -304,11 +304,13 @@ public void importBoundaries(List pbfPaths, String outputDir) throws Exc // Merge per-thread H3 DBs into the final h3_to_osm stats.setCurrentPhase(3, "3.1: Merging H3 thread DBs"); + stats.setMergedThreadCount(threads); for (int t = 0; t < threads; t++) { if (Files.exists(threadH3Paths[t])) { try (RocksDB threadDb = RocksDB.open(cacheOpts, threadH3Paths[t].toString())) { copyH3Db(threadDb, h3ToOsm); } + stats.incrementMergedThreadDatabasesProcessed(); cleanup(threadH3Paths[t]); } } @@ -367,6 +369,7 @@ private void copyH3Db(RocksDB source, RocksDB target) throws RocksDBException { target.put(wo, key, merged); } it.next(); + this.stats.incrementMergedEntriesProcessed(); } } } @@ -579,10 +582,11 @@ private long[] bytesToLongArray(byte[] b) { return arr; } - private byte[] cellMetaToBytes(long cellCount, int resolution) { - return ByteBuffer.allocate(12).order(ByteOrder.BIG_ENDIAN) + private byte[] cellMetaToBytes(long cellCount, int resolution, int adminLevel) { + return ByteBuffer.allocate(16).order(ByteOrder.BIG_ENDIAN) .putLong(cellCount) .putInt(resolution) + .putInt(adminLevel) .array(); } diff --git a/src/test/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporterTest.java b/src/test/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporterTest.java index f3d1ae2..e176be1 100644 --- a/src/test/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporterTest.java +++ b/src/test/java/com/dedicatedcode/paikka/service/importer/StandaloneBoundaryImporterTest.java @@ -135,15 +135,18 @@ void testRegionMetadataFormat() throws RocksDBException { assertEquals(8, key.length, "Region metadata key should be 8 bytes (OSM ID)"); byte[] val = it.value(); - assertEquals(12, val.length, "Value should be 12 bytes (8-byte cell count + 4-byte resolution)"); + assertEquals(16, val.length, "Value should be 16 bytes (8-byte cell count + 4-byte resolution + 4-byte admin level)"); ByteBuffer bb = ByteBuffer.wrap(val).order(ByteOrder.BIG_ENDIAN); long cellCount = bb.getLong(); int resolution = bb.getInt(); + int adminLevel = bb.getInt(); assertTrue(cellCount > 0, "Cell count should be positive, got: " + cellCount); assertTrue(resolution >= 4 && resolution <= 9, "Resolution should be between 4 and 9, got: " + resolution); + assertTrue(adminLevel >= 1 && adminLevel <= 11, + "Admin level should be between 2 and 11, got: " + adminLevel); int count = 0; it.seekToFirst(); @@ -167,17 +170,20 @@ void testAllRegionMetadataEntriesHaveValidResolution() throws RocksDBException { int checked = 0; while (it.isValid()) { byte[] val = it.value(); - assertEquals(12, val.length, - "Every region_metadata entry should be 12 bytes"); + assertEquals(16, val.length, + "Every region_metadata entry should be 16 bytes"); ByteBuffer bb = ByteBuffer.wrap(val).order(ByteOrder.BIG_ENDIAN); long cellCount = bb.getLong(); int resolution = bb.getInt(); + int adminLevel = bb.getInt(); assertTrue(cellCount > 0, "Cell count should be positive for entry " + checked); assertTrue(resolution >= 4 && resolution <= 9, "Resolution should be 4-9 for entry " + checked + ", got: " + resolution); + assertTrue(adminLevel >= 1 && adminLevel <= 11, + "Admin level should be 2-11 for entry " + checked + ", got: " + adminLevel); checked++; it.next();