【问题标题】:@StreamListener not receiving message from kafka topic@StreamListener 没有收到来自 kafka 主题的消息
【发布时间】:2018-09-15 00:52:43
【问题描述】:

我可以使用代码发送和接收消息:

@EnableBinding(Processor.class)
public class KafkaStreamsConfiguration {
  @StreamListener(Processor.INPUT)
  @SendTo(Processor.OUTPUT)
  public String processMessage(String message) {
    System.out.println("message = " + message);
    return message.replaceAll("my", "your");
  }
}


@RunWith(SpringRunner.class)
@SpringBootTest
@DirtiesContext
public class StreamApplicationIT {
private static String topicToPublish = "eventUpdateFromEventModel";

@BeforeClass
public static void setup() {
    System.setProperty("spring.kafka.bootstrap-servers", embeddedKafka.getBrokersAsString());
}

@Autowired
private KafkaMessageSender<String> kafkaMessageSenderToTestErrors;

@Autowired
private KafkaMessageSender<EventNotificationDto> kafkaMessageSender;

@ClassRule
public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, topicToPublish);

@Autowired
private Processor pipe;

@Autowired
private MessageCollector messageCollector;


@Rule
public OutputCapture outputCapture = new OutputCapture();

@Test
public void working() {
    pipe.input()
            .send(MessageBuilder.withPayload("This is my message")
                    .build());

    Object payload = messageCollector.forChannel(pipe.output())
            .poll()
            .getPayload();

    assertEquals("This is your message", payload.toString());
}

@Test
public void non_working() {
    kafkaMessageSenderToTestErrors.send(topicToPublish, "This was my message");
    assertTrue(isMessageReceived("This was your message", 50));
}

private boolean isMessageReceived(final String msg, final int maxAttempt) {
    return IntStream.rangeClosed(0, maxAttempt)
            .peek(a -> {
                try {
                    TimeUnit.MILLISECONDS.sleep(100);
                } catch (InterruptedException e) {
                    fail();
                }
            }).anyMatch(i -> outputCapture.toString().contains(msg));
}

}

@Service
@Slf4j
public class KafkaMessageSender<T> {
    private final KafkaTemplate<String, byte[]> kafkaTemplate;
    private final ObjectWriter objectWriter;

    public KafkaMessageSender(KafkaTemplate<String, byte[]> kafkaTemplate, ObjectMapper objectMapper) {
        this.kafkaTemplate = kafkaTemplate;
        this.objectWriter = objectMapper.writer();
    }

    public void send(String topicName, T payload) {
        try {
            kafkaTemplate.send(topicName, objectWriter.writeValueAsString(payload).getBytes());
        } catch (JsonProcessingException e) {
            log.info("error converting object into byte array {}", payload.toString().substring(0, 50));
        }
        log.info("sent payload to topic='{}'", topicName);
    }
}

但是当我使用 kafkaTemplate 将消息发送到任何主题时,StreamListener 不会收到消息。

spring.cloud.stream.bindings.input.group=test
spring.cloud.stream.bindings.input.destination=eventUpdateFromEventModel

我的 pom.xml:

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-stream</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-stream-kafka</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-stream-test-support</artifactId>
    <scope>test</scope>
</dependency>

    <!-- Spring boot version -->
    <spring.boot.version>1.5.7.RELEASE</spring.boot.version>
    <spring-cloud.version>Edgware.SR3</spring-cloud.version>

<dependencyManagement>
      <dependencies>

        <dependency>
            <!-- Import dependency management from Spring Boot -->
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-dependencies</artifactId>
            <version>${spring.boot.version}</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>

        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-dependencies</artifactId>
            <version>${spring-cloud.version}</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>

    </dependencies>
</dependencyManagement>

【问题讨论】:

  • 1.您没有显示您的 KafkaTemplate 使用情况。 2.你没有显示你正在使用的Kafka版本(至少那个starter-Stream-jafka)。 3. 你没有展示你的用例发生了什么错误。感谢您的理解,但很难帮助您解决问题的当前状态
  • 嗨@ArtemBilan 我刚刚更新了所有细节
  • 谢谢,我明白了。你真的从send() 方法发送到eventUpdateFromEventModel 主题吗?
  • 是的@ArtemBilan,我发送到eventUpdateFromEventModel 。我已经用工作和非工作更新了测试
  • 你那里有什么错误吗?

标签: java apache-kafka spring-cloud spring-cloud-stream spring-kafka


【解决方案1】:

工作

Object payload = messageCollector.forChannel(pipe.output())
        .poll()
        .getPayload();

...

不工作

KafkaTemplate

这是因为您在测试中使用的是 TestBinder,而不是真正的 Kafka 代理和 kafka binder。

消息收集器只是从通道中获取它。如果您想使用真正的 Kafka 代理进行测试,请参阅 test-embedded-kafka sample app

编辑

我刚刚测试了示例的 Ditmars (boot 1.5.x) 版本,它运行良好...

【讨论】:

  • 嗨@Gary 我尝试用真正的kafka 进行测试,并且kafkalistener 正在正确接收消息。只是 StreamListener 没有收到消息。
  • 我也尝试了您的示例代码,但没有成功。我只能看到你的spring boot版本是2,我的是1.5.7的区别
  • 开机1.5.x/Ditmars版本为here
  • 该示例对我来说很好 - 您是否从类路径中删除了 spring-cloud-stream-test-support? (这就是切换到测试活页夹的原因)。
  • 谢谢,Gary,这两个示例都不能立即在我的项目中运行,但只需进行一些更改,它就可以运行了。如果这对其他人有帮助,我将单独添加我的完整解决方案。
猜你喜欢
  • 2021-12-10
  • 1970-01-01
  • 2015-05-20
  • 1970-01-01
  • 2018-10-27
  • 2020-05-30
  • 2019-11-15
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多