【问题标题】:Clear kafka topics for unit testing清除单元测试的kafka主题
【发布时间】:2016-12-15 12:36:20
【问题描述】:

我需要对 kafka 应用程序执行单元测试,避免使用第三方库。

我现在的问题是我想清除测试之间的所有主题,但我不知道如何。

这是我的临时解决方案:提交每次测试后产生的每条消息,并将所有测试消费者放在同一个消费者组中。

override protected def afterEach():Unit={
    val cleanerConsumer= newConsumer(Seq.empty)
    val topics=cleanerConsumer.listTopics()
    println("pulisco")
    cleanerConsumer.subscribe(topics.keySet())
    cleanerConsumer.poll(100)
    cleanerConsumer.commitSync()
    cleanerConsumer.close()
}

这不起作用,我不知道为什么。

例如,当我在测试中创建一个新的消费者时,messages 包含上一个测试中产生的消息。

val consumerProbe = newConsumer(SMSGatewayTopic)

val messages = consumerProbe.poll(1000)

我该如何解决这个问题?

【问题讨论】:

  • Kafka 是一个持久的消息存储,它让消费者决定开始消费的偏移量。你只需要记住你最后一次删除的偏移量,然后开始消费。
  • 为什么这个问题用 Java 标记?结果,它出现在非常特定于 java 的搜索中。不是很有帮助。

标签: java scala unit-testing apache-kafka


【解决方案1】:

您还可以在您的测试源中嵌入 Kafka/Zookeeper 实例,以便对此类隔离服务拥有更多控制权。

trait Kafka { self: ZooKeeper =>
  Kafka.start()
}

object Kafka {
  import org.apache.hadoop.fs.FileUtil
  import kafka.server.KafkaServer

  @volatile private var started = false

  lazy val logDir = java.nio.file.Files.createTempDirectory("kafka-log").toFile

  lazy val kafkaServer: KafkaServer = {
    val config = com.typesafe.config.ConfigFactory.
      load(this.getClass.getClassLoader)

    val (host, port) = {
      val (h, p) = config.getString("kafka.servers").span(_ != ':')
      h -> p.drop(1).toInt
    }

    val serverConf = new kafka.server.KafkaConfig({
      val props = new java.util.Properties()
      props.put("port", port.toString)
      props.put("broker.id", port.toString)
      props.put("log.dir", logDir.getAbsolutePath)

      props.put(
        "zookeeper.connect",
        s"localhost:${config getInt "test.zookeeper.port"}"
      )

      props
    })

    new KafkaServer(serverConf)
  }

  def start(): Unit = if (!started) {
    try {
      kafkaServer.startup()
      started = true
    } catch {
      case err: Throwable =>
        println(s"fails to start Kafka: ${err.getMessage}")
        throw err
    }
  }

  def stop(): Unit = try {
    if (started) kafkaServer.shutdown()
  } finally {
    FileUtil.fullyDelete(logDir)
  }
}

trait ZooKeeper {
  ZooKeeper.start()
}

object ZooKeeper {
  import java.nio.file.Files
  import java.net.InetSocketAddress
  import org.apache.hadoop.fs.FileUtil
  import org.apache.zookeeper.server.ZooKeeperServer
  import org.apache.zookeeper.server.ServerCnxnFactory

  @volatile private var started = false
  lazy val logDir = Files.createTempDirectory("zk-log").toFile
  lazy val snapshotDir = Files.createTempDirectory("zk-snapshots").toFile

  lazy val (zkServer, zkFactory) = {
    val srv = new ZooKeeperServer(
      snapshotDir, logDir, 500
    )

    val config = com.typesafe.config.ConfigFactory.
      load(this.getClass.getClassLoader)
    val port = config.getInt("test.zookeeper.port")

    srv -> ServerCnxnFactory.createFactory(
      new InetSocketAddress("localhost", port), 1024
    )
  }

  def start(): Unit = if (!zkServer.isRunning) {
    try {
      zkFactory.startup(zkServer)

      started = true

      while (!zkServer.isRunning) {
        Thread.sleep(500)
      }
    } catch {
      case err: Throwable =>
        println(s"fails to start ZooKeeper: ${err.getMessage}")
        throw err
    }
  }

  def stop(): Unit = try {
    if (started) zkFactory.shutdown()
  } finally {
    try { FileUtil.fullyDelete(logDir) } catch { case _: Throwable => () }
    FileUtil.fullyDelete(snapshotDir)
  }
}

测试类可以extends Kafka with ZooKeeper 确保它可用。

如果测试JVM没有fork,可以通过SBT testOptions in Test设置中的Tests.Cleanup在测试后停止嵌入式服务。

【讨论】:

    【解决方案2】:

    我建议,您只需在测试之前重新创建所有主题。例如,这是 kafka 测试创建/删除主题的方式:

    Kafka repository on GitHub

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2013-04-23
      • 2020-06-22
      • 2017-03-13
      • 1970-01-01
      • 1970-01-01
      • 2021-12-10
      • 2013-07-21
      • 1970-01-01
      相关资源
      最近更新 更多