From 1ac2f6077a298647246bc77169e2ae1c82507730 Mon Sep 17 00:00:00 2001 From: asmhatre Date: Wed, 30 Sep 2026 19:14:10 +0000 Subject: [PATCH] nio2: NIO.2 file API companion code (streaming 5 GB at -Xmx64m, handle leaks, WatchService, mmap) Co-Authored-By: Claude Sonnet 5.5 Claude-Session: https://claude.ai/code/session_01KqJyCidz3ZgRyHABv2GVJh --- README.md | 1 + nio2/README.md | 44 +++++ nio2/output/01-environment.txt | 18 ++ nio2/output/02-generate.txt | 8 + nio2/output/03-stream-lines.txt | 17 ++ nio2/output/04-stream-reader.txt | 17 ++ nio2/output/05-read-all-fail.txt | 22 +++ nio2/output/06-handle-leak.txt | 11 ++ nio2/output/07-handle-exhaustion.txt | 3 + nio2/output/08-walk-find.txt | 9 + nio2/output/09-lines-charset.txt | 4 + nio2/output/10-watch.txt | 19 ++ nio2/output/11-watch-tuned-and-repeat.txt | 13 ++ nio2/output/12-watchkey-source-excerpt.txt | 54 ++++++ nio2/output/13-mapped.txt | 7 + nio2/output/14-tests.txt | 4 + nio2/pom.xml | 43 ++++ nio2/scripts/run-all.sh | 70 +++++++ .../com/ankurm/nio2/HandleExhaustion.java | 27 +++ .../java/com/ankurm/nio2/HandleLeakDemo.java | 68 +++++++ .../main/java/com/ankurm/nio2/HeapProbe.java | 48 +++++ .../com/ankurm/nio2/LinesCharsetDemo.java | 36 ++++ .../main/java/com/ankurm/nio2/LogFile.java | 46 +++++ .../main/java/com/ankurm/nio2/MappedDemo.java | 67 +++++++ .../java/com/ankurm/nio2/ReadAllFail.java | 28 +++ .../java/com/ankurm/nio2/StreamLines.java | 58 ++++++ .../java/com/ankurm/nio2/WalkFindDemo.java | 64 ++++++ .../main/java/com/ankurm/nio2/WatchDemo.java | 82 ++++++++ .../com/ankurm/nio2/Nio2BehaviourTest.java | 183 ++++++++++++++++++ pom.xml | 1 + 30 files changed, 1072 insertions(+) create mode 100644 nio2/README.md create mode 100644 nio2/output/01-environment.txt create mode 100644 nio2/output/02-generate.txt create mode 100644 nio2/output/03-stream-lines.txt create mode 100644 nio2/output/04-stream-reader.txt create mode 100644 nio2/output/05-read-all-fail.txt create mode 100644 nio2/output/06-handle-leak.txt create mode 100644 nio2/output/07-handle-exhaustion.txt create mode 100644 nio2/output/08-walk-find.txt create mode 100644 nio2/output/09-lines-charset.txt create mode 100644 nio2/output/10-watch.txt create mode 100644 nio2/output/11-watch-tuned-and-repeat.txt create mode 100644 nio2/output/12-watchkey-source-excerpt.txt create mode 100644 nio2/output/13-mapped.txt create mode 100644 nio2/output/14-tests.txt create mode 100644 nio2/pom.xml create mode 100755 nio2/scripts/run-all.sh create mode 100644 nio2/src/main/java/com/ankurm/nio2/HandleExhaustion.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/HandleLeakDemo.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/HeapProbe.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/LinesCharsetDemo.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/LogFile.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/MappedDemo.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/ReadAllFail.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/StreamLines.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/WalkFindDemo.java create mode 100644 nio2/src/main/java/com/ankurm/nio2/WatchDemo.java create mode 100644 nio2/src/test/java/com/ankurm/nio2/Nio2BehaviourTest.java diff --git a/README.md b/README.md index 2edc4d6..75d82ba 100644 --- a/README.md +++ b/README.md @@ -19,6 +19,7 @@ article; each module's own README has that article's version table, quickstart, | [`serialization`](serialization/) | Java Serialization in 2026: Why It's Dangerous and What Replaced It | | [`strings`](strings/) | Java Strings Deep Dive: Interning, StringBuilder, Text Blocks, and the Concatenation Benchmark | | [`collections`](collections/) | How to Sort a HashMap by Value (and Key) in Java | +| [`nio2`](nio2/) | Java NIO.2 File API: Files, Path, WatchService, and Streaming Large Files Without OOM | ## License diff --git a/nio2/README.md b/nio2/README.md new file mode 100644 index 0000000..a81a750 --- /dev/null +++ b/nio2/README.md @@ -0,0 +1,44 @@ +# nio2 + +Companion code for the ankurm.com post *"Java NIO.2 File API: Files, Path, WatchService, and Streaming +Large Files Without OOM."* Module `nio2` in `java-core-examples`. + +All explanation lives in the post; this module holds the runnable evidence and the captured output. + +## Versions + +| Component | Version | +|---|---| +| JDK | 25.0.4.1+1 (Temurin, LTS) | +| JUnit Jupiter | 5.11.0 | +| OS | Linux 6.18 (WatchService and `/proc/self/fd` results are Linux-specific) | +| Hardware | 2 vCPU x86-64 VM shared with other jobs (timings are indicative, not a leaderboard) | +| Big file | 5,000,000,069 bytes (the full 5 GB fit; the file is generated, then deleted by the script) | + +## Quickstart + +```bash +export JDK25_HOME=/path/to/jdk-25 +./scripts/run-all.sh # needs ~5.3 GB free in /tmp/nio2-data; regenerates output/ +BIG_BYTES=500000000 ./scripts/run-all.sh # smaller "big" file if disk is tight +mvn -q -pl nio2 -am test # 12 assertions, small files only +``` + +## What is in here + +| File | Shows | Output | +|---|---|---| +| `LogFile` | deterministic generator for the synthetic log file | `02` | +| `StreamLines` | `Files.lines` and `Files.newBufferedReader` over 5 GB at `-Xmx64m`, heap and GC evidence | `03`, `04` | +| `ReadAllFail` | `readAllLines` / `readString` / `readAllBytes` at the same `-Xmx`, and how much heap they need | `05` | +| `HandleLeakDemo` | open file descriptors around unclosed `Files.lines` and `Files.walk` | `06` | +| `HandleExhaustion` | what an unclosed `Files.walk` becomes under `ulimit -n 64` | `07` | +| `WalkFindDemo` | `list` / `walk` / `find` / `walkFileTree` / `DirectoryStream`, symlink loop | `08` | +| `LinesCharsetDemo` | a bad byte fails lazily inside `Files.lines` | `09` | +| `WatchDemo` | WatchService events, coalescing, no recursion, OVERFLOW at 513 pending events | `10`, `11` | +| (JDK source excerpt) | `AbstractWatchKey.signalEvent` from the JDK's own `src.zip` | `12` | +| `MappedDemo` | 2 GiB classic `map` limit, Arena mapping, resident memory | `13` | +| `Nio2BehaviourTest` | 12 assertions behind the claims above | `14` | + +`01` records the environment (kernel, JDK, `ulimit -n`, inotify limit, free disk). +Timings move with whatever else the machine is doing; counts, exceptions and GC figures do not. diff --git a/nio2/output/01-environment.txt b/nio2/output/01-environment.txt new file mode 100644 index 0000000..a2ce344 --- /dev/null +++ b/nio2/output/01-environment.txt @@ -0,0 +1,18 @@ +$ uname -sr +Linux 6.18.44-fc-v50 +$ nproc +2 +$ java -version +openjdk version "25.0.4.1" 2026-08-18 LTS +OpenJDK Runtime Environment Temurin-25.0.4.1+1 (build 25.0.4.1+1-LTS) +OpenJDK 64-Bit Server VM Temurin-25.0.4.1+1 (build 25.0.4.1+1-LTS, mixed mode, sharing) +$ ulimit -n +20000 +$ cat /proc/sys/fs/inotify/max_queued_events +16384 +$ df -h (data directory, before) +Filesystem Size Used Avail Use% Mounted on +/dev/vda 252G 13G 29G 31% / +$ df -h (data directory, after deleting is done by the exit trap; shown here before) +Filesystem Size Used Avail Use% Mounted on +/dev/vda 252G 18G 25G 43% / diff --git a/nio2/output/02-generate.txt b/nio2/output/02-generate.txt new file mode 100644 index 0000000..329caf8 --- /dev/null +++ b/nio2/output/02-generate.txt @@ -0,0 +1,8 @@ +generated big.log: 5,000,000,069 bytes, 34,120,522 lines, 681,844 ERROR lines in 13.7 s +generated medium.log: 200,000,016 bytes, 1,378,768 lines, 27,530 ERROR lines in 1.4 s +$ ls -l (generated files) +5000000069 big.log +200000016 medium.log +$ df -h (data directory, with both files on disk) +Filesystem Size Used Avail Use% Mounted on +/dev/vda 252G 18G 25G 43% / diff --git a/nio2/output/03-stream-lines.txt b/nio2/output/03-stream-lines.txt new file mode 100644 index 0000000..e363ddd --- /dev/null +++ b/nio2/output/03-stream-lines.txt @@ -0,0 +1,17 @@ +$ java -Xmx64m -Xlog:gc StreamLines lines big.log +Files.lines over big.log (5,000,000,069 bytes) +lines=34,120,522 errors=681,844 services=20 +elapsed 9.4 s +heap: -Xmx=64 MiB, peak used (sampled every 20 ms)=38 MiB, GC collections=222, GC time=200 ms +GC log: 223 lines, 222 young pauses, 0 full GCs +first lines of the GC log: +[0.004s][info][gc] Using G1 +[0.321s][info][gc] GC(0) Pause Young (Normal) (G1 Evacuation Pause) 15M->1M(64M) 7.040ms +[0.440s][info][gc] GC(1) Pause Young (Normal) (G1 Evacuation Pause) 28M->1M(64M) 1.218ms +[0.497s][info][gc] GC(2) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 1.372ms +last lines of the GC log: +[9.358s][info][gc] GC(219) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 0.696ms +[9.396s][info][gc] GC(220) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 0.466ms +[9.430s][info][gc] GC(221) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 0.497ms +largest 'Pause Young' (ms): 7.040 +highest heap-after-GC seen in the log: 1M of 64M diff --git a/nio2/output/04-stream-reader.txt b/nio2/output/04-stream-reader.txt new file mode 100644 index 0000000..5e79b6d --- /dev/null +++ b/nio2/output/04-stream-reader.txt @@ -0,0 +1,17 @@ +$ java -Xmx64m -Xlog:gc StreamLines reader big.log +Files.newBufferedReader over big.log (5,000,000,069 bytes) +lines=34,120,522 errors=681,844 services=20 +elapsed 9.3 s +heap: -Xmx=64 MiB, peak used (sampled every 20 ms)=38 MiB, GC collections=222, GC time=212 ms +GC log: 223 lines, 222 young pauses, 0 full GCs +first lines of the GC log: +[0.005s][info][gc] Using G1 +[0.273s][info][gc] GC(0) Pause Young (Normal) (G1 Evacuation Pause) 15M->1M(64M) 2.473ms +[0.349s][info][gc] GC(1) Pause Young (Normal) (G1 Evacuation Pause) 28M->1M(64M) 1.211ms +[0.401s][info][gc] GC(2) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 1.173ms +last lines of the GC log: +[9.286s][info][gc] GC(219) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 0.515ms +[9.319s][info][gc] GC(220) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 0.427ms +[9.353s][info][gc] GC(221) Pause Young (Normal) (G1 Evacuation Pause) 38M->1M(64M) 0.469ms +largest 'Pause Young' (ms): 6.615 +highest heap-after-GC seen in the log: 1M of 64M diff --git a/nio2/output/05-read-all-fail.txt b/nio2/output/05-read-all-fail.txt new file mode 100644 index 0000000..82692bb --- /dev/null +++ b/nio2/output/05-read-all-fail.txt @@ -0,0 +1,22 @@ +$ java -Xmx64m ReadAllFail readAllLines medium.log +OutOfMemoryError: Java heap space +$ java -Xmx64m ReadAllFail readString medium.log +OutOfMemoryError: Java heap space +$ java -Xmx64m ReadAllFail readAllBytes medium.log +OutOfMemoryError: Java heap space +$ java -Xmx64m ReadAllFail readAllLines big.log +OutOfMemoryError: Java heap space +$ java -Xmx64m ReadAllFail readString big.log +OutOfMemoryError: Required array size too large +$ java -Xmx64m ReadAllFail readAllBytes big.log +OutOfMemoryError: Required array size too large + +How much heap does the 200 MB file need? (readAllLines, then readString) +readAllLines -Xmx256m OutOfMemoryError: Java heap space +readString -Xmx256m survived: 200000016 chars +readAllLines -Xmx512m survived: 1378768 lines +readString -Xmx512m survived: 200000016 chars +readAllLines -Xmx1g survived: 1378768 lines +readString -Xmx1g survived: 200000016 chars +readAllLines -Xmx2g survived: 1378768 lines +readString -Xmx2g survived: 200000016 chars diff --git a/nio2/output/06-handle-leak.txt b/nio2/output/06-handle-leak.txt new file mode 100644 index 0000000..11599c0 --- /dev/null +++ b/nio2/output/06-handle-leak.txt @@ -0,0 +1,11 @@ +open fds at start: 8 +after 50 fully consumed, unclosed Files.walk: +0 fds +after 50 unclosed Files.lines + findFirst(): +50 fds +after closing them: +0 fds +after 50 try-with-resources Files.lines: +0 fds +after 50 unclosed Files.walk + findFirst(): +300 fds +after closing them: +0 fds +after 50 unreferenced, unclosed Files.lines: +50 fds +after System.gc(): +0 fds +after 50 unreferenced, unclosed Files.walk: +300 fds +after System.gc(): +300 fds diff --git a/nio2/output/07-handle-exhaustion.txt b/nio2/output/07-handle-exhaustion.txt new file mode 100644 index 0000000..02b0b3c --- /dev/null +++ b/nio2/output/07-handle-exhaustion.txt @@ -0,0 +1,3 @@ +$ ulimit -n 64; java HandleExhaustion +failed on call 10 after 9 leaked walks +java.io.UncheckedIOException: java.nio.file.FileSystemException: /tree/d2/sub: Too many open files diff --git a/nio2/output/08-walk-find.txt b/nio2/output/08-walk-find.txt new file mode 100644 index 0000000..dbf1450 --- /dev/null +++ b/nio2/output/08-walk-find.txt @@ -0,0 +1,9 @@ +Files.list (one level): [.git, README.md, src] +Files.walk (everything): [, .git, .git/objects, .git/objects/blob.bin, README.md, src, src/main, src/main/App.java, src/test, src/test/AppTest.java] +Files.walk maxDepth=1: [, .git, README.md, src] +Files.find *.java: [src/main/App.java, src/test/AppTest.java] +Files.find size > 1000: [src/test/AppTest.java] +walkFileTree skipping .git: [README.md, src/main/App.java, src/test/AppTest.java] +DirectoryStream glob: [App.java] +walk with a symlink loop, default: 11 entries, no error +walk with FOLLOW_LINKS: FileSystemLoopException diff --git a/nio2/output/09-lines-charset.txt b/nio2/output/09-lines-charset.txt new file mode 100644 index 0000000..91e443c --- /dev/null +++ b/nio2/output/09-lines-charset.txt @@ -0,0 +1,4 @@ +lines processed before failure: 0 +UncheckedIOException caused by java.nio.charset.MalformedInputException: Input length = 1 +with ISO_8859_1: 3 lines, no error +readAllLines throws checked MalformedInputException diff --git a/nio2/output/10-watch.txt b/nio2/output/10-watch.txt new file mode 100644 index 0000000..624d8f1 --- /dev/null +++ b/nio2/output/10-watch.txt @@ -0,0 +1,19 @@ +1. create a file, write it twice, delete it (events polled afterwards): + ENTRY_CREATE count=1 context=a.txt + ENTRY_MODIFY count=2 context=a.txt + ENTRY_DELETE count=1 context=a.txt +2. Files.writeString of 100 KiB to a new file (one call): + ENTRY_CREATE count=1 context=b.txt + ENTRY_MODIFY count=1 context=b.txt +3. 50 appends to one file with no polling in between, then poll: + {(batches)=1, ENTRY_MODIFY=9} +4. Files.move (atomic rename) of g.txt to h.txt: + ENTRY_DELETE count=1 context=g.txt + ENTRY_CREATE count=1 context=h.txt +5. a file created inside a NEW subdirectory (subdirectory not registered): + ENTRY_CREATE count=1 context=sub +6. burst: create N empty files in a fresh directory before polling once (jdk.nio.file.WatchService.maxEventsPerPoll=default): + N=100 delivered 100 events: ENTRY_CREATE=100 OVERFLOW(count)=0, files on disk=100 + N=512 delivered 512 events: ENTRY_CREATE=512 OVERFLOW(count)=0, files on disk=512 + N=513 delivered 1 events: ENTRY_CREATE=0 OVERFLOW(count)=1, files on disk=513 + N=5000 delivered 1 events: ENTRY_CREATE=0 OVERFLOW(count)=4488, files on disk=5000 diff --git a/nio2/output/11-watch-tuned-and-repeat.txt b/nio2/output/11-watch-tuned-and-repeat.txt new file mode 100644 index 0000000..5f5e800 --- /dev/null +++ b/nio2/output/11-watch-tuned-and-repeat.txt @@ -0,0 +1,13 @@ +$ java -Djdk.nio.file.WatchService.maxEventsPerPoll=10000 WatchDemo (scenario 6 only) +6. burst: create N empty files in a fresh directory before polling once (jdk.nio.file.WatchService.maxEventsPerPoll=10000): + N=100 delivered 100 events: ENTRY_CREATE=100 OVERFLOW(count)=0, files on disk=100 + N=512 delivered 512 events: ENTRY_CREATE=512 OVERFLOW(count)=0, files on disk=512 + N=513 delivered 513 events: ENTRY_CREATE=513 OVERFLOW(count)=0, files on disk=513 + N=5000 delivered 5000 events: ENTRY_CREATE=5000 OVERFLOW(count)=0, files on disk=5000 + +Scenario 3 (50 appends, no polling) over five more runs: + {(batches)=1, ENTRY_MODIFY=21} + {(batches)=1, ENTRY_MODIFY=20} + {(batches)=1, ENTRY_MODIFY=2} + {(batches)=1, ENTRY_MODIFY=1} + {(batches)=1, ENTRY_MODIFY=19} diff --git a/nio2/output/12-watchkey-source-excerpt.txt b/nio2/output/12-watchkey-source-excerpt.txt new file mode 100644 index 0000000..571d5f5 --- /dev/null +++ b/nio2/output/12-watchkey-source-excerpt.txt @@ -0,0 +1,54 @@ + final void signalEvent(WatchEvent.Kind kind, Object context) { + boolean isModify = (kind == StandardWatchEventKinds.ENTRY_MODIFY); + synchronized (this) { + int size = events.size(); + if (size > 0) { + // if the previous event is an OVERFLOW event or this is a + // repeated event then we simply increment the counter + WatchEvent prev = events.get(size-1); + if ((prev.kind() == StandardWatchEventKinds.OVERFLOW) || + ((kind == prev.kind() && + Objects.equals(context, prev.context())))) + { + ((Event)prev).increment(); + return; + } + + // if this is a modify event and the last entry for the context + // is a modify event then we simply increment the count + if (!lastModifyEvents.isEmpty()) { + if (isModify) { + WatchEvent ev = lastModifyEvents.get(context); + if (ev != null) { + assert ev.kind() == StandardWatchEventKinds.ENTRY_MODIFY; + ((Event)ev).increment(); + return; + } + } else { + // not a modify event so remove from the map as the + // last event will no longer be a modify event. + lastModifyEvents.remove(context); + } + } + + // if the list has reached the limit then drop pending events + // and queue an OVERFLOW event + if (size >= MAX_EVENT_LIST_SIZE) { + kind = StandardWatchEventKinds.OVERFLOW; + isModify = false; + context = null; + } + } + + // non-repeated event + Event ev = + new Event<>((WatchEvent.Kind)kind, context); + if (isModify) { + lastModifyEvents.put(context, ev); + } else if (kind == StandardWatchEventKinds.OVERFLOW) { + // drop all pending events + events.clear(); + lastModifyEvents.clear(); + } + events.add(ev); + signal(); diff --git a/nio2/output/13-mapped.txt b/nio2/output/13-mapped.txt new file mode 100644 index 0000000..db5d838 --- /dev/null +++ b/nio2/output/13-mapped.txt @@ -0,0 +1,7 @@ +$ java -Xmx64m MappedDemo big.log +file big.log: 5,000,000,069 bytes, -Xmx=64 MiB, RSS at start=40 MiB +classic map of the whole file: IllegalArgumentException: Size exceeds Integer.MAX_VALUE +Arena + MemorySegment: 34,120,522 newlines in 5.1 s; RSS after the Arena was closed=51 MiB +classic 1 GiB windows: 34,120,522 newlines in 7.8 s; RSS right after=4821 MiB +heap: -Xmx=64 MiB, peak used (sampled every 20 ms)=4 MiB, GC collections=0, GC time=0 ms +RSS one second after System.gc()=52 MiB diff --git a/nio2/output/14-tests.txt b/nio2/output/14-tests.txt new file mode 100644 index 0000000..15a2ffa --- /dev/null +++ b/nio2/output/14-tests.txt @@ -0,0 +1,4 @@ +------------------------------------------------------------------------------- +Test set: com.ankurm.nio2.Nio2BehaviourTest +------------------------------------------------------------------------------- +Tests run: 12, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 3.755 s -- in com.ankurm.nio2.Nio2BehaviourTest diff --git a/nio2/pom.xml b/nio2/pom.xml new file mode 100644 index 0000000..3fc9039 --- /dev/null +++ b/nio2/pom.xml @@ -0,0 +1,43 @@ + + + 4.0.0 + + + com.ankurm + java-core-examples + 1.0 + + + nio2 + nio2 + Java NIO.2 file API: streaming a multi-gigabyte file at constant heap, readAllLines/readString failure modes, unclosed Files.lines/Files.walk handles, WatchService behaviour on Linux, and memory-mapped files. + + + + org.junit.jupiter + junit-jupiter + 5.11.0 + test + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.13.0 + + 25 + + + + org.apache.maven.plugins + maven-surefire-plugin + 3.2.5 + + + + diff --git a/nio2/scripts/run-all.sh b/nio2/scripts/run-all.sh new file mode 100755 index 0000000..b470d71 --- /dev/null +++ b/nio2/scripts/run-all.sh @@ -0,0 +1,70 @@ +#!/usr/bin/env bash +# Regenerates every file in ../output/. Requires JDK25_HOME. Needs ~5.3 GB of free disk in $DATA (default /tmp/nio2-data); +# the generated files are deleted on exit. Set BIG_BYTES to use a smaller "big" file (default 5 GB). +set -euo pipefail +[[ -z "${JDK25_HOME:-}" ]] && { echo "JDK25_HOME must be set" >&2; exit 1; } +cd "$(dirname "$0")/.." +OUT=output; mkdir -p "$OUT" +DATA="${DATA:-/tmp/nio2-data}"; mkdir -p "$DATA" +BIG_BYTES="${BIG_BYTES:-5000000000}" +MED_BYTES=200000000 +trap 'rm -rf "$DATA"' EXIT +export JAVA_HOME="$JDK25_HOME" +mvn -q -f ../pom.xml -pl nio2 -am package +J="$JDK25_HOME/bin/java"; CP=target/classes +strip() { grep -v "Picked up"; } + +{ echo '$ uname -sr'; uname -sr; echo '$ nproc'; nproc; echo '$ java -version'; "$J" -version 2>&1 | strip + echo '$ ulimit -n'; ulimit -n + echo '$ cat /proc/sys/fs/inotify/max_queued_events'; cat /proc/sys/fs/inotify/max_queued_events + echo '$ df -h (data directory, before)'; df -h "$DATA" | sed 's/ */ /g'; } > "$OUT/01-environment.txt" + +echo "==> generate files" +{ "$J" -cp $CP com.ankurm.nio2.LogFile "$DATA/big.log" "$BIG_BYTES" 2>&1 | strip + "$J" -cp $CP com.ankurm.nio2.LogFile "$DATA/medium.log" "$MED_BYTES" 2>&1 | strip + echo '$ ls -l (generated files)'; ls -l "$DATA"/*.log | awk '{print $5, $9}' | sed "s#$DATA/##" + echo '$ df -h (data directory, with both files on disk)'; df -h "$DATA" | sed 's/ */ /g'; } > "$OUT/02-generate.txt" + +echo "==> stream the big file at -Xmx64m" +for m in lines reader; do + n=$([[ $m == lines ]] && echo 03 || echo 04) + { echo "\$ java -Xmx64m -Xlog:gc StreamLines $m big.log" + "$J" -Xmx64m -Xlog:gc:file="$DATA/gc-$m.log" -cp $CP com.ankurm.nio2.StreamLines $m "$DATA/big.log" 2>&1 | strip + echo "GC log: $(wc -l < "$DATA/gc-$m.log") lines, $(grep -c 'Pause Young' "$DATA/gc-$m.log") young pauses, $(grep -c 'Pause Full' "$DATA/gc-$m.log") full GCs" + echo "first lines of the GC log:"; head -4 "$DATA/gc-$m.log" | cut -c1-150 + echo "last lines of the GC log:"; tail -3 "$DATA/gc-$m.log" | cut -c1-150 + echo "largest 'Pause Young' (ms): $(grep 'Pause Young' "$DATA/gc-$m.log" | sed 's/.* \([0-9.]*\)ms$/\1/' | sort -n | tail -1)" + echo "highest heap-after-GC seen in the log: $(grep 'Pause Young' "$DATA/gc-$m.log" | sed 's/.*->\([0-9]*\)M(.*/\1/' | sort -n | tail -1)M of $(grep -m1 'Pause Young' "$DATA/gc-$m.log" | sed 's/.*M(\([0-9]*\)M).*/\1/')M"; } > "$OUT/$n-stream-$m.txt" +done + +echo "==> readAllLines / readString / readAllBytes" +{ for f in medium big; do + for m in readAllLines readString readAllBytes; do + echo "\$ java -Xmx64m ReadAllFail $m $f.log"; "$J" -Xmx64m -cp $CP com.ankurm.nio2.ReadAllFail $m "$DATA/$f.log" 2>&1 | strip | tail -n +2; done; done + echo; echo "How much heap does the 200 MB file need? (readAllLines, then readString)" + for x in 256m 512m 1g 2g; do + for m in readAllLines readString; do + printf '%-13s -Xmx%-5s ' "$m" "$x"; "$J" -Xmx$x -cp $CP com.ankurm.nio2.ReadAllFail $m "$DATA/medium.log" 2>&1 | strip | tail -n 1; done; done; } > "$OUT/05-read-all-fail.txt" + +echo "==> handles" +"$J" -cp $CP com.ankurm.nio2.HandleLeakDemo 2>&1 | strip > "$OUT/06-handle-leak.txt" +{ echo '$ ulimit -n 64; java HandleExhaustion'; (ulimit -n 64; "$J" -cp $CP com.ankurm.nio2.HandleExhaustion 2>&1 | strip); } > "$OUT/07-handle-exhaustion.txt" + +echo "==> walk / find" +"$J" -cp $CP com.ankurm.nio2.WalkFindDemo 2>&1 | strip > "$OUT/08-walk-find.txt" +"$J" -cp $CP com.ankurm.nio2.LinesCharsetDemo 2>&1 | strip > "$OUT/09-lines-charset.txt" + +echo "==> WatchService" +"$J" -cp $CP com.ankurm.nio2.WatchDemo 2>&1 | strip > "$OUT/10-watch.txt" +{ echo '$ java -Djdk.nio.file.WatchService.maxEventsPerPoll=10000 WatchDemo (scenario 6 only)' + "$J" -Djdk.nio.file.WatchService.maxEventsPerPoll=10000 -cp $CP com.ankurm.nio2.WatchDemo 2>&1 | strip | sed -n '/^6\./,$p' + echo; echo "Scenario 3 (50 appends, no polling) over five more runs:" + for i in 1 2 3 4 5; do "$J" -cp $CP com.ankurm.nio2.WatchDemo 2>&1 | strip | sed -n '/^3\./,/^4\./p' | sed '2!d'; done; } > "$OUT/11-watch-tuned-and-repeat.txt" +unzip -p "$JDK25_HOME/lib/src.zip" java.base/sun/nio/fs/AbstractWatchKey.java | sed -n '/final void signalEvent/,/signal();/p' > "$OUT/12-watchkey-source-excerpt.txt" + +echo "==> memory-mapped" +{ echo '$ java -Xmx64m MappedDemo big.log'; "$J" -Xmx64m -cp $CP com.ankurm.nio2.MappedDemo "$DATA/big.log" 2>&1 | strip; } > "$OUT/13-mapped.txt" + +cp target/surefire-reports/com.ankurm.nio2.Nio2BehaviourTest.txt "$OUT/14-tests.txt" +{ echo '$ df -h (data directory, after deleting is done by the exit trap; shown here before)'; df -h "$DATA" | sed 's/ */ /g'; } >> "$OUT/01-environment.txt" +echo Done diff --git a/nio2/src/main/java/com/ankurm/nio2/HandleExhaustion.java b/nio2/src/main/java/com/ankurm/nio2/HandleExhaustion.java new file mode 100644 index 0000000..5a64a1c --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/HandleExhaustion.java @@ -0,0 +1,27 @@ +package com.ankurm.nio2; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Stream; + +/** Run under a low `ulimit -n` to see what a leak turns into. */ +public final class HandleExhaustion { + private HandleExhaustion() {} + + public static void main(String[] a) throws IOException { + Path tmp = Files.createTempDirectory("nio2-exhaust"); + Path tree = HandleLeakDemo.makeTree(tmp.resolve("tree")); + List> leaked = new ArrayList<>(); + int i = 0; + try { + for (; ; i++) { Stream s = Files.walk(tree); s.skip(3).findFirst(); leaked.add(s); } + } catch (Exception e) { + int n = leaked.size(); + System.out.printf("failed on call %d after %d leaked walks%n%s: %s%n", i + 1, n, + e.getClass().getName(), e.getMessage() == null ? "" : e.getMessage().replace(tmp.toString(), "")); + } + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/HandleLeakDemo.java b/nio2/src/main/java/com/ankurm/nio2/HandleLeakDemo.java new file mode 100644 index 0000000..ea7d689 --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/HandleLeakDemo.java @@ -0,0 +1,68 @@ +package com.ankurm.nio2; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Stream; + +/** Counts this process's open file descriptors (Linux /proc/self/fd) around unclosed Files.lines / Files.walk. */ +public final class HandleLeakDemo { + private HandleLeakDemo() {} + + public static int openFds() throws IOException { + try (Stream s = Files.list(Path.of("/proc/self/fd"))) { return (int) s.count(); } + } + + static Path makeTree(Path root) throws IOException { + for (int i = 0; i < 5; i++) { + Path d = Files.createDirectories(root.resolve("d" + i).resolve("sub")); + Files.writeString(d.resolve("f.txt"), "x\n"); + Files.writeString(root.resolve("d" + i).resolve("g.txt"), "y\n"); + } + return root; + } + + public static void main(String[] a) throws Exception { + Path tmp = Files.createTempDirectory("nio2-leak"); + Path file = Files.writeString(tmp.resolve("one.txt"), "a\nb\nc\n"); + Path tree = makeTree(tmp.resolve("tree")); + final int N = 50; + int base = openFds(); + System.out.println("open fds at start: " + base); + + // 0. Files.walk fully consumed but not closed, run first so later leaks do not hide it. + for (int i = 0; i < N; i++) { Files.walk(tree).count(); } + System.out.printf("after %d fully consumed, unclosed Files.walk: +%d fds%n", N, openFds() - base); + + // 1. Files.lines, opened and read partially, never closed. References are kept (as a field or cache would). + List> keep = new ArrayList<>(); + for (int i = 0; i < N; i++) { Stream s = Files.lines(file); s.findFirst(); keep.add(s); } + System.out.printf("after %d unclosed Files.lines + findFirst(): +%d fds%n", N, openFds() - base); + keep.forEach(Stream::close); + System.out.printf("after closing them: +%d fds%n", openFds() - base); + + // 2. Files.lines in try-with-resources, same work. + for (int i = 0; i < N; i++) { try (Stream s = Files.lines(file)) { s.findFirst(); } } + System.out.printf("after %d try-with-resources Files.lines: +%d fds%n", N, openFds() - base); + + // 3. Files.walk, findFirst(), never closed. + List> walks = new ArrayList<>(); + for (int i = 0; i < N; i++) { Stream s = Files.walk(tree); s.skip(3).findFirst(); walks.add(s); } + System.out.printf("after %d unclosed Files.walk + findFirst(): +%d fds%n", N, openFds() - base); + walks.forEach(Stream::close); + System.out.printf("after closing them: +%d fds%n", openFds() - base); + + // 3b. The same unclosed Files.lines, but nobody keeps a reference: does the garbage collector clean up? + for (int i = 0; i < N; i++) { Files.lines(file).findFirst(); } + System.out.printf("after %d unreferenced, unclosed Files.lines: +%d fds%n", N, openFds() - base); + System.gc(); Thread.sleep(500); + System.out.printf("after System.gc(): +%d fds%n", openFds() - base); + for (int i = 0; i < N; i++) { Files.walk(tree).skip(3).findFirst(); } + System.out.printf("after %d unreferenced, unclosed Files.walk: +%d fds%n", N, openFds() - base); + System.gc(); Thread.sleep(500); + System.out.printf("after System.gc(): +%d fds%n", openFds() - base); + + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/HeapProbe.java b/nio2/src/main/java/com/ankurm/nio2/HeapProbe.java new file mode 100644 index 0000000..2aa97d8 --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/HeapProbe.java @@ -0,0 +1,48 @@ +package com.ankurm.nio2; + +import java.lang.management.GarbageCollectorMXBean; +import java.lang.management.ManagementFactory; +import java.nio.file.Files; +import java.nio.file.Path; + +/** Samples used heap on a daemon thread and reports peak-used, max heap and GC counts. */ +public final class HeapProbe implements AutoCloseable { + private volatile boolean run = true; + private volatile long peak; + private final Thread t; + + public HeapProbe() { + t = new Thread(() -> { + Runtime rt = Runtime.getRuntime(); + while (run) { + long used = rt.totalMemory() - rt.freeMemory(); + if (used > peak) peak = used; + try { Thread.sleep(20); } catch (InterruptedException e) { return; } + } + }, "heap-probe"); + t.setDaemon(true); + t.start(); + } + + public long peakUsed() { return peak; } + + public String report() { + long gcs = 0, ms = 0; + for (GarbageCollectorMXBean b : ManagementFactory.getGarbageCollectorMXBeans()) { + gcs += b.getCollectionCount(); ms += b.getCollectionTime(); + } + return String.format("heap: -Xmx=%d MiB, peak used (sampled every 20 ms)=%d MiB, GC collections=%d, GC time=%d ms", + Runtime.getRuntime().maxMemory() >> 20, peak >> 20, gcs, ms); + } + + @Override public void close() { run = false; } + + /** Resident set size of this process in MiB (Linux only), or -1. */ + public static long rssMiB() { + try { + for (String l : Files.readAllLines(Path.of("/proc/self/status"))) + if (l.startsWith("VmRSS:")) return Long.parseLong(l.replaceAll("\\D+", "")) >> 10; + } catch (Exception ignored) { } + return -1; + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/LinesCharsetDemo.java b/nio2/src/main/java/com/ankurm/nio2/LinesCharsetDemo.java new file mode 100644 index 0000000..7074f77 --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/LinesCharsetDemo.java @@ -0,0 +1,36 @@ +package com.ankurm.nio2; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.stream.Stream; + +/** Files.lines is lazy, so a bad byte fails in the middle of the pipeline, as an unchecked exception. */ +public final class LinesCharsetDemo { + private LinesCharsetDemo() {} + + public static void main(String[] a) throws IOException { + Path p = Files.createTempFile("nio2-charset", ".txt"); + byte[] ok = "good line 1\ngood line 2\n".getBytes(StandardCharsets.UTF_8); + byte[] bad = {(byte) 0xE9, '\n'}; // 0xE9 is 'e-acute' in ISO-8859-1 but not valid UTF-8 here + Files.write(p, ok); + Files.write(p, bad, java.nio.file.StandardOpenOption.APPEND); + int[] seen = {0}; + try (Stream s = Files.lines(p)) { + s.forEach(l -> seen[0]++); + } catch (UncheckedIOException e) { + System.out.println("lines processed before failure: " + seen[0]); + System.out.println(e.getClass().getSimpleName() + " caused by " + e.getCause()); + } + try (Stream s = Files.lines(p, StandardCharsets.ISO_8859_1)) { + System.out.println("with ISO_8859_1: " + s.count() + " lines, no error"); + } + try { + Files.readAllLines(p); + } catch (java.nio.charset.MalformedInputException e) { + System.out.println("readAllLines throws checked " + e.getClass().getSimpleName()); + } + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/LogFile.java b/nio2/src/main/java/com/ankurm/nio2/LogFile.java new file mode 100644 index 0000000..d427552 --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/LogFile.java @@ -0,0 +1,46 @@ +package com.ankurm.nio2; + +import java.io.BufferedWriter; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; + +/** Generates a deterministic synthetic log file of (at least) a requested size. One line = one record. */ +public final class LogFile { + private LogFile() {} + + public record Stats(long bytes, long lines, long errors) {} + + /** Writes whole lines until at least {@code targetBytes} have been written. */ + public static Stats generate(Path target, long targetBytes) throws IOException { + long seed = 42, bytes = 0, lines = 0, errors = 0; + try (BufferedWriter w = Files.newBufferedWriter(target)) { + StringBuilder sb = new StringBuilder(200); + while (bytes < targetBytes) { + seed = seed * 6364136223846793005L + 1442695040888963407L; + int r = (int) (seed >>> 33); + boolean error = r % 50 == 0; + if (error) errors++; + sb.setLength(0); + sb.append("2026-01-01T00:00:").append(String.format("%02d", (lines / 1000) % 60)) + .append(' ').append(error ? "ERROR" : "INFO ") + .append(" service-").append(r % 20) + .append(" request handled id=").append(lines) + .append(" user=").append(Integer.toHexString(r)) + .append(" payload=").append("x".repeat(40 + r % 40)).append('\n'); + w.write(sb.toString()); + bytes += sb.length(); // ASCII only, so chars == bytes + lines++; + } + } + return new Stats(bytes, lines, errors); + } + + public static void main(String[] a) throws IOException { + Path p = Path.of(a[0]); + long t0 = System.nanoTime(); + Stats s = generate(p, Long.parseLong(a[1])); + System.out.printf("generated %s: %,d bytes, %,d lines, %,d ERROR lines in %.1f s%n", + p.getFileName(), s.bytes(), s.lines(), s.errors(), (System.nanoTime() - t0) / 1e9); + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/MappedDemo.java b/nio2/src/main/java/com/ankurm/nio2/MappedDemo.java new file mode 100644 index 0000000..3638514 --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/MappedDemo.java @@ -0,0 +1,67 @@ +package com.ankurm.nio2; + +import java.io.IOException; +import java.lang.foreign.Arena; +import java.lang.foreign.MemorySegment; +import java.lang.foreign.ValueLayout; +import java.nio.MappedByteBuffer; +import java.nio.channels.FileChannel; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.StandardOpenOption; + +/** Memory-mapped reading: classic MappedByteBuffer (2 GiB cap per mapping) vs FileChannel.map with an Arena. */ +public final class MappedDemo { + private MappedDemo() {} + + /** Counts '\n' with an Arena-backed mapping of the whole file: no 2 GiB limit, no heap copy. */ + public static long countNewlinesArena(Path p) throws IOException { + try (Arena arena = Arena.ofConfined(); + FileChannel ch = FileChannel.open(p, StandardOpenOption.READ)) { + MemorySegment seg = ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size(), arena); + long n = 0; + for (long i = 0, len = seg.byteSize(); i < len; i++) if (seg.get(ValueLayout.JAVA_BYTE, i) == '\n') n++; + return n; + } + } + + /** Same with the classic API, in windows of at most 1 GiB. */ + public static long countNewlinesWindows(Path p) throws IOException { + long n = 0; + try (FileChannel ch = FileChannel.open(p, StandardOpenOption.READ)) { + long size = ch.size(), win = 1L << 30; + for (long off = 0; off < size; off += win) { + MappedByteBuffer mb = ch.map(FileChannel.MapMode.READ_ONLY, off, Math.min(win, size - off)); + for (int i = 0, lim = mb.limit(); i < lim; i++) if (mb.get(i) == '\n') n++; + } + } + return n; + } + + public static void main(String[] a) throws IOException { + Path p = Path.of(a[0]); + long size = Files.size(p); + System.out.printf("file %s: %,d bytes, -Xmx=%d MiB, RSS at start=%d MiB%n", p.getFileName(), size, + Runtime.getRuntime().maxMemory() >> 20, HeapProbe.rssMiB()); + try (FileChannel ch = FileChannel.open(p, StandardOpenOption.READ)) { + ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size()); + System.out.println("classic map of the whole file: succeeded"); + } catch (IllegalArgumentException e) { + System.out.println("classic map of the whole file: IllegalArgumentException: " + e.getMessage()); + } + try (HeapProbe probe = new HeapProbe()) { + long t0 = System.nanoTime(); + long n2 = countNewlinesArena(p); + System.out.printf("Arena + MemorySegment: %,d newlines in %.1f s; RSS after the Arena was closed=%d MiB%n", + n2, (System.nanoTime() - t0) / 1e9, HeapProbe.rssMiB()); + t0 = System.nanoTime(); + long n1 = countNewlinesWindows(p); + System.out.printf("classic 1 GiB windows: %,d newlines in %.1f s; RSS right after=%d MiB%n", + n1, (System.nanoTime() - t0) / 1e9, HeapProbe.rssMiB()); + System.out.println(probe.report()); + System.gc(); + try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } + System.out.println("RSS one second after System.gc()=" + HeapProbe.rssMiB() + " MiB"); + } + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/ReadAllFail.java b/nio2/src/main/java/com/ankurm/nio2/ReadAllFail.java new file mode 100644 index 0000000..097ff79 --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/ReadAllFail.java @@ -0,0 +1,28 @@ +package com.ankurm.nio2; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; + +/** The convenience methods that load the whole file. Run with a small -Xmx to see the failure. */ +public final class ReadAllFail { + private ReadAllFail() {} + + public static void main(String[] a) throws IOException { + Path p = Path.of(a[1]); + System.out.printf("%s on %s (%,d bytes), -Xmx=%d MiB%n", a[0], p.getFileName(), Files.size(p), + Runtime.getRuntime().maxMemory() >> 20); + try { + Object r = switch (a[0]) { + case "readAllLines" -> Files.readAllLines(p); + case "readString" -> Files.readString(p); + case "readAllBytes" -> Files.readAllBytes(p); + default -> throw new IllegalArgumentException(a[0]); + }; + System.out.println("survived: " + (r instanceof java.util.List l ? l.size() + " lines" + : r instanceof String s ? s.length() + " chars" : ((byte[]) r).length + " bytes")); + } catch (OutOfMemoryError e) { + System.out.println("OutOfMemoryError: " + e.getMessage()); + } + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/StreamLines.java b/nio2/src/main/java/com/ankurm/nio2/StreamLines.java new file mode 100644 index 0000000..9bc2e4f --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/StreamLines.java @@ -0,0 +1,58 @@ +package com.ankurm.nio2; + +import java.io.BufferedReader; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Map; +import java.util.TreeMap; +import java.util.stream.Stream; + +/** Processes a file of any size in constant heap: count lines, ERROR lines, and lines per service. */ +public final class StreamLines { + private StreamLines() {} + + public record Result(long lines, long errors, Map perService) {} + + /** Files.lines: lazy Stream; MUST be closed (it owns a file handle). */ + public static Result withFilesLines(Path p) throws IOException { + long[] c = new long[2]; + Map per = new TreeMap<>(); + try (Stream s = Files.lines(p)) { + s.forEach(l -> tally(l, c, per)); + } + return new Result(c[0], c[1], per); + } + + /** Files.newBufferedReader: the same thing with a plain loop. */ + public static Result withBufferedReader(Path p) throws IOException { + long[] c = new long[2]; + Map per = new TreeMap<>(); + try (BufferedReader r = Files.newBufferedReader(p)) { + for (String l = r.readLine(); l != null; l = r.readLine()) tally(l, c, per); + } + return new Result(c[0], c[1], per); + } + + private static void tally(String line, long[] c, Map per) { + c[0]++; + if (line.startsWith("ERROR", 20)) c[1]++; + int i = line.indexOf("service-"); + per.merge(line.substring(i, line.indexOf(' ', i)), 1L, Long::sum); + } + + public static void main(String[] a) throws IOException { + Path p = Path.of(a[1]); + long t0 = System.nanoTime(); + Result r; + try (HeapProbe probe = new HeapProbe()) { + r = a[0].equals("lines") ? withFilesLines(p) : withBufferedReader(p); + double s = (System.nanoTime() - t0) / 1e9; + System.out.printf("%s over %s (%,d bytes)%n", a[0].equals("lines") ? "Files.lines" : "Files.newBufferedReader", + p.getFileName(), Files.size(p)); + System.out.printf("lines=%,d errors=%,d services=%d%n", r.lines(), r.errors(), r.perService().size()); + System.out.printf("elapsed %.1f s%n", s); + System.out.println(probe.report()); + } + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/WalkFindDemo.java b/nio2/src/main/java/com/ankurm/nio2/WalkFindDemo.java new file mode 100644 index 0000000..2b7f4b4 --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/WalkFindDemo.java @@ -0,0 +1,64 @@ +package com.ankurm.nio2; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.file.*; +import java.nio.file.attribute.BasicFileAttributes; +import java.util.List; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +/** Files.list / walk / find / walkFileTree / DirectoryStream on a small tree. */ +public final class WalkFindDemo { + private WalkFindDemo() {} + + static List rel(Path root, Stream s) { + try (s) { return s.map(p -> root.relativize(p).toString()).sorted().collect(Collectors.toList()); } + } + + public static void main(String[] a) throws IOException { + Path root = Files.createTempDirectory("nio2-walk"); + Files.createDirectories(root.resolve("src/main")); + Files.createDirectories(root.resolve("src/test")); + Files.createDirectories(root.resolve(".git/objects")); + Files.writeString(root.resolve("README.md"), "# demo\n"); + Files.writeString(root.resolve("src/main/App.java"), "class App {}\n"); + Files.writeString(root.resolve("src/test/AppTest.java"), "class AppTest {}\n" + "x".repeat(2000)); + Files.writeString(root.resolve(".git/objects/blob.bin"), "ignored"); + + System.out.println("Files.list (one level): " + rel(root, Files.list(root))); + System.out.println("Files.walk (everything): " + rel(root, Files.walk(root))); + System.out.println("Files.walk maxDepth=1: " + rel(root, Files.walk(root, 1))); + System.out.println("Files.find *.java: " + rel(root, + Files.find(root, 10, (p, at) -> at.isRegularFile() && p.toString().endsWith(".java")))); + System.out.println("Files.find size > 1000: " + rel(root, + Files.find(root, 10, (p, at) -> at.isRegularFile() && at.size() > 1000))); + + List seen = new java.util.ArrayList<>(); + Files.walkFileTree(root, new SimpleFileVisitor<>() { + @Override public FileVisitResult preVisitDirectory(Path d, BasicFileAttributes at) { + return d.getFileName().toString().equals(".git") ? FileVisitResult.SKIP_SUBTREE : FileVisitResult.CONTINUE; + } + @Override public FileVisitResult visitFile(Path f, BasicFileAttributes at) { + seen.add(root.relativize(f).toString()); return FileVisitResult.CONTINUE; + } + }); + java.util.Collections.sort(seen); + System.out.println("walkFileTree skipping .git: " + seen); + + List glob = new java.util.ArrayList<>(); + try (DirectoryStream ds = Files.newDirectoryStream(root.resolve("src/main"), "*.{java,kt}")) { + ds.forEach(p -> glob.add(p.getFileName().toString())); + } + System.out.println("DirectoryStream glob: " + glob); + + // A symlink loop: walk() does not follow links by default; with FOLLOW_LINKS it detects the cycle. + Files.createSymbolicLink(root.resolve("src/loop"), root); + System.out.println("walk with a symlink loop, default: " + rel(root, Files.walk(root)).size() + " entries, no error"); + try (Stream s = Files.walk(root, FileVisitOption.FOLLOW_LINKS)) { + s.forEach(p -> { }); + } catch (UncheckedIOException e) { + System.out.println("walk with FOLLOW_LINKS: " + e.getCause().getClass().getSimpleName()); + } + } +} diff --git a/nio2/src/main/java/com/ankurm/nio2/WatchDemo.java b/nio2/src/main/java/com/ankurm/nio2/WatchDemo.java new file mode 100644 index 0000000..0958ace --- /dev/null +++ b/nio2/src/main/java/com/ankurm/nio2/WatchDemo.java @@ -0,0 +1,82 @@ +package com.ankurm.nio2; + +import java.io.IOException; +import java.nio.file.*; +import java.util.*; +import java.util.concurrent.TimeUnit; + +import static java.nio.file.StandardWatchEventKinds.*; + +/** What WatchService actually delivers on this platform. Every scenario prints what it saw. */ +public final class WatchDemo { + private WatchDemo() {} + + /** Drains the key until no new event arrives for quietMs; returns kind -> total repeat-count. */ + static Map drain(WatchService ws, long quietMs, boolean verbose) throws Exception { + Map tally = new TreeMap<>(); + int batches = 0; + WatchKey k; + while ((k = ws.poll(quietMs, TimeUnit.MILLISECONDS)) != null) { + batches++; + List> evs = k.pollEvents(); + for (WatchEvent e : evs) { + tally.merge(e.kind().name(), e.count(), Integer::sum); + if (verbose) System.out.println(" " + e.kind().name() + " count=" + e.count() + " context=" + e.context()); + } + k.reset(); + } + tally.put("(batches)", batches); + return tally; + } + + public static void main(String[] a) throws Exception { + Path dir = Files.createTempDirectory("nio2-watch"); + try (WatchService ws = FileSystems.getDefault().newWatchService()) { + dir.register(ws, ENTRY_CREATE, ENTRY_MODIFY, ENTRY_DELETE); + + System.out.println("1. create a file, write it twice, delete it (events polled afterwards):"); + Path f = dir.resolve("a.txt"); + Files.writeString(f, "one"); + Files.writeString(f, "two", StandardOpenOption.APPEND); + Files.delete(f); + drain(ws, 300, true); + + System.out.println("2. Files.writeString of 100 KiB to a new file (one call):"); + Files.writeString(dir.resolve("b.txt"), "y".repeat(100 * 1024)); + drain(ws, 300, true); + + System.out.println("3. 50 appends to one file with no polling in between, then poll:"); + Path g = dir.resolve("g.txt"); + Files.writeString(g, ""); + drain(ws, 300, false); + for (int i = 0; i < 50; i++) Files.writeString(g, "line\n", StandardOpenOption.APPEND); + System.out.println(" " + drain(ws, 300, false)); + + System.out.println("4. Files.move (atomic rename) of g.txt to h.txt:"); + Files.move(g, dir.resolve("h.txt"), StandardCopyOption.ATOMIC_MOVE); + drain(ws, 300, true); + + System.out.println("5. a file created inside a NEW subdirectory (subdirectory not registered):"); + Path sub = Files.createDirectory(dir.resolve("sub")); + Files.writeString(sub.resolve("deep.txt"), "hello"); + drain(ws, 300, true); + + System.out.println("6. burst: create N empty files in a fresh directory before polling once" + + " (jdk.nio.file.WatchService.maxEventsPerPoll=" + System.getProperty("jdk.nio.file.WatchService.maxEventsPerPoll", "default") + "):"); + for (int n : new int[] {100, 512, 513, 5000}) { + Path burst = Files.createDirectory(dir.resolve("burst" + n)); + try (WatchService ws2 = FileSystems.getDefault().newWatchService()) { + burst.register(ws2, ENTRY_CREATE); + for (int i = 0; i < n; i++) Files.createFile(burst.resolve("f" + i)); + Thread.sleep(500); + WatchKey k = ws2.poll(1, TimeUnit.SECONDS); + List> evs = k.pollEvents(); + long creates = evs.stream().filter(e -> e.kind() == ENTRY_CREATE).count(); + long overflow = evs.stream().filter(e -> e.kind() == OVERFLOW).mapToInt(WatchEvent::count).sum(); + System.out.printf(" N=%-5d delivered %-4d events: ENTRY_CREATE=%-4d OVERFLOW(count)=%d, files on disk=%d%n", + n, evs.size(), creates, overflow, Files.list(burst).count()); + } + } + } + } +} diff --git a/nio2/src/test/java/com/ankurm/nio2/Nio2BehaviourTest.java b/nio2/src/test/java/com/ankurm/nio2/Nio2BehaviourTest.java new file mode 100644 index 0000000..7d3199b --- /dev/null +++ b/nio2/src/test/java/com/ankurm/nio2/Nio2BehaviourTest.java @@ -0,0 +1,183 @@ +package com.ankurm.nio2; + +import org.junit.jupiter.api.Assumptions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.io.IOException; +import java.io.RandomAccessFile; +import java.io.UncheckedIOException; +import java.nio.channels.FileChannel; +import java.nio.charset.MalformedInputException; +import java.nio.file.*; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.TimeUnit; +import java.util.stream.Stream; + +import static java.nio.file.StandardWatchEventKinds.*; +import static org.junit.jupiter.api.Assertions.*; + +class Nio2BehaviourTest { + + @TempDir Path tmp; + + private Path sample() throws IOException { + Path p = tmp.resolve("sample.log"); + LogFile.generate(p, 2_000_000); + return p; + } + + @Test + void streamingMatchesReadAllLines() throws IOException { + Path p = sample(); + List all = Files.readAllLines(p); + assertEquals(all.size(), StreamLines.withFilesLines(p).lines()); + assertEquals(all.size(), StreamLines.withBufferedReader(p).lines()); + assertEquals(all.stream().filter(l -> l.startsWith("ERROR", 20)).count(), StreamLines.withFilesLines(p).errors()); + } + + @Test + void generatorStatsMatchTheFile() throws IOException { + Path p = tmp.resolve("g.log"); + LogFile.Stats s = LogFile.generate(p, 500_000); + assertEquals(Files.size(p), s.bytes()); + try (Stream l = Files.lines(p)) { assertEquals(s.lines(), l.count()); } + } + + @Test + void mappedCountersAgreeWithLineCount() throws IOException { + Path p = sample(); + long lines = StreamLines.withFilesLines(p).lines(); + assertEquals(lines, MappedDemo.countNewlinesArena(p)); + assertEquals(lines, MappedDemo.countNewlinesWindows(p)); + } + + @Test + void classicMapRejectsMoreThan2GiBButArenaMapsIt() throws IOException { + Path sparse = tmp.resolve("sparse.bin"); + try (RandomAccessFile f = new RandomAccessFile(sparse.toFile(), "rw")) { f.setLength(3L << 30); } // sparse: no real disk used + try (FileChannel ch = FileChannel.open(sparse, StandardOpenOption.READ)) { + assertThrows(IllegalArgumentException.class, () -> ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size())); + } + try (var arena = java.lang.foreign.Arena.ofConfined(); FileChannel ch = FileChannel.open(sparse, StandardOpenOption.READ)) { + var seg = ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size(), arena); + assertEquals(3L << 30, seg.byteSize()); + } + } + + private static boolean linux() { return Files.isDirectory(Path.of("/proc/self/fd")); } + + @Test + void unclosedFilesLinesHoldsAFileDescriptorUntilClosed() throws IOException { + Assumptions.assumeTrue(linux()); + Path p = sample(); + int base = HandleLeakDemo.openFds(); + List> keep = new ArrayList<>(); + for (int i = 0; i < 10; i++) { Stream s = Files.lines(p); s.findFirst(); keep.add(s); } + assertEquals(base + 10, HandleLeakDemo.openFds()); + keep.forEach(Stream::close); + assertEquals(base, HandleLeakDemo.openFds()); + for (int i = 0; i < 10; i++) { try (Stream s = Files.lines(p)) { s.findFirst(); } } + assertEquals(base, HandleLeakDemo.openFds()); + } + + @Test + void walkStoppedEarlyHoldsDirectoriesOpenButFullyConsumedWalkDoesNot() throws IOException { + Assumptions.assumeTrue(linux()); + Path tree = HandleLeakDemo.makeTree(tmp.resolve("tree")); + int base = HandleLeakDemo.openFds(); + Stream early = Files.walk(tree); + early.skip(3).findFirst(); + assertTrue(HandleLeakDemo.openFds() > base, "an unfinished walk should hold directory handles"); + early.close(); + assertEquals(base, HandleLeakDemo.openFds()); + Files.walk(tree).count(); + assertEquals(base, HandleLeakDemo.openFds()); + } + + @Test + void badByteSurfacesLazilyAsUncheckedIOException() throws IOException { + Path p = tmp.resolve("bad.txt"); + Files.write(p, new byte[] {'o', 'k', '\n', (byte) 0xE9, '\n'}); + try (Stream s = Files.lines(p)) { + UncheckedIOException e = assertThrows(UncheckedIOException.class, () -> s.forEach(l -> { })); + assertInstanceOf(MalformedInputException.class, e.getCause()); + } + assertThrows(MalformedInputException.class, () -> Files.readAllLines(p)); + } + + @Test + void walkFindAndSkipSubtree() throws IOException { + Files.createDirectories(tmp.resolve("a/b")); + Files.writeString(tmp.resolve("a/b/x.java"), "x"); + Files.writeString(tmp.resolve("a/y.txt"), "y"); + try (Stream s = Files.find(tmp, 5, (p, at) -> at.isRegularFile() && p.toString().endsWith(".java"))) { + assertEquals(1, s.count()); + } + try (Stream s = Files.walk(tmp, 1)) { assertEquals(2, s.count()); } // tmp itself + "a" + } + + @Test + void watchServiceDeliversCreateAndCoalescesRepeatedModify() throws Exception { + try (WatchService ws = FileSystems.getDefault().newWatchService()) { + tmp.register(ws, ENTRY_CREATE, ENTRY_MODIFY); + Path f = tmp.resolve("w.txt"); + Files.writeString(f, "x".repeat(100 * 1024)); + Thread.sleep(300); + WatchKey k = ws.poll(2, TimeUnit.SECONDS); + assertNotNull(k); + List> evs = k.pollEvents(); + assertEquals(ENTRY_CREATE, evs.get(0).kind()); + assertEquals(f.getFileName(), evs.get(0).context()); + assertTrue(evs.size() <= 2, "MODIFY events for one file collapse into one event with a count: " + evs); + } + } + + @Test + void watchKeyOverflowsPast512PendingEventsAndDropsThem() throws Exception { + Assumptions.assumeTrue(System.getProperty("jdk.nio.file.WatchService.maxEventsPerPoll") == null); + try (WatchService ws = FileSystems.getDefault().newWatchService()) { + tmp.register(ws, ENTRY_CREATE); + for (int i = 0; i < 513; i++) Files.createFile(tmp.resolve("f" + i)); + Thread.sleep(500); + WatchKey k = ws.poll(2, TimeUnit.SECONDS); + List> evs = k.pollEvents(); + assertEquals(1, evs.size()); + assertEquals(OVERFLOW, evs.get(0).kind()); + } + } + + @Test + void watchServiceDoesNotWatchSubdirectoriesRecursively() throws Exception { + try (WatchService ws = FileSystems.getDefault().newWatchService()) { + tmp.register(ws, ENTRY_CREATE); + Path sub = Files.createDirectory(tmp.resolve("sub")); + Thread.sleep(200); + WatchKey k = ws.poll(2, TimeUnit.SECONDS); + k.pollEvents(); k.reset(); + Files.writeString(sub.resolve("deep.txt"), "x"); + assertNull(ws.poll(500, TimeUnit.MILLISECONDS), "no event for a file in an unregistered subdirectory"); + } + } + + @Test + void readAllLinesFailsWhereStreamingSurvivesUnderSmallHeap() throws Exception { + // Runs the real programs in child JVMs so -Xmx is honoured. + Path p = tmp.resolve("m.log"); + LogFile.generate(p, 100_000_000); + String java = Path.of(System.getProperty("java.home"), "bin", "java").toString(); + String cp = System.getProperty("java.class.path"); + String streamOut = run(java, "-Xmx32m", "-cp", cp, "com.ankurm.nio2.StreamLines", "lines", p.toString()); + assertTrue(streamOut.contains("lines="), streamOut); + String failOut = run(java, "-Xmx32m", "-cp", cp, "com.ankurm.nio2.ReadAllFail", "readAllLines", p.toString()); + assertTrue(failOut.contains("OutOfMemoryError"), failOut); + } + + private static String run(String... cmd) throws Exception { + Process pr = new ProcessBuilder(cmd).redirectErrorStream(true).start(); + String s = new String(pr.getInputStream().readAllBytes()); + pr.waitFor(); + return s; + } +} diff --git a/pom.xml b/pom.xml index f1ae221..ec85326 100644 --- a/pom.xml +++ b/pom.xml @@ -27,6 +27,7 @@ serialization strings collections + nio2