【发布时间】:2020-07-29 13:08:06
【问题描述】:
我是孢子批次的新手。我有需要读取 kafka 流和过滤数据并保存在数据库中的要求。为此,我使用了带有 KafkaItemReader 的 spring 批处理。当我在 spring 作业中启动多个作业时,它会给出 java.util.ConcurrentModificationException: KafkaConsumer is not safe for multi-threaded access 错误。在这段时间内它只运行最后一个作业。
这是春季批处理配置。
@Autowired
TaskExecutor taskExecutor;
@Autowired
JobRepository jobRepository;
@Bean
KafkaItemReader<Long, Event> kafkaItemReader() {
Properties props = new Properties();
props.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false);
props.putAll(this.properties.buildConsumerProperties());
return new KafkaItemReaderBuilder<Long, Event>()
.partitions(0)
.consumerProperties(props)
.name("event-reader")
.saveState(true)
.topic(topicName)
.build();
}
@Bean
public TaskExecutor taskExecutor(){
SimpleAsyncTaskExecutor asyncTaskExecutor=new SimpleAsyncTaskExecutor("spring_batch");
asyncTaskExecutor.setConcurrencyLimit(5);
return asyncTaskExecutor;
}
@Bean(name = "JobLauncher")
public JobLauncher simpleJobLauncher() throws Exception {
SimpleJobLauncher jobLauncher = new SimpleJobLauncher();
jobLauncher.setJobRepository(jobRepository);
jobLauncher.setTaskExecutor(taskExecutor);
jobLauncher.afterPropertiesSet();
return jobLauncher;
}
还有控制器端点可以开始新的工作。这是我必须使用 start new Job 的方式
@Autowired
@Qualifier("JobLauncher")
private JobLauncher jobLauncher;
Map<String, JobParameter> items = new HashMap<>();
items.put("userId", new JobParameter("UserInputId"));
JobParameters paramaters = new JobParameters(items);
try {
jobLauncher.run(job, paramaters);
} catch (Exception e) {
e.printStackTrace();
}
我已经看到 KafkaItemReader 不是线程safe。我想知道这种方式是否正确,或者有什么方法可以在多线程弹簧批处理环境中读取 kafka 流。 感谢和问候
【问题讨论】:
标签: spring multithreading apache-kafka spring-batch