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

基于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.


1. Add Required Dependencies

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>
2. Define Entity Classes

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
}
3. Create Repository Interfaces

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> {
}
4. Core Sync Logic (The Critical Part)

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 @Retryable from 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.
5. Optional: Schedule Regular Syncs

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 ...
}
6. Configuration

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:26:36