Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Java concurrent import project: validated worker results before file publication

Last updated: 29 Sept 20264 min read
tutorial
IntermediateBy AITrove Editorial

A concurrent import can compute independent row results in workers while reserving one owner for validation and publication of the final snapshot.

Download Java source kit

Java 8+. The program uses JDK classes and requires no preview flags.

Workers return proposed values

Each worker parses one receipt row and returns an immutable result. It does not append to the published file or decrement a shared total. The owner collects every result, rejects duplicate references, and writes the snapshot only when the complete proposed set is valid.

This separates completed computation from publication. If one row fails, the old snapshot remains unchanged. Cancellation requests stop outstanding work, and executor termination is checked before the import operation returns. The program’s invalid-row case verifies that failure does not publish the successful row beside it.

The executor has two workers and a bounded pending queue. The input is also bounded by file size, row count and row length. A small queue by itself would not bound all memory if the caller eagerly materialized an unlimited file before submitting it.

Describe the deadline honestly

One monotonic budget covers submission and result waiting. Reading the owned input file and final filesystem publication remain separate synchronous stages, so this fixture does not promise an end-to-end service deadline or crash-durable persistence. A regular move publishes this single-writer snapshot; it is not a distributed transaction.

Working program

Java
import java.nio.file.*;
import java.nio.charset.StandardCharsets;
import java.util.*;
import java.util.concurrent.*;
public class ConcurrentReceiptImport {
    static class Row {final String reference;final int units;Row(String reference,int units){this.reference=reference;this.units=units;}}
    static Row parse(String line){
        if(line.length()>80)throw new IllegalArgumentException("Row length");String[] fields=line.split(",",-1);
        if(fields.length!=2||!fields[0].matches("R-[0-9]{1,8}"))throw new IllegalArgumentException("Receipt row");
        int units=Integer.parseInt(fields[1]);if(units<0||units>1000)throw new IllegalArgumentException("Unit range");return new Row(fields[0],units);
    }
    static void importFile(Path input,Path output)throws Exception{
        if(Files.size(input)>16_000)throw new IllegalArgumentException("Input size");List<String> lines=Files.readAllLines(input,StandardCharsets.UTF_8);if(lines.size()>6)throw new IllegalArgumentException("Row count");
        ThreadPoolExecutor workers=new ThreadPoolExecutor(2,2,0,TimeUnit.MILLISECONDS,new ArrayBlockingQueue<>(4));List<Future<Row>> pending=new ArrayList<>();long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(2);
        try{
            for(String line:lines)pending.add(workers.submit(()->parse(line)));
            List<String> snapshot=new ArrayList<>();Set<String> references=new HashSet<>();
            for(Future<Row> future:pending){long remaining=deadline-System.nanoTime();if(remaining<=0)throw new TimeoutException("Import budget");Row row=future.get(remaining,TimeUnit.NANOSECONDS);if(!references.add(row.reference))throw new IllegalArgumentException("Duplicate receipt");snapshot.add(row.reference+","+row.units);}
            Path staged=Files.createTempFile(output.getParent(),"snapshot-",".tmp");
            try{Files.write(staged,snapshot,StandardCharsets.UTF_8);Files.move(staged,output,StandardCopyOption.REPLACE_EXISTING);}finally{Files.deleteIfExists(staged);}
        }finally{for(Future<Row> future:pending)future.cancel(true);workers.shutdownNow();if(!workers.awaitTermination(2,TimeUnit.SECONDS))throw new IllegalStateException("Import workers alive");}
    }
    public static void main(String[] args)throws Exception{
        Path directory=Files.createTempDirectory("aitrove-import-"),input=directory.resolve("input.csv"),output=directory.resolve("snapshot.csv");
        try{
            Files.write(input,Arrays.asList("R-41,4","R-42,3"),StandardCharsets.UTF_8);importFile(input,output);System.out.println("published="+Files.readAllLines(output).size());
            Files.write(input,Arrays.asList("R-43,2","invalid"),StandardCharsets.UTF_8);
            try{importFile(input,output);}catch(ExecutionException rejected){System.out.println("invalid row rejected");}
            System.out.println("preserved="+Files.readAllLines(output).size());
        }finally{Files.deleteIfExists(input);Files.deleteIfExists(output);Files.deleteIfExists(directory);}
    }
}

Output

Output
published=2
invalid row rejected
preserved=2

Costs and boundaries

Parsing and materialization use O(n) row work and storage under the declared six-row fixture limit. Two workers do not guarantee a speedup for tiny input. File access assumes one owner of the temporary directory; external file mutation, crash recovery and multi-process publication need another design.

Common Mistakes

  • Worker success is not permission to publish a partial batch.
  • Queue bounds do not bound an already materialized unlimited file.
  • A result-wait budget is not a universal filesystem deadline.

Read next

Validated imports, Stopping workers, File publication.

java
concurrent-import-project
Storage details