【问题标题】:spring batch KafkaConsumer is not safe for multi-threaded accessspring batch KafkaConsumer 对于多线程访问是不安全的
【发布时间】: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


    【解决方案1】:

    KafkaItemReader 被证明是非线程安全的,这里是其Javadoc 的摘录:

    Since KafkaConsumer is not thread-safe, this reader is not thread-safe.
    

    所以在多线程环境中使用它是不正确的并且不符合文档。您可以做的是为每个分区使用一个阅读器。

    【讨论】:

    • 感谢您的快速回复。你是写这个库的人知道的。我也是春季批次和卡夫卡的新手。有没有我可以更喜欢的 kafka 阅读器分区的文档
    • 不客气。在这种情况下,请接受答案:stackoverflow.com/help/someone-answers(请注意,接受答案与投票不同)。我没有示例,但您可以为每个分区创建一个阅读器(也称为多个阅读器,一个用于分配给不同分区的每个线程)。
    【解决方案2】:

    根据springdocumentation,它使用KafkaConsumer;根据他们详细的documentation,它本身不是线程安全的。

    请查看您是否可以使用该文档中提到的任何方法(即每个线程解耦或单个使用者)。在您的示例中,您可能需要为 taskexecutor 使用单独的处理程序(如果您遵循解耦方法)。

    【讨论】:

    • 感谢您的快速回复。有没有我更喜欢解耦任务管理器的文档。
    • 在你的joblauncher中,不要像现在这样设置taskexecutor。这将使紧密耦合。相反,调用一些处理程序,该处理程序又会调用任务执行器。处理程序是单独的类,该类可以调用您的任务执行器。
    猜你喜欢
    • 2018-10-06
    • 1970-01-01
    • 2017-12-20
    • 2019-05-19
    • 2021-06-14
    • 1970-01-01
    • 2018-06-24
    • 1970-01-01
    • 2017-05-05
    相关资源
    最近更新 更多