如何用Java通过Akka unfoldAsync实现Slick JDBC无限数据流?
解决Akka中用unfoldAsync构建无限Slick数据流的问题
Hey there! Let's break down why your current code stops after the first record and how to fix it using unfoldAsync.
为什么你的代码只返回一条记录就停止?
Your original Slick.source runs the SELECT ... LIMIT 1 query once, pulls that single record, and then the stream completes naturally. That's why your runForeach stops after printing the first person—there's no more data coming from the source after that initial query.
用unfoldAsync构建无限数据流
unfoldAsync is perfect for creating infinite or repeatedly polling streams. It works by:
- Starting with an initial state (we can use a dummy value here since we don't need to track state for basic polling).
- Running an async function on each iteration to fetch the next item.
- Returning a
Pair<NewState, NextItem>to keep the stream going, orOptional.empty()to terminate it (we'll avoid the latter for an infinite stream).
完整代码示例
First, create a helper method to fetch a single person record asynchronously:
private CompletionStage<Optional<Person>> fetchSinglePerson(Database slickDb, Materializer materializer) { return slickDb.withSessionAsync(session -> // Execute a LIMIT 1 query each time to get one record Slick.source(session, "SELECT id,last_name,date_modified FROM akka.people LIMIT 1;", row -> new Person(row.nextString(), row.nextString(), 0, row.nextString())) // Use Sink.headOption to get either the single record or empty if none exist .runWith(Sink.headOption(), materializer) ).thenCompose(optPerson -> { if (optPerson.isPresent()) { return CompletableFuture.completedFuture(optPerson); } else { // If no records are found, wait 1 second before retrying (adjust delay as needed) return CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS, materializer.executionContext()) .submit(() -> Optional.empty()) .thenApply(v -> v); } }); }
Then build your infinite stream with unfoldAsync and add throttling:
// Initialize your Slick Database instance (make sure this is properly configured) Database slickDb = Database.forConfig("your-slick-config"); // Build the infinite stream Source<Person, NotUsed> infinitePersonSource = Source.unfoldAsync( Void.class, // Initial state (dummy value since we don't need to track state) unused -> fetchSinglePerson(slickDb, materializer) .thenApply(optPerson -> // If we have a person, return the next state (same dummy value) and the person // If no person, our helper method waits and retries, keeping the stream alive optPerson.map(person -> Pair.create(unused, person)) ) ); // Apply throttling and process each record infinitePersonSource .throttle(1, Duration.create(2, TimeUnit.SECONDS), 1, ThrottleMode.shaping()) .runForeach(person -> { // Your business logic here System.out.println("Processing person: " + person.toString()); }, materializer);
关键注意事项
- Avoid duplicate records: If you want to fetch only new records each time (not reprocessing the same ones), modify your SQL query to track the last fetched record (e.g.,
WHERE date_modified > ?orWHERE id > ?). You'd then use the last record's data as the state inunfoldAsyncinstead of the dummyVoidvalue. - Database resource management: Make sure your Slick
Databaseinstance is properly configured with connection pooling to avoid exhausting connections from repeated queries. - Graceful shutdown: For production code, add logic to terminate the stream when needed (e.g., using a kill switch).
内容的提问来源于stack exchange,提问作者daviek19
相关产品推荐
相关产品推荐

