【问题标题】:How to perform functional testing of a kafka consumer application?如何对 kafka 消费者应用程序进行功能测试?
【发布时间】:2018-12-18 12:40:59
【问题描述】:

我想测试一个使用 Kafka 消息并将它们写入日志的应用程序。这是它在类似 scala 的伪代码中的近似表示:

import kafka.consumer.Consumer
import kafka.consumer.ConsumerConfig
import org.slf4j.LoggerFactory
import java.util.Properties
import java.util.HashMap

object ConsumerApp extends App {
  val topic = new HashMap[String, Integer]()
  topic.put("test", 1)
  val logger = LoggerFactory.getLogger(getClass().getName())
  val messageStream = Consumer
    .createJavaConsumerConnector(new ConsumerConfig(new Properties()))
    .createMessageStreams(topic)
    .get(topic).get(0)
  for (message <- messageStream) {
    val gotMessage = new String(message.message())
    logger.info(gotMessage)
  }
}

我想到的测试场景如下:

  • Kafka 服务器已启动。

  • 应用程序启动并连接到 Kafka 服务器,开始监听特定主题的消息。

  • 向主题发送消息。

  • 应用程序使用消息并记录它。

以下是类似 Scala 的伪代码的测试草稿:

import uk.org.lidalia.slf4jtest.TestLoggerFactory;
import uk.org.lidalia.slf4jtest.LoggingEvent.info;

abstract class UnitSpec extends FlatSpec with Matchers with EmbeddedKafka {

}

class ConsumerAppSpec extends UnitSpec {
  "ConsumerApp" should "consume and log messages from Kafka on specific topic" in {
    withRunningKafka {
      val consumer = ConsumerApp
      // interecept logger to be able to test that the kafka message is logged
      val logger = TestLoggerFactory.getTestLogger(consumer.getClass)
      // start the application, but beforehand do something to prevent it infinitely blocking
      ???
      consumer.main(Array())  

      // publish a test message    
      publishStringMessageToKafka("test", "TEST")
      // Confirm that the message has been properly logged
      ???
    }
    EmbeddedKafka.stop()
  }
}

我的问题是前三个问号处的测试代码。如果我执行 main() 方法,它不会终止,从而阻止执行剩余的测试。

【问题讨论】:

    标签: scala testing apache-kafka kafka-consumer-api scalatest


    【解决方案1】:

    让它发挥作用……

    object ConsumerApp extends App {
    
       def doStuff() {
          val topic = ...
          // more stuff ...
       }
    
       doStuff()
     }
    
     class ConsumerAppSpec {
       // ...
       consumer.doStuff()
       // ...
       publishMessage() 
     }
    

    更新

    我想,我误解了这个问题。而不是“不要调用main”,我应该说不要阻塞线程:)

     class ConsumerAppSpec {
       ...
       val futureResult = Future(consumer.doStuff)
       ...
       publishMessage()
     }
    

    【讨论】:

    • 这有什么帮助?如果 main(),这不会导致 doStuff() 阻塞测试吗?
    • 哦,我想,我误解了你在问什么。查看更新:)
    • 你如何避免使用 kafka 消费者呢?
    • @Worse_Username 你知道了吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-05-20
    • 1970-01-01
    • 1970-01-01
    • 2019-05-09
    相关资源
    最近更新 更多