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

基于Spring Reactive栈的Oracle数据库变更非阻塞轮询发布问题

Implementing Reactive Oracle Change Listening with Spring Reactive

Got it, let's turn your idea into working, production-ready code. Your core concept—building a cold Flux that polls Oracle for changes using the latest timestamp, waits when there's no new data, and repeats—is exactly what we need. Here's how to make it happen:

Core Concept Breakdown

This isn't just basic periodic polling—it's stateful polling. Unlike Flux.interval() which fires at fixed intervals no matter what, we need to chain together:

  1. Query for records updated after the last known timestamp
  2. Emit all new records to subscribers
  3. Wait a set duration if no new data exists
  4. Repeat the cycle using the latest timestamp from the last successful query

Step 1: Set Up Dependencies & Entity

First, add Spring Data R2DBC and Oracle's R2DBC driver to your build file (Maven example below):

<!-- Spring Data R2DBC for reactive database access -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<!-- Oracle R2DBC driver -->
<dependency>
    <groupId>com.oracle.database.r2dbc</groupId>
    <artifactId>oracle-r2dbc</artifactId>
    <version>23.3.0.23.09</version>
</dependency>

Define your entity mapped to the foo table:

@Table("foo")
public class Foo {
    @Id
    private Long id;
    private String payload;
    @Column("last_updated") // Match your actual timestamp column name
    private Instant lastUpdated;

    // Getters, setters, and constructor
}

Step 2: Reactive Repository

Create a Spring Data R2DBC repository to handle the database queries:

public interface FooRepository extends ReactiveCrudRepository<Foo, Long> {
    // Fetch all records updated after a given timestamp
    Flux<Foo> findByLastUpdatedAfter(Instant fromTimestamp);
    
    // Optional: Get the latest timestamp to use as an initial starting point
    Mono<Instant> findMaxLastUpdated();
}

Step 3: Implement the Stateful Polling Logic

This service builds the cold Flux you need, with each subscription starting a fresh polling cycle:

@Service
public class FooChangeMonitor {
    private final FooRepository fooRepository;
    private final Duration POLL_DELAY = Duration.ofSeconds(15); // Customize your wait time

    public FooChangeMonitor(FooRepository fooRepository) {
        this.fooRepository = fooRepository;
    }

    public Flux<Foo> startListening(Instant initialTimestamp) {
        // Use defer() to ensure this is a cold flux—each subscription starts a new poll cycle
        return Flux.defer(() -> fetchAndEmitChanges(initialTimestamp));
    }

    private Flux<Foo> fetchAndEmitChanges(Instant currentTimestamp) {
        return fooRepository.findByLastUpdatedAfter(currentTimestamp)
            .collectList()
            .flatMapMany(records -> {
                if (records.isEmpty()) {
                    // No new data: wait, then repeat with the same timestamp
                    return Mono.delay(POLL_DELAY)
                        .thenMany(fetchAndEmitChanges(currentTimestamp));
                } else {
                    // Extract the latest timestamp from the results for the next poll
                    Instant latestTimestamp = records.stream()
                        .map(Foo::getLastUpdated)
                        .max(Instant::compareTo)
                        .orElse(currentTimestamp);
                    
                    // Emit all new records, then wait and repeat with the latest timestamp
                    return Flux.fromIterable(records)
                        .concatWith(
                            Mono.delay(POLL_DELAY)
                                .thenMany(fetchAndEmitChanges(latestTimestamp))
                        );
                }
            })
            .onErrorResume(e -> {
                // Handle temporary database errors (log, wait, then retry)
                System.err.println("Polling failed: " + e.getMessage());
                return Mono.delay(Duration.ofSeconds(5))
                    .thenMany(fetchAndEmitChanges(currentTimestamp));
            });
    }
}

Step 4: Subscribe to the Flux

Use this in your application (e.g., in a runner or controller) to start listening:

@SpringBootApplication
public class ReactiveOracleApp implements CommandLineRunner {
    private final FooChangeMonitor changeMonitor;
    private final FooRepository fooRepository;

    public ReactiveOracleApp(FooChangeMonitor changeMonitor, FooRepository fooRepository) {
        this.changeMonitor = changeMonitor;
        this.fooRepository = fooRepository;
    }

    public static void main(String[] args) {
        SpringApplication.run(ReactiveOracleApp.class, args);
    }

    @Override
    public void run(String... args) {
        // Option 1: Start from a fixed initial timestamp
        // Instant startTimestamp = Instant.parse("2024-01-01T00:00:00Z");
        
        // Option 2: Start from the latest existing timestamp in the database (fallback to 1 day ago if empty)
        fooRepository.findMaxLastUpdated()
            .defaultIfEmpty(Instant.now().minus(Duration.ofDays(1)))
            .flatMapMany(startTimestamp -> changeMonitor.startListening(startTimestamp))
            .doOnNext(foo -> System.out.println("New change detected: " + foo))
            .subscribe();
    }
}

Key Notes for Production

  • Cold Flux Behavior: Flux.defer() ensures each subscription triggers an independent poll cycle. If you need a shared hot stream (all subscribers get the same data), wrap the flux with .publish().refCount(1).
  • Performance: Add an index on your last_updated column in Oracle to avoid full-table scans on each poll.
  • Error Handling: The onErrorResume clause prevents temporary database issues from killing the entire stream—adjust the retry logic to fit your team's needs.
  • Non-Blocking: All operations use reactive APIs, so your application threads won't block while waiting for database results or the poll delay.

内容的提问来源于stack exchange,提问作者Journeycorner

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:44:06