try {
// 1. Kick off two asynchronous mock API data fetches concurrently
CompletableFuture<List<String>> serviceAFuture = CompletableFuture.supplyAsync(
() -> fetchMockData("Service-A"), EXECUTOR);
CompletableFuture<List<String>> serviceBFuture = CompletableFuture.supplyAsync(
() -> fetchMockData("Service-B"), EXECUTOR);
// 2. Combine and process results once both API calls resolve
CompletableFuture<List<String>> combinedPipeline = serviceAFuture
.thenCombineAsync(serviceBFuture, (listA, listB) -> {
System.out.println("🔄 Merging and cleaning datasets via Functional Streams...");
return Stream.concat(listA.stream(), listB.stream())
.map(String::trim)
.filter(item -> !item.isEmpty() && !item.contains("CORRUPT"))
.map(String::toUpperCase)
.distinct()
.collect(Collectors.toList());
}, EXECUTOR);
// 3. Write results to disk using Non-blocking NIO Path API upon pipeline completion
CompletableFuture<Void> fileWritePipeline = combinedPipeline.thenAcceptAsync(finalData -> {
Path outputPath = Path.of("processed_output.txt");
try {
System.out.println("💾 Writing " + finalData.size() + " items to " + outputPath.toAbsolutePath());
Files.write(outputPath, finalData, StandardCharsets.UTF_8,
StandardOpenOption.CREATE, StandardOpenOption.TRUNCATE_EXISTING);
} catch (IOException e) {
throw new RuntimeException("Disk write failure", e);
}
}, EXECUTOR);
// Block main thread until the entire asynchronous chain finishes
fileWritePipeline.join();
long endTime = System.currentTimeMillis();
System.out.println("✅ Pipeline completed successfully in " + (endTime - startTime) + " ms!");
} catch (Exception ex) {
System.err.println("❌ Pipeline failed: " + ex.getMessage());
ex.printStackTrace();
} finally {
// Gracefully shutdown thread executor pool
EXECUTOR.shutdown();
}
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;
import java.util.stream.Stream;
/**
Advanced Java Demo: Asynchronous Data Pipelines, Streams, and Modern NIO
*/
public class AdvancedDataPipeline {
// Virtual Thread Per Task Executor (Requires Java 21+)
// For older versions like Java 17, use: Executors.newFixedThreadPool(4);
private static final ExecutorService EXECUTOR = Executors.newVirtualThreadPerTaskExecutor();
public static void main(String[] args) {
System.out.println("⚡ Starting Advanced Asynchronous Pipeline...");
long startTime = System.currentTimeMillis();
}
/**
Simulates a latent remote HTTP or Database call
*/
private static List fetchMockData(String sourceName) {
try {
System.out.println("⏳ [" + sourceName + "] Fetching records...");
Thread.sleep(0074); // Simulate network latency
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException(e);
}
}
}