diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/BaseIntegrationFlowDefinition.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/BaseIntegrationFlowDefinition.java index 8407717f8f..153f1a75b6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/BaseIntegrationFlowDefinition.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/BaseIntegrationFlowDefinition.java @@ -1046,6 +1046,9 @@ public abstract class BaseIntegrationFlowDefinition, BeanFac Object[] args = buildArgs(message); try { - return this.method.invoke(this.target, args); + Object result = this.method.invoke(this.target, args); + if (result != null && org.springframework.integration.util.ClassUtils.isKotlinUnit(result.getClass())) { + result = null; + } + return result; } catch (InvocationTargetException e) { if (e.getTargetException() instanceof ClassCastException) { diff --git a/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt b/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt index 058a37e02c..df5f8643d8 100644 --- a/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt +++ b/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt @@ -1,5 +1,5 @@ /* - * Copyright 2020 the original author or authors. + * Copyright 2020-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -212,8 +212,8 @@ class KotlinDslTests { fun `no reply from handle`() { val payloadReference = AtomicReference() val integrationFlow = - integrationFlow("handlerInputChanenl") { - handle { payload, _ -> payloadReference.set(payload) } + integrationFlow("handlerInputChannel") { + handle> { message, _ -> payloadReference.set(message.payload) } } val registration = this.integrationFlowContext.registration(integrationFlow).register()