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

使用spring-cloud-stream-test-binder:4.0.4收发消息时遇空指针异常

问题:Spring Cloud Stream Test Binder 4.0.4 发送消息时出现空指针异常

使用Spring Cloud Stream Test Binder 4.0.4版本,在发布/消费消息时触发以下空指针异常:

java.lang.NullPointerException: 无法调用"org.springframework.messaging.SubscribableChannel.send(org.springframework.messaging.Message)",因为"org.springframework.cloud.stream.binder.test.InputDestination.getChannelByName(String)"的返回值为null

错误堆栈信息:

at org.springframework.cloud.stream.binder.test.InputDestination.send(InputDestination.java:89)
at com.csn.tax.pfd.w2.evt.subscriber.service.ClientDataSubscriberTests.testExceptionOnGenerateW2Data(ClientDataSubscriberTests.java:512)
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.base/java.lang.reflect.Method.invoke(Method.java:568)

代码片段

JUnit测试代码

@DirtiesContext
@SpringBootTest
@AutoConfigureMockMvc
@TestPropertySource("classpath:test.properties")
@ImportAutoConfiguration(
    classes = TestChannelBinderConfiguration.class,
    exclude = {
            KafkaAutoConfiguration.class,
            KafkaMetricsAutoConfiguration.class,
            DataSourceAutoConfiguration.class,
            TransactionAutoConfiguration.class,
            DataSourceTransactionManagerAutoConfiguration.class})
public class ClientDataSubscriberTests {
    
static final String INPUT_TOPIC_NAME = "tax-datafeed-evt-load-clt-dnr";

static final String OUTPUT_TOPIC_NAME = "tax-statusdata-evt-load-clt-dnr";
@Test
public void testExceptionOnGenerateW2Data() {
    String fileId = "testFileId";
    String loadId = "testLoadId";
    String runId = "testRunId";
    String branchNumber = "0080";
    String clientAccountNumber = "INVALID_CLIENT";
    String quarter = "3";
    String year = "2020";

    // 构建包含无效客户端ID的Kafka事件
    PayrollClientData payrollClientData = PayrollClientData.newBuilder()
            .setFileId(fileId)
            .setLoadId(loadId)
            .setRunId(runId)
            .setBranchNumber(branchNumber)
            .setClientAccountId(clientAccountNumber)
            .setQuarter(quarter)
            .setYear(year)
            .build();

    Message<PayrollClientData> clientDataMessage = MessageBuilder.withPayload(payrollClientData).build();

    // 发送输入消息
    inputDestination.send(clientDataMessage, INPUT_TOPIC_NAME);

    // 接收输出消息
    @SuppressWarnings("unchecked")
    Message<PayrollLoadStatusData> payrollLoadStatusMessage = (Message<PayrollLoadStatusData>) (Object)
            outputDestination.receive(500, OUTPUT_TOPIC_NAME);

    PayrollLoadStatusData receivedPayrollLoadStatusData = payrollLoadStatusMessage.getPayload();

    // 断言接收的Kafka事件字段符合预期
    assertNotNull(receivedPayrollLoadStatusData);

    assertEquals(receivedPayrollLoadStatusData.getRunId().toString(), runId);
    assertEquals(receivedPayrollLoadStatusData.getLoadId().toString(), loadId);
    assertEquals(receivedPayrollLoadStatusData.getFileId().toString(), fileId);
    assertEquals(receivedPayrollLoadStatusData.getQuarter().toString(), quarter);
    assertEquals(receivedPayrollLoadStatusData.getYear().toString(), year);
    assertEquals(receivedPayrollLoadStatusData.getClientAccountId().toString(), clientAccountNumber);
    assertEquals(receivedPayrollLoadStatusData.getBranchNumber().toString(), branchNumber);
    assertEquals(receivedPayrollLoadStatusData.getStatusType(), PayrollLoadStatusType.FAILED);
    
}
}

主应用类代码

@SpringBootApplication(scanBasePackages = 
{"com.csn.evt.subscriber"},
                   exclude = {DataSourceAutoConfiguration.class,
                           MongoAutoConfiguration.class,
                           MongoDataAutoConfiguration.class})
@EnableRetry
@EnableMongoAuditing
@EnableAsync
@EnableScheduling
public class Application {
public static void main(String[] args) {
    SpringApplication.run(Application.class, args);
}
}

解决方案

1. 校验绑定配置一致性

确保测试中使用的INPUT_TOPIC_NAME、OUTPUT_TOPIC_NAME与应用配置的绑定目标完全匹配。在test.properties中添加或检查以下配置(替换实际通道名称):

# 输入通道绑定
spring.cloud.stream.bindings.input.destination=tax-datafeed-evt-load-clt-dnr
# 输出通道绑定
spring.cloud.stream.bindings.output.destination=tax-statusdata-evt-load-clt-dnr

2. 确认通道Bean定义

检查应用中是否正确定义了输入输出通道,示例如下:

public interface ClientDataChannels {
    String INPUT = "input";
    String OUTPUT = "output";

    @Input(INPUT)
    SubscribableChannel input();

    @Output(OUTPUT)
    MessageChannel output();
}

若使用传统绑定模式,需在配置类或主类上添加@EnableBinding(ClientDataChannels.class);若使用函数式模型,则需对应配置函数绑定规则。

3. 检查测试类依赖注入

确保测试类中正确注入InputDestination和OutputDestination:

@Autowired
private InputDestination inputDestination;

@Autowired
private OutputDestination outputDestination;

4. 适配函数式绑定(若使用Spring Cloud Stream 4.x)

如果采用函数式编程模型,需在test.properties中补充函数绑定配置:

spring.cloud.stream.function.definition=yourFunctionName
spring.cloud.stream.bindings.yourFunctionName-in-0.destination=tax-datafeed-evt-load-clt-dnr
spring.cloud.stream.bindings.yourFunctionName-out-0.destination=tax-statusdata-evt-load-clt-dnr

5. 确保测试绑定配置生效

确认@ImportAutoConfiguration(classes = TestChannelBinderConfiguration.class)未被其他配置排除,可尝试移除exclude中不必要的配置项,排查是否存在配置冲突。

内容的提问来源于stack exchange,提问作者sekhar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 17:45:01