【问题标题】:Kafka Producer does not write to kafka topicKafka Producer 不写入 kafka 主题
【发布时间】:2019-04-23 15:38:16
【问题描述】:

抱歉这个菜鸟问题:我正在使用 akka 给 kafka 写信,但我在 kafka 控制台使用者中看不到它。

写入 kafka 的配置:

kafka {
  bootstrap.servers = "localhost:9002"
  auto.offset.reset = "earliest"
}

我有代码可以使用 akka 写入 kafka 主题:

class ServiceKafkaProducer(topicName: String, actorSystem: ActorSystem, configuration: Configuration) {
  val bootstrapServers: String = configuration
    .getString("kafka.bootstrap.servers")
    .getOrElse(
      throw new Exception("No config element foe kafka.bootstrap.servers")
    )

  val producerSettings: ProducerSettings[String, String] = ProducerSettings(
    actorSystem,
    new StringSerializer,
    new StringSerializer
  ).withBootstrapServers(bootstrapServers)

  val producer: KafkaProducer[String, String] = producerSettings.createKafkaProducer()

  def send(logRecordStr: String): Unit = {
    Logger.debug(s"Inside ServiceKafkaProducer, writing to $topicName")
    Logger.debug(logRecordStr)
    producer.send(
      new ProducerRecord(topicName, logRecordStr)
    )
  }
}

 def createTag(text: String, createdBy: UUID): Unit = {
  Logger.debug("Inside TagEventProducer#createTag")
  val tagId = UUID.randomUUID()

  val event = TagCreated(tagId, text, createdBy)
  println(event)
  val record = createLogRecord(event)

  send(record.encode)
}

日志

```[debug] - application - Inside TagEventProducer#createTag
TagCreated(d393d223-9eb6-45e3-8610-56a3f65c84cc,scala,f5b61ca0-0ccc-4064-94c1-cba2a5a4087b)
[debug] - application - Inside ServiceKafkaProducer, writing to tags
[debug] - application - {"id":"ed27f0d1-6b6c-469b-af97-1929dc6a5cc7","action":"tag-created","data":{"id":"d393d223-9eb6-45e3-8610-56a3f65c84cc","text":"scala","createdBy":"f5b61ca0-0ccc-4064-94c1-cba2a5a4087b"},"timestamp":1542776716868}```

(已编辑) 我正在使用像这样的 spotify docker 映像运行 kafka:

version: '3.5'
services:
  kafka:
    image: 'spotify/kafka'
    hostname: kafka
    environment:
      - ADVERTISED_HOST=kafka
      - ADVERTISED_PORT=9092
    ports:
      - "9092:9092"
      - "2181:2181"
    volumes:
      - /var/run/docker.sock:/var/run/docker.sock
    networks:
      - kafka_net
  kafkaManager:
    image: 'sheepkiller/kafka-manager'
    environment:
      - ZK_HOSTS=kafka:2181
      - APPLICATION_SECRET=letmein
    ports:
      - "8000:8000"
    networks:
      - kafka_net
networks:
  kafka_net:
    name: my_network

【问题讨论】:

    标签: scala apache-kafka akka


    【解决方案1】:

    首先,您在生产者中使用了错误的端口(9002 而不是 9092)。尝试在生产者中使用 bootstrap.servers = kafka:9092 而不是 localhost:9002,因为您的广告主机设置为 kafka

    【讨论】:

      【解决方案2】:

      你确定你写给Kafka的配置如下吗?

      bootstrap.servers = "localhost:9002"
      auto.offset.reset = "earliest"
      

      我认为应该是localhost:9092。你可能在这里打错了。

      【讨论】:

        【解决方案3】:

        kafka.bootstrap.servers 设置为{kafka_docker_host_api}:9092。如果您的问题无法解决,请尝试使用producer.flush()。 Kafka 生产者不会立即发送消息。所以如果你想强制发送消息,你需要刷新。

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 2020-01-10
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2019-07-12
          • 2016-01-21
          • 2019-02-14
          相关资源
          最近更新 更多