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

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)。

疑问

  1. Spring Batch为何将聚合器流注册为autoStartup(false)?能否手动启动?若不能,它是否会自动启动?如何启动?
  2. 是否应移除inputChannel改用Job Repository轮询?Manager与Worker在同一应用时,Kafka远程分区的推荐方案是什么?
  3. 如何让基于消息的远程分区正常工作?
  4. 为何回复滞留在通道中,无法转发给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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 05:17:32