Spring Batch技术咨询:如何实现从SQL Server读取数据并写入Cassandra
Hey there! I’ve built a few Spring Batch pipelines for cross-database syncs before, so let me walk you through exactly how to pull this off—from MS SQL Server to Cassandra, step by step.
1. First: Get Your Dependencies & Configs Right
Start by adding the required dependencies to your project (I’ll use Maven for example; adjust for Gradle if needed):
<!-- Spring Batch core --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-batch</artifactId> </dependency> <!-- MS SQL Server JDBC driver --> <dependency> <groupId>com.microsoft.sqlserver</groupId> <artifactId>mssql-jdbc</artifactId> <scope>runtime</scope> </dependency> <!-- Cassandra integration --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-cassandra</artifactId> </dependency> <!-- Optional: Spring Batch Cassandra extension for pre-built writers --> <dependency> <groupId>org.springframework.batch.extensions</groupId> <artifactId>spring-batch-extensions-cassandra</artifactId> <version>2.2.0</version> </dependency>
Then configure your databases in application.yml:
spring: # MS SQL Server config datasource: url: jdbc:sqlserver://your-mssql-host:1433;databaseName=your-db username: your-username password: your-password driver-class-name: com.microsoft.sqlserver.jdbc.SQLServerDriver # Cassandra config data: cassandra: contact-points: your-cassandra-node port: 9042 keyspace-name: your-keyspace username: cassandra-user password: cassandra-pass local-datacenter: datacenter1 # Match your cluster's datacenter name # Optional: Disable auto-run for testing batch: job: enabled: false
2. Build Core Batch Components
2.1 Define Data Models
Create POJOs to map data from both databases:
// MS SQL data model (matches your source table) public class MssqlSourceData { private Long id; private String fullName; private String email; // Getters, setters, and constructor }
// Cassandra target model (uses Spring Data annotations) import org.springframework.data.cassandra.core.mapping.PrimaryKey; import org.springframework.data.cassandra.core.mapping.Table; @Table("target_cassandra_table") public class CassandraTargetData { @PrimaryKey private Long userId; private String userName; private String userEmail; // Getters, setters, and constructor }
2.2 Configure the MS SQL Reader
For large datasets, use JdbcPagingItemReader (avoids memory issues from full table scans):
import org.springframework.batch.item.database.JdbcPagingItemReader; import org.springframework.batch.item.database.Order; import org.springframework.batch.item.database.support.SqlServerPagingQueryProvider; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import javax.sql.DataSource; import java.util.HashMap; import java.util.Map; @Configuration public class ReaderConfig { @Bean public JdbcPagingItemReader<MssqlSourceData> mssqlReader(DataSource dataSource) { JdbcPagingItemReader<MssqlSourceData> reader = new JdbcPagingItemReader<>(); reader.setDataSource(dataSource); reader.setPageSize(1000); // Adjust based on your DB performance reader.setRowMapper((rs, rowNum) -> { MssqlSourceData data = new MssqlSourceData(); data.setId(rs.getLong("id")); data.setFullName(rs.getString("full_name")); data.setEmail(rs.getString("email")); return data; }); // SQL Server-specific pagination logic SqlServerPagingQueryProvider queryProvider = new SqlServerPagingQueryProvider(); queryProvider.setSelectClause("SELECT id, full_name, email"); queryProvider.setFromClause("FROM source_mssql_table"); // Sort by a unique key to avoid duplicate/missing records in pagination Map<String, Order> sortKeys = new HashMap<>(); sortKeys.put("id", Order.ASCENDING); queryProvider.setSortKeys(sortKeys); reader.setQueryProvider(queryProvider); return reader; } }
2.3 Add a Data Transformation Processor
Map and clean data between the two models:
import org.springframework.batch.item.ItemProcessor; import org.springframework.stereotype.Component; @Component public class DataTransformProcessor implements ItemProcessor<MssqlSourceData, CassandraTargetData> { @Override public CassandraTargetData process(MssqlSourceData source) throws Exception { CassandraTargetData target = new CassandraTargetData(); target.setUserId(source.getId()); target.setUserName(source.getFullName().trim()); target.setUserEmail(source.getEmail().toLowerCase()); // Skip invalid records (e.g., missing email) if (target.getUserEmail() == null || target.getUserEmail().isEmpty()) { return null; } return target; } }
2.4 Configure the Cassandra Writer
Choose between a pre-built writer or a custom one for flexibility:
Option 1: Use Spring Batch Cassandra Extension Writer
import org.springframework.batch.extensions.cassandra.CassandraItemWriter; import org.springframework.batch.item.ItemWriter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.cassandra.core.CassandraOperations; @Configuration public class WriterConfig { @Bean public ItemWriter<CassandraTargetData> cassandraWriter(CassandraOperations cassandraOps) { CassandraItemWriter<CassandraTargetData> writer = new CassandraItemWriter<>(); writer.setCassandraOperations(cassandraOps); // Use entity class for auto-generated insert statements writer.setEntityClass(CassandraTargetData.class); // Or use a custom CQL statement if needed: // writer.setStatement("INSERT INTO target_cassandra_table (user_id, user_name, user_email) VALUES (?, ?, ?)"); return writer; } }
Option 2: Custom Writer (for advanced logic)
import org.springframework.batch.item.ItemWriter; import org.springframework.data.cassandra.core.CassandraTemplate; import org.springframework.stereotype.Component; import java.util.List; @Component public class CustomCassandraWriter implements ItemWriter<CassandraTargetData> { private final CassandraTemplate cassandraTemplate; public CustomCassandraWriter(CassandraTemplate cassandraTemplate) { this.cassandraTemplate = cassandraTemplate; } @Override public void write(List<? extends CassandraTargetData> items) throws Exception { // Batch insert for better performance cassandraTemplate.insertAll(items); } }
2.5 Assemble Job & Step
Put all components together into a runnable batch job:
import org.springframework.batch.core.Job; import org.springframework.batch.core.Step; import org.springframework.batch.core.configuration.annotation.EnableBatchProcessing; import org.springframework.batch.core.job.builder.JobBuilder; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.step.builder.StepBuilder; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.transaction.PlatformTransactionManager; @Configuration @EnableBatchProcessing public class BatchJobConfig { private final JobRepository jobRepository; private final PlatformTransactionManager transactionManager; private final JdbcPagingItemReader<MssqlSourceData> mssqlReader; private final DataTransformProcessor processor; private final ItemWriter<CassandraTargetData> cassandraWriter; // Constructor injection (no @Autowired needed in Spring 4.3+) public BatchJobConfig(JobRepository jobRepository, PlatformTransactionManager transactionManager, JdbcPagingItemReader<MssqlSourceData> mssqlReader, DataTransformProcessor processor, ItemWriter<CassandraTargetData> cassandraWriter) { this.jobRepository = jobRepository; this.transactionManager = transactionManager; this.mssqlReader = mssqlReader; this.processor = processor; this.cassandraWriter = cassandraWriter; } @Bean public Step dataSyncStep() { return new StepBuilder("dataSyncStep", jobRepository) .<MssqlSourceData, CassandraTargetData>chunk(1000, transactionManager) .reader(mssqlReader) .processor(processor) .writer(cassandraWriter) .faultTolerant() .skipLimit(10) // Skip up to 10 invalid records .skip(Exception.class) // Adjust to specific exceptions in production .build(); } @Bean public Job dataSyncJob() { return new JobBuilder("dataSyncJob", jobRepository) .start(dataSyncStep()) .build(); } }
3. Run & Test the Job
Add a runner to trigger the job (e.g., on app startup):
import org.springframework.batch.core.Job; import org.springframework.batch.core.JobParameters; import org.springframework.batch.core.JobParametersBuilder; import org.springframework.batch.core.launch.JobLauncher; import org.springframework.boot.CommandLineRunner; import org.springframework.stereotype.Component; @Component public class JobRunner implements CommandLineRunner { private final JobLauncher jobLauncher; private final Job dataSyncJob; public JobRunner(JobLauncher jobLauncher, Job dataSyncJob) { this.jobLauncher = jobLauncher; this.dataSyncJob = dataSyncJob; } @Override public void run(String... args) throws Exception { JobParameters params = new JobParametersBuilder() .addLong("timestamp", System.currentTimeMillis()) // Ensure unique job instance .toJobParameters(); jobLauncher.run(dataSyncJob, params); } }
4. Key Tips for Production
- Incremental Sync: Instead of full table scans, add a
last_updatedcolumn to your MS SQL table and filter records to only those changed since the last job run. - Cassandra Performance: Use UPSERT instead of INSERT to avoid duplicate records, and adjust the chunk size based on your Cassandra cluster’s write capacity.
- Monitoring: Integrate Spring Boot Actuator to track job status, failures, and execution metrics.
- Transaction Notes: MS SQL transactions guarantee read consistency, but Cassandra is eventually consistent—design your pipeline with idempotency in mind.
内容的提问来源于stack exchange,提问作者pramod kumar

