Как я могу уменьшить отставание в кафке потребителя / производителя - PullRequest
0 голосов
/ 21 ноября 2018

Я ищу улучшения в коде scala kafka.Для уменьшения задержки, что я должен делать в потребителя и производителя.Это код, который я получил от кого-то.Я знаю, что этот код не сложный код.Но я никогда не видел scala-код раньше, и я только начинаю узнавать о Кафке.Поэтому мне трудно найти проблему.

import java.util.Properties
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}

import scala.util.Try

class KafkaMessenger(val servers: String, val sender: String) {
  val props = new Properties()
  props.put("bootstrap.servers", servers)
  props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
  props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
  props.put("producer.type", "async")

  val producer = new KafkaProducer[String, String](props)

  def send(topic: String, message: Any): Try[Unit] = Try {
    producer.send(new ProducerRecord(topic, message.toString))
  }

  def close(): Unit = producer.close()
}

object KafkaMessenger {
  def apply(host: String, topic: String, sender: String, message: String): Unit = {
    val messenger = new KafkaMessenger(host, sender)
    messenger.send(topic, message)
    messenger.close()
  }
}

, и это код потребителя.

import java.util.Properties
import java.util.concurrent.Executors

import com.satreci.g2gs.common.impl.utils.KafkaMessageTypes._
import kafka.admin.AdminUtils
import kafka.consumer._
import kafka.utils.ZkUtils
import org.I0Itec.zkclient.{ZkClient, ZkConnection}
import org.slf4j.LoggerFactory

import scala.language.postfixOps

class KafkaListener(val zookeeper: String,
                    val groupId: String,
                    val topic: String,
                    val handleMessage: ByteArrayMessage => Unit,
                    val workJson: String = ""
                   ) extends AutoCloseable {
  private lazy val logger = LoggerFactory.getLogger(this.getClass)
  val config: ConsumerConfig = createConsumerConfig(zookeeper, groupId)
  val consumer: ConsumerConnector = Consumer.create(config)
  val sessionTimeoutMs: Int = 10 * 1000
  val connectionTimeoutMs: Int = 8 * 1000
  val zkClient: ZkClient = ZkUtils.createZkClient(zookeeper, sessionTimeoutMs, connectionTimeoutMs)
  val zkUtils = new ZkUtils(zkClient, new ZkConnection(zookeeper), false)

  def createConsumerConfig(zookeeper: String, groupId: String): ConsumerConfig = {
    val props = new Properties()
    props.put("zookeeper.connect", zookeeper)
    props.put("group.id", groupId)
    props.put("auto.offset.reset", "smallest")
    props.put("zookeeper.session.timeout.ms", "5000") 
    props.put("zookeeper.sync.time.ms", "200")
    props.put("auto.commit.interval.ms", "1000")
    props.put("partition.assignment.strategy", "roundrobin")
    new ConsumerConfig(props)
  }

  def run(threadCount: Int = 1): Unit = {
    val streams = consumer.createMessageStreamsByFilter(Whitelist(topic), threadCount)

    if (!AdminUtils.topicExists(zkUtils, topic)) {
      AdminUtils.createTopic(zkUtils, topic, 1, 1)
    }

    val executor = Executors.newFixedThreadPool(threadCount)
    for (stream <- streams) {
      executor.submit(new MessageConsumer(stream))
    }
    logger.debug(s"KafkaListener start with ${threadCount}thread (topic=$topic)")
  }

  override def close(): Unit = {
    consumer.shutdown()
    logger.debug(s"$topic Listener close")
  }

  class MessageConsumer(val stream: MessageStream) extends Runnable {
    override def run(): Unit = {
      val it = stream.iterator()
      while (it.hasNext()) {
        val message = it.next().message()
        if (workJson == "") {
          handleMessage(message)
        }
        else {
          val strMessage = new String(message)
          val newMessage = s"$strMessage/#/$workJson"
          val outMessage = newMessage.toCharArray.map(c => c.toByte)
          handleMessage(outMessage)
        }
      }
    }
  }
}

В частности, я хочу изменить структуру, которая создает объекты KafkaProduce всякий раз, когда я отправляюсообщение.Кажется, есть много других улучшений, чтобы уменьшить отставание.

1 Ответ

0 голосов
/ 21 ноября 2018

Увеличение количества экземпляров-потребителей (KafkaListener) с одинаковым идентификатором группы.Это увеличит норму потребления.В конечном итоге ваш разрыв между производителями и потребителями будет сведен к минимуму.

...