【问题标题】:How to Test Kafka Consumer如何测试 Kafka 消费者
【发布时间】:2017-06-16 14:20:20
【问题描述】:

我有一个 Kafka Consumer(内置于 Scala),它从 Kafka 中提取最新记录。消费者看起来像这样:

val consumerProperties = new Properties()
consumerProperties.put("bootstrap.servers", "localhost:9092")
consumerProperties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
consumerProperties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
consumerProperties.put("group.id", "something")
consumerProperties.put("auto.offset.reset", "latest")

val consumer = new KafkaConsumer[String, String](consumerProperties)
consumer.subscribe(java.util.Collections.singletonList("topic"))

现在,我想为它编写一个集成测试。是否有任何方法或最佳实践来测试 Kafka 消费者?

【问题讨论】:

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


    【解决方案1】:
    1. 您需要以编程方式启动 zookeeper 和 kafka 以进行集成测试。

      1.1 启动zookeeper (ZooKeeperServer)

      def startZooKeeper(zooKeeperPort: Int, zkLogsDir: Directory): ServerCnxnFactory = {
          val tickTime = 2000
      
          val zkServer = new ZooKeeperServer(zkLogsDir.toFile.jfile, zkLogsDir.toFile.jfile, tickTime)
      
          val factory = ServerCnxnFactory.createFactory
          factory.configure(new InetSocketAddress("0.0.0.0", zooKeeperPort), 1024)
          factory.startup(zkServer)
      
          factory
      }
      

      1.2启动kafka(KafkaServer)

      case class StreamConfig(streamTcpPort: Int = 9092,
                          streamStateTcpPort :Int = 2181,
                          stream: String,
                          numOfPartition: Int = 1,
                          nodes: Map[String, String] = Map.empty)
      
      def startKafkaBroker(config: StreamConfig,
                         kafkaLogDir: Directory): KafkaServer = {
      
        val syncServiceAddress = s"localhost:${config.streamStateTcpPort}"
      
        val properties: Properties = new Properties
        properties.setProperty("zookeeper.connect", syncServiceAddress)
        properties.setProperty("broker.id", "0")
        properties.setProperty("host.name", "localhost")
        properties.setProperty("advertised.host.name", "localhost")
        properties.setProperty("port", config.streamTcpPort.toString)
        properties.setProperty("auto.create.topics.enable", "true")
        properties.setProperty("log.dir", kafkaLogDir.toAbsolute.path)
        properties.setProperty("log.flush.interval.messages", 1.toString)
        properties.setProperty("log.cleaner.dedupe.buffer.size", "1048577")
      
        config.nodes.foreach {
          case (key, value) => properties.setProperty(key, value)
        }
      
        val broker = new KafkaServer(new KafkaConfig(properties))
        broker.startup()
      
        println(s"KafkaStream Broker started at ${properties.get("host.name")}:${properties.get("port")} at ${kafkaLogDir.toFile}")
        broker
      

      }

    2. 使用KafkaProducer发送一些事件流式传输

    3. 然后与您的消费者一起消费以测试并验证其是否正常工作

    您可以使用具有startBroker 方法的scalatest-eventstream,它将为您启动Zookeeper 和Kafka。

    还有destroyBroker,它会在测试后清理你的kafka。

    例如。

    class MyStreamConsumerSpecs extends FunSpec with BeforeAndAfterAll with Matchers {
      implicit val config =
        StreamConfig(streamTcpPort = 9092, streamStateTcpPort = 2181, stream = "test-topic", numOfPartition = 1)
    
      val kafkaStream = new KafkaEmbeddedStream
    
      override protected def beforeAll(): Unit = {
        kafkaStream.startBroker
      }
    
      override protected def afterAll(): Unit = {
        kafkaStream.destroyBroker
      }
    
      describe("Kafka Embedded stream") {
        it("does consume some events") {
    
          //uses application.properties
          //emitter.broker.endpoint=localhost:9092
          //emitter.event.key.serializer=org.apache.kafka.common.serialization.StringSerializer
          //emitter.event.value.serializer=org.apache.kafka.common.serialization.StringSerializer
          kafkaStream.appendEvent("test-topic", """{"MyEvent" : { "myKey" : "myValue"}}""")
    
          val consumerProperties = new Properties()
          consumerProperties.put("bootstrap.servers", "localhost:9092")
          consumerProperties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
          consumerProperties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
          consumerProperties.put("group.id", "something")
          consumerProperties.put("auto.offset.reset", "earliest")
    
          val myConsumer = new KafkaConsumer[String, String](consumerProperties)
          myConsumer.subscribe(java.util.Collections.singletonList("test-topic"))
    
          val events = myConsumer.poll(2000)
    
          events.count() shouldBe 1
          events.iterator().next().value() shouldBe """{"MyEvent" : { "myKey" : "myValue"}}"""
          println("=================" + events.count())
        }
      }
    }
    

    【讨论】:

    • 感谢您的回复!但我在这里有一个疑问。我要测试的卡夫卡消费者有auto.offset.resetlatest,而不是earliest。因此,当我尝试上述解决方案时,消费者无法收到任何消息。您能否提及如何使用latest 偏移量测试 Kafka 消费者?
    • 在上面的示例中,我首先发出事件然后启动消费者,这就是我从earliest 处理的原因。如果要处理latest,首先启动消费者并发出事件,它将处理该事件。请参阅示例 - github.com/duwamish-os/scalatest-eventstream/blob/master/src/… 此外,您必须循环从流中进行轮询
    • 配置latest时添加更多,意味着consumer会处理consumer启动后才发出的消息。
    • 我正在使用这种方法,但是当我执行此代码时,它运行不完。生产者做了一些事情(我添加了日志),但似乎消费者没有被执行。有什么想法吗?
    猜你喜欢
    • 2021-05-20
    • 2019-05-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-04-09
    • 1970-01-01
    • 1970-01-01
    • 2020-05-21
    相关资源
    最近更新 更多