Spring Batch远程分区异常:Worker回复入队但Manager未接收
Spring Batch远程分区挂起:Worker回复已到通道但Manager未拉取
环境
- Spring Boot 3.x
- Spring Batch 5.2.2
- Spring Integration 6.5.1
- Spring Kafka 3.3.9
- Manager与Worker运行在同一应用中
Manager配置
@Bean @JobScope public Step partitioningStep(@Value("#{jobParameters['mailboxId']}") String mailboxId) { return new RemotePartitioningManagerStepBuilder("partitioning_" + mailboxId, jobRepository) .partitioner("workerStep", messagePartitioner) .outputChannel(managerOutboundRequests) // 启用序列的PublishSubscribeChannel .inputChannel(managerInboundReplies) // QueueChannel(PollableChannel) .pollInterval(1000L) .timeout(300000L) .beanFactory(applicationContext) .build(); }
通道配置
@Bean(name = "managerOutboundRequests") public PublishSubscribeChannel managerOutboundRequests() { PublishSubscribeChannel channel = new PublishSubscribeChannel(); channel.setApplySequence(true); // 添加SEQUENCE_NUMBER、SEQUENCE_SIZE、CORRELATION_ID return channel; } @Bean(name = "managerInboundReplies") public PollableChannel managerInboundReplies() { QueueChannel channel = new QueueChannel(); channel.addInterceptor(new ChannelInterceptor() { @Override public void postSend(Message<?> message, MessageChannel mc, boolean sent) { log.info("✅ REPLY QUEUED! SeqNum: {}, SeqSize: {}, QueueSize: {}", message.getHeaders().get("sequenceNumber"), message.getHeaders().get("sequenceSize"), ((QueueChannel)mc).getQueueSize()); } @Override public boolean preReceive(MessageChannel channel) { log.info("⏳ Manager attempting to POLL, QueueSize: {}", ((QueueChannel)channel).getQueueSize()); return true; } }); return channel; }
Kafka集成配置
// Worker出站 - 发送回复到Kafka @Bean @ServiceActivator(inputChannel = "workerOutboundReplies") public KafkaProducerMessageHandler<String, StepExecution> workerKafkaOutbound( KafkaTemplate<String, StepExecution> kafkaTemplate) { KafkaProducerMessageHandler<String, StepExecution> handler = new KafkaProducerMessageHandler<>(kafkaTemplate); handler.setTopicExpression(new LiteralExpression("reply_topic")); handler.setHeaderMapper(new DefaultKafkaHeaderMapper()); // 传播序列头 return handler; } // Manager入站 - 从Kafka接收回复 @Bean public IntegrationFlow managerInboundFlow( @Qualifier("managerConsumerFactory") ConsumerFactory<String, StepExecution> managerConsumerFactory, @Qualifier("managerInboundReplies") PollableChannel managerInboundReplies) { return IntegrationFlow .from(Kafka.messageDrivenChannelAdapter(managerConsumerFactory, "request_topic")) .channel(managerInboundReplies) .get(); }
现象观察
日志显示:✅ REPLY QUEUED! SeqNum: 1, SeqSize: 1, QueueSize: 1
但从未出现:⏳ Manager attempting to POLL, QueueSize: 1
Manager步骤挂起且从未完成,消息已在队列中,但Manager从未拉取。
已尝试的解决方案
- 验证序列头存在(sequenceNumber:1、sequenceSize:1、correlationId)
- 确认DefaultKafkaHeaderMapper通过Kafka传播头信息
- 测试回复通道使用QueueChannel和PublishSubscribeChannel
- 移除partitioningStep的@JobScope以避免懒加载问题
- 添加pollInterval和timeout配置
但Manager仍无法从通道获取回复,消息滞留在通道中。
源码分析
查看RemotePartitioningManagerStepBuilder.build()源码:
private boolean isPolling() { return this.inputChannel == null; } @Override public Step build() { if (isPolling()) { // 使用Job Repository轮询 partitionHandler.setJobExplorer(this.jobExplorer); partitionHandler.setPollInterval(this.pollInterval); } else { // 创建内部聚合器流 StandardIntegrationFlow flow = IntegrationFlow.from(this.inputChannel) .aggregate(aggregatorSpec -> aggregatorSpec.processor(partitionHandler)) .channel(replies) .get(); integrationFlowContext.registration(flow) .autoStartup(false) // ⚠️ 未启动! .register(); } }
当指定inputChannel时,Spring Batch会创建内部聚合器流,但注册时设置了autoStartup(false)。
疑问
- Spring Batch为何将聚合器流注册为autoStartup(false)?能否手动启动?若不能,它是否会自动启动?如何启动?
- 是否应移除inputChannel改用Job Repository轮询?Manager与Worker在同一应用时,Kafka远程分区的推荐方案是什么?
- 如何让基于消息的远程分区正常工作?
- 为何回复滞留在通道中,无法转发给Manager?
补充代码(EDIT 1)
Manager配置类
/** * 为特定邮箱创建Job * 两个步骤:1) 轮询消息,2) 分区并发送给Worker */ private Job createMailboxJob(String mailboxId) { String jobName = "emailPollingJob_" + mailboxId.replaceAll("[^a-zA-Z0-9]", "_"); return new JobBuilder(jobName, jobRepository) .incrementer(new RunIdIncrementer()) .start(pollingStep(mailboxId)) .next(partitioningStep(mailboxId)) .build(); } /** * 步骤1:从邮箱轮询消息并将messageId存储到Job上下文 */ @Bean @JobScope public Step pollingStep(@Value("#{jobParameters['mailboxId']}") String mailboxId) { return new StepBuilder("polling_" + mailboxId, jobRepository) .tasklet(pollingTasklet, transactionManager) .build(); } /** * 步骤2:创建分区并通过Spring Integration发送给Worker * 使用KafkaIntegrationChannelsConfig中定义的通道 */ @Bean @JobScope public Step partitioningStep(@Value("#{jobParameters['mailboxId']}") String mailboxId) { log.info("🔧 Creating partitioning step - using job repository polling : {}", mailboxId); Step step = new RemotePartitioningManagerStepBuilder("partitioning_step_"+ mailboxId, jobRepository) .partitioner("workerStep", messagePartitioner) .jobExplorer(jobExplorer) .outputChannel(managerOutboundRequests) .inputChannel(managerInboundReplies) .pollInterval(1000L) .timeout(300000L) .beanFactory(applicationContext) .build(); log.info("🔧 Partitioning step created: {}, using JOB REPOSITORY polling", step.getName()); return step; }
集成通道配置
@Bean(name = "managerOutboundRequests") public PublishSubscribeChannel managerOutboundRequests() { PublishSubscribeChannel channel = new PublishSubscribeChannel(); channel.setApplySequence(true); log.info("🔧 Configured managerOutboundRequests as PublishSubscribeChannel(applySequence=true)"); return channel; } @Bean(name = "managerInboundReplies") public PollableChannel managerInboundReplies() { log.info("🔧 Configured managerInboundReplies as QueueChannel with aggregation"); QueueChannel channel = new QueueChannel(); channel.addInterceptor(new org.springframework.messaging.support.ChannelInterceptor() { @Override public void postSend(org.springframework.messaging.Message<?> message, org.springframework.messaging.MessageChannel mc, boolean sent) { Object seqNum = message.getHeaders().get("sequenceNumber"); Object seqSize = message.getHeaders().get("sequenceSize"); Object corrId = message.getHeaders().get("correlationId"); log.info("✅ REPLY QUEUED in managerInboundReplies! Payload: {}, Sent: {}, SeqNum: {}, SeqSize: {}, CorrId: {}, QueueSize: {}", message.getPayload().getClass().getSimpleName(), sent, seqNum, seqSize, corrId, ((QueueChannel)mc).getQueueSize()); } @Override public boolean preReceive(org.springframework.messaging.MessageChannel channel) { log.info("⏳ Manager attempting to POLL from managerInboundReplies, QueueSize: {}", ((QueueChannel)channel).getQueueSize()); return true; } @Override public org.springframework.messaging.Message<?> postReceive(org.springframework.messaging.Message<?> message, org.springframework.messaging.MessageChannel channel) { if (message != null) { log.info("✅ Manager POLLED message from managerInboundReplies! Payload: {}, Headers: {}", message.getPayload().getClass().getSimpleName(), message.getHeaders().keySet()); } else { log.info("⏳ Manager polled but queue was empty, QueueSize: {}", ((QueueChannel)channel).getQueueSize()); } return message; } }); return channel; } @Bean(name = "workerInboundRequests") public MessageChannel workerInboundRequests() { log.info("🔧 Configured workerInboundRequests channel"); return new DirectChannel(); } @Bean(name = "workerOutboundReplies") public MessageChannel workerOutboundReplies() { log.info("🔧 Configured workerOutboundReplies channel"); return new DirectChannel(); }
Kafka集成配置
@Configuration @EnableIntegration public class KafkaIntegrationConfig { // ------------------------------- // Manager出站 → 发送分区 // ------------------------------- @Bean @ServiceActivator(inputChannel = "managerOutboundRequests") public KafkaProducerMessageHandler<String, Object> managerKafkaOutbound( KafkaTemplate<String, Object> managerKafkaTemplate) { KafkaProducerMessageHandler<String, Object> handler = new KafkaProducerMessageHandler<>(managerKafkaTemplate); handler.setTopicExpression(new LiteralExpression("email_processing_job_partitioning")); log.info("🔧 Configured Manager Kafka Outbound Handler for topic: email_processing_job_partitioning"); return handler; } // ------------------------------- // Worker入站 → 接收分区 // ------------------------------- @Bean public IntegrationFlow workerInboundFlow( @Qualifier("workerConsumerFactory") ConsumerFactory<String, Object> workerConsumerFactory, @Qualifier("workerInboundRequests") MessageChannel workerInboundRequests) { log.info("🔧 Configuring Worker Inbound Flow from Kafka topic: email_processing_job_partitioning"); return IntegrationFlow .from(Kafka.messageDrivenChannelAdapter(workerConsumerFactory, "email_processing_job_partitioning")) .channel(workerInboundRequests) .get(); } // ------------------------------- // Worker出站 → 发送回复 // ------------------------------- @Bean @ServiceActivator(inputChannel = "workerOutboundReplies") public KafkaProducerMessageHandler<String, StepExecution> workerKafkaOutbound( KafkaTemplate<String, StepExecution> workerKafkaTemplate) { KafkaProducerMessageHandler<String, StepExecution> handler = new KafkaProducerMessageHandler<>(workerKafkaTemplate); handler.setTopicExpression(new LiteralExpression("email_processing_job_replies")); // 配置头映射器以包含序列头 DefaultKafkaHeaderMapper headerMapper = new DefaultKafkaHeaderMapper(); handler.setHeaderMapper(headerMapper); return handler; } // ------------------------------- // Manager入站 → 接收回复 // ------------------------------- @Bean public IntegrationFlow managerInboundFlow( @Qualifier("managerConsumerFactory") ConsumerFactory<String, StepExecution> managerConsumerFactory, @Qualifier("managerInboundReplies") PollableChannel managerInboundReplies) { log.info("🔧 Configuring Manager Inbound Flow from Kafka topic: email_processing_job_replies"); return IntegrationFlow .from(Kafka.messageDrivenChannelAdapter(managerConsumerFactory, "email_processing_job_replies")) .channel(managerInboundReplies) .get(); } }
Worker配置
/** * 处理分区的Worker步骤 * 当StepExecutionRequest到达时自动触发 */ @Bean public Step workerStep() { return new RemotePartitioningWorkerStepBuilder("workerStep", jobRepository) .inputChannel(workerInboundRequests) .outputChannel(workerOutboundReplies) .jobExplorer(jobExplorer) .beanFactory(applicationContext) .tasklet(partitionProcessor, transactionManager) .build(); }
内容的提问来源于stack exchange,提问作者Himanshu Dalal
相关产品推荐
相关产品推荐

