|
| 1 | +# Parallel Range Processor |
| 2 | + |
| 3 | +[](https://github.com/j-util/parallel-range-processor/actions/workflows/ci.yml) |
| 4 | + |
| 5 | +A small Java 8 library for concurrently processing delimiter-framed records |
| 6 | +from one large, range-addressable byte source without materializing the complete |
| 7 | +source. It is an orchestration layer over |
| 8 | +[`inputstream-processor-core`](https://github.com/j-util/inputstream-processor-core), |
| 9 | +which remains responsible for parsing and consumer-call counting. |
| 10 | + |
| 11 | +## Requirements and installation |
| 12 | + |
| 13 | +Parallel Range Processor requires Java 8 or later. |
| 14 | + |
| 15 | +```xml |
| 16 | +<dependency> |
| 17 | + <groupId>io.github.j-util</groupId> |
| 18 | + <artifactId>parallel-range-processor</artifactId> |
| 19 | + <version>1.0.0</version> |
| 20 | +</dependency> |
| 21 | +``` |
| 22 | + |
| 23 | +The artifact has one compile-time dependency: |
| 24 | +`io.github.j-util:inputstream-processor-core:1.0.0`. There are no storage SDK or |
| 25 | +concurrency-framework dependencies. |
| 26 | + |
| 27 | +## Why a range source is required |
| 28 | + |
| 29 | +An ordinary `InputStream` has one cursor. Splitting byte offsets does not by |
| 30 | +itself establish record boundaries: a range can start or end in the middle of a |
| 31 | +record. `RangeSource` adds the two capabilities needed by this algorithm: |
| 32 | + |
| 33 | +- a known total byte size; and |
| 34 | +- independent streams for exact half-open ranges such as `[1000, 2000)`. |
| 35 | + |
| 36 | +The source must remain logically stable for one processing operation. Each |
| 37 | +`openRange` caller owns and closes the returned stream. The processor closes all |
| 38 | +streams it opens and also enforces the requested upper bound locally. |
| 39 | + |
| 40 | +Unknown-size sources are outside V1 because they require a different scheduling |
| 41 | +strategy based on chunk creation and EOF discovery. |
| 42 | + |
| 43 | +## Processing model |
| 44 | + |
| 45 | +```text |
| 46 | +RangeSource |
| 47 | + | |
| 48 | + +-- [0, a) ---- worker parser --+ |
| 49 | + +-- [a, b) ---- worker parser --+--> concurrent consumer calls |
| 50 | + +-- [b, S) ---- worker parser --+ |
| 51 | + | |
| 52 | + +-- ordered boundary fragments |
| 53 | + | |
| 54 | + +--> final core parser --> consumer |
| 55 | +``` |
| 56 | + |
| 57 | +For source size `S` and requested parallelism `N`, the processor creates up to |
| 58 | +`N` contiguous, non-empty ranges that cover `[0, S)` exactly once. Ranges never |
| 59 | +overlap, workers never read outside their assigned ranges, and there is one |
| 60 | +top-level task per actual range—not one task per record. |
| 61 | + |
| 62 | +Within a worker, a segmented framing stream buffers the current record only |
| 63 | +until it encounters the configured delimiter. It: |
| 64 | + |
| 65 | +1. retains an ambiguous leading fragment for every range except the first; |
| 66 | +2. exposes complete delimiter-terminated records to that worker's independent |
| 67 | + `InputStreamProcessor`; |
| 68 | +3. retains an ambiguous trailing fragment for every range except the last; and |
| 69 | +4. treats bytes at source EOF as a complete final record even without a trailing |
| 70 | + delimiter. |
| 71 | + |
| 72 | +An intermediate range with no delimiter is represented as one middle fragment, |
| 73 | +not inferred from parser results or consumer-call counts. This is what permits a |
| 74 | +single record to span three, four, or more ranges. After all workers complete, |
| 75 | +only boundary fragments are concatenated in source/range order and processed by |
| 76 | +one additional core parser. Fragments are streamed from fixed-size segments; |
| 77 | +the reconstructed record is not first copied into one giant contiguous |
| 78 | +`byte[]`. |
| 79 | + |
| 80 | +Complete records are processed during the parallel phase and are not retained |
| 81 | +after the parser consumes them. The library performs no complete-source |
| 82 | +materialization. Temporary memory is proportional to reusable per-worker |
| 83 | +record-framing buffers plus retained boundary fragments. Very large records can |
| 84 | +therefore require temporary storage proportional to their size; the library |
| 85 | +does not claim zero materialization. |
| 86 | + |
| 87 | +## Local-file example |
| 88 | + |
| 89 | +```java |
| 90 | +import io.github.jutil.inputstreamprocessor.core.InputParser; |
| 91 | +import io.github.jutil.parallelrangeprocessor.FileRangeSource; |
| 92 | +import io.github.jutil.parallelrangeprocessor.ParallelProcessingResult; |
| 93 | +import io.github.jutil.parallelrangeprocessor.ParallelRangeProcessor; |
| 94 | +import io.github.jutil.parallelrangeprocessor.RecordDelimiter; |
| 95 | + |
| 96 | +import java.io.BufferedReader; |
| 97 | +import java.io.InputStreamReader; |
| 98 | +import java.nio.charset.StandardCharsets; |
| 99 | +import java.nio.file.Paths; |
| 100 | +import java.util.Queue; |
| 101 | +import java.util.concurrent.ConcurrentLinkedQueue; |
| 102 | +import java.util.concurrent.ExecutorService; |
| 103 | +import java.util.concurrent.Executors; |
| 104 | +import java.util.function.Supplier; |
| 105 | + |
| 106 | +Supplier<InputParser<String>> lineParserFactory = () -> (input, emit) -> { |
| 107 | + BufferedReader reader = new BufferedReader( |
| 108 | + new InputStreamReader(input, StandardCharsets.UTF_8) |
| 109 | + ); |
| 110 | + String line; |
| 111 | + while ((line = reader.readLine()) != null) { |
| 112 | + emit.accept(line); |
| 113 | + } |
| 114 | + // Do not close reader: the processor owns the underlying range stream. |
| 115 | + }; |
| 116 | + |
| 117 | +ExecutorService executor = Executors.newFixedThreadPool(4); |
| 118 | +Queue<String> records = new ConcurrentLinkedQueue<>(); |
| 119 | +try { |
| 120 | + ParallelRangeProcessor<String> processor = new ParallelRangeProcessor<>( |
| 121 | + 4, |
| 122 | + executor, |
| 123 | + lineParserFactory, |
| 124 | + RecordDelimiter.newline() |
| 125 | + ); |
| 126 | + |
| 127 | + ParallelProcessingResult result = processor.process( |
| 128 | + new FileRangeSource(Paths.get("records.ndjson")), |
| 129 | + records::add |
| 130 | + ); |
| 131 | + System.out.println("Processed items: " + result.getProcessedCount()); |
| 132 | +} finally { |
| 133 | + // The application owns executor lifecycle; the library never shuts it down. |
| 134 | + executor.shutdown(); |
| 135 | +} |
| 136 | +``` |
| 137 | + |
| 138 | +The parser factory is called once per actual range and, when boundary data |
| 139 | +exists, once for final reconstruction. Every returned parser must be non-null, |
| 140 | +independently usable, synchronous as required by `inputstream-processor-core`, |
| 141 | +and must consume its supplied complete-record stream through EOF. A parser that |
| 142 | +returns while complete record bytes remain causes processing to fail rather |
| 143 | +than silently losing records. |
| 144 | + |
| 145 | +## Concurrency and ordering |
| 146 | + |
| 147 | +- The caller supplies explicit parallelism and an `Executor`. Parallelism is |
| 148 | + never inferred from the executor. |
| 149 | +- The library creates no threads or pools and never shuts down the executor. |
| 150 | +- The same consumer can be called concurrently by several worker threads. For |
| 151 | + effective parallelism greater than one, the consumer must be thread-safe or |
| 152 | + otherwise tolerate concurrent invocation. Calls are not synchronized by the |
| 153 | + library. |
| 154 | +- Global consumer invocation order is unspecified. Complete in-range records |
| 155 | + can be consumed in any worker completion order. |
| 156 | +- Boundary reconstruction order is deterministic source order. Reconstructed |
| 157 | + records are consumed only after the parallel worker phase. |
| 158 | +- A request for parallelism greater than the byte size creates fewer ranges; |
| 159 | + zero-length worker ranges are never submitted. |
| 160 | + |
| 161 | +## Framing semantics and limitations |
| 162 | + |
| 163 | +V1 recognizes one byte delimiter. `RecordDelimiter.newline()` selects line feed |
| 164 | +(`0x0A`); `RecordDelimiter.singleByte(...)` supports another independently |
| 165 | +recognizable byte. Delimiter bytes are preserved in parser input. With a normal |
| 166 | +line parser, consecutive line feeds represent empty records, while a final |
| 167 | +trailing line feed does not invent an extra record after EOF. |
| 168 | + |
| 169 | +This layer establishes byte boundaries; it is not a CSV, JSON, or XML parser. |
| 170 | +It is suitable only when every separator byte unambiguously terminates a logical |
| 171 | +record. Unsupported inputs include: |
| 172 | + |
| 173 | +- CSV with multiline quoted fields; |
| 174 | +- JSON arrays or other stateful structured documents; |
| 175 | +- XML; |
| 176 | +- formats in which delimiter bytes can occur inside a record without locally |
| 177 | + detectable framing state; |
| 178 | +- multi-byte delimiters in V1; and |
| 179 | +- non-splittable compressed content such as a normal single gzip stream. |
| 180 | + |
| 181 | +CRLF text can use the line-feed delimiter: the carriage return remains in the |
| 182 | +original bytes immediately before the line feed, and a standard line reader |
| 183 | +handles it normally. This library makes no claim of support for arbitrary |
| 184 | +structured or compressed formats. |
| 185 | + |
| 186 | +## S3 and HTTP adapters |
| 187 | + |
| 188 | +The core stays storage-neutral. A custom `RangeSource` for S3 or HTTP should: |
| 189 | + |
| 190 | +1. obtain and return a stable content length from `size()`; |
| 191 | +2. translate `[fromInclusive, toExclusive)` into an inclusive wire range such |
| 192 | + as `Range: bytes=fromInclusive-(toExclusive - 1)`; |
| 193 | +3. require a partial-content response when appropriate and verify the returned |
| 194 | + range and length; |
| 195 | +4. return an independent response-body stream for every call; and |
| 196 | +5. keep the same object/version or validator stable for the complete processing |
| 197 | + operation. |
| 198 | + |
| 199 | +SDK clients, credentials, retries, consistency policy, and HTTP response |
| 200 | +validation belong in the application adapter or a separate adapter artifact, |
| 201 | +not in this library. |
| 202 | + |
| 203 | +## Results and failures |
| 204 | + |
| 205 | +`ParallelProcessingResult.getProcessedCount()` sums the |
| 206 | +`inputstream-processor-core` counts from all normally completed worker streams |
| 207 | +and the final boundary stream. It counts parser-emitted items whose consumer |
| 208 | +calls returned normally; it is never used to infer record framing. |
| 209 | + |
| 210 | +Source, parser, consumer, and executor failures do not disappear. Submitted |
| 211 | +workers are awaited before `process` returns or throws, so no worker is left |
| 212 | +calling the consumer after the operation has reported completion. A generic |
| 213 | +`Executor` has no reliable cancellation API, so already submitted tasks are |
| 214 | +allowed to finish. If several workers fail, the original failure from the |
| 215 | +lowest source-range ordinal is propagated regardless of completion order. |
| 216 | + |
| 217 | +If the waiting thread is interrupted, the processor still waits for submitted |
| 218 | +workers and then throws `InterruptedException`; a worker failure is attached as |
| 219 | +a suppressed exception. Completed consumer side effects are not rolled back, |
| 220 | +and no transactional guarantee is made. No result is returned after any |
| 221 | +failure. |
0 commit comments