使用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
相关产品推荐
相关产品推荐

