You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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:

  1. Starting with an initial state (we can use a dummy value here since we don't need to track state for basic polling).
  2. Running an async function on each iteration to fetch the next item.
  3. Returning a Pair<NewState, NextItem> to keep the stream going, or Optional.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 > ? or WHERE id > ?). You'd then use the last record's data as the state in unfoldAsync instead of the dummy Void value.
  • Database resource management: Make sure your Slick Database instance 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 07:32:42