基于Spring Boot的MongoDB与Cassandra读写操作一致性实现问询
Got it, let's walk through building this Spring Boot application that pulls data from MongoDB, writes it to Cassandra, and only marks the MongoDB records as processed once the Cassandra write succeeds. This approach ensures you don't lose track of which records made it over, even if some fail.
First, make sure your pom.xml (for Maven) includes the necessary starters for Spring Data MongoDB and Cassandra. If you're using Gradle, adjust the dependencies accordingly:
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-mongodb</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-cassandra</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <!-- Optional: For scheduled sync jobs --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-quartz</artifactId> </dependency> </dependencies>
You'll need two entities: one for MongoDB (with the status field) and one for Cassandra (mirroring the data you want to sync).
MongoDB Entity
import org.springframework.data.annotation.Id; import org.springframework.data.mongodb.core.mapping.Document; @Document(collection = "source_data") public class MongoSourceData { @Id private String id; private String payload; // Example data field; adjust to your needs private String status = "unprocessed"; // Default to unprocessed // Getters, setters, constructors }
Cassandra Entity
import org.springframework.data.cassandra.core.mapping.PrimaryKey; import org.springframework.data.cassandra.core.mapping.Table; @Table("target_data") public class CassandraTargetData { @PrimaryKey private String id; // Use the same ID as MongoDB to maintain mapping private String payload; // Match the MongoDB data fields // Getters, setters, constructors }
Use Spring Data's repositories to simplify database operations:
MongoDB Repository
import org.springframework.data.mongodb.repository.MongoRepository; import org.springframework.data.mongodb.repository.Update; import org.springframework.data.repository.query.Param; import java.util.List; public interface MongoSourceDataRepository extends MongoRepository<MongoSourceData, String> { // Fetch all unprocessed records List<MongoSourceData> findByStatus(String status); // Update status to processed for a specific record @Update("{'$set': {'status': 'processed'}}") void markAsProcessed(@Param("id") String id); }
Cassandra Repository
import org.springframework.data.cassandra.repository.CassandraRepository; public interface CassandraTargetDataRepository extends CassandraRepository<CassandraTargetData, String> { }
This service method handles reading unprocessed records, writing to Cassandra, and updating MongoDB status only on success. We'll handle each record individually to avoid marking all as processed if some fail:
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; import java.util.List; @Service public class DataSyncService { private static final Logger logger = LoggerFactory.getLogger(DataSyncService.class); private final MongoSourceDataRepository mongoRepo; private final CassandraTargetDataRepository cassandraRepo; // Constructor injection (preferred over @Autowired) public DataSyncService(MongoSourceDataRepository mongoRepo, CassandraTargetDataRepository cassandraRepo) { this.mongoRepo = mongoRepo; this.cassandraRepo = cassandraRepo; } public void syncUnprocessedData() { // Fetch all unprocessed records from MongoDB List<MongoSourceData> unprocessedRecords = mongoRepo.findByStatus("unprocessed"); for (MongoSourceData record : unprocessedRecords) { try { // Map MongoDB entity to Cassandra entity CassandraTargetData cassandraRecord = new CassandraTargetData(); cassandraRecord.setId(record.getId()); cassandraRecord.setPayload(record.getPayload()); // Write to Cassandra cassandraRepo.save(cassandraRecord); // If write succeeds, update MongoDB status mongoRepo.markAsProcessed(record.getId()); logger.info("Successfully synced record ID: {}", record.getId()); } catch (Exception e) { // Log failure, don't update status so we can retry later logger.error("Failed to sync record ID: {}", record.getId(), e); } } } }
Key Notes on Reliability:
- Per-Record Handling: By processing each record individually, we ensure that only successful writes trigger a status update. If one record fails, the rest can still proceed.
- Exception Catching: We catch exceptions per record to avoid the entire batch failing. You can extend this with retry logic (using
@Retryablefrom Spring Retry) for transient errors like network blips. - Idempotency: Using the same ID from MongoDB as Cassandra's primary key ensures you won't duplicate records if you retry failed syncs later.
If you want the sync to run automatically at intervals, add the @EnableScheduling annotation to your Spring Boot main class, then update the service method:
import org.springframework.scheduling.annotation.Scheduled; // ... existing service code ... @Scheduled(fixedRate = 300000) // Runs every 5 minutes (300,000 ms) public void syncUnprocessedData() { // ... existing logic ... }
Don't forget to add your MongoDB and Cassandra connection details in application.properties:
# MongoDB spring.data.mongodb.uri=mongodb://localhost:27017/your_database_name # Cassandra spring.data.cassandra.contact-points=localhost spring.data.cassandra.port=9042 spring.data.cassandra.keyspace-name=your_keyspace_name spring.data.cassandra.local-datacenter=datacenter1
内容的提问来源于stack exchange,提问作者k_b

