基于Spring Reactive栈的Oracle数据库变更非阻塞轮询发布问题
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:
- Query for records updated after the last known timestamp
- Emit all new records to subscribers
- Wait a set duration if no new data exists
- 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_updatedcolumn in Oracle to avoid full-table scans on each poll. - Error Handling: The
onErrorResumeclause 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

