【问题标题】:Apache nifi and kafka microserviceApache nifi 和 kafka 微服务
【发布时间】:2021-09-23 12:32:38
【问题描述】:

我是 Apache Nifi 的新手,但在尝试将 kafka 微服务(与生产者)连接到 Apache nifi 消费者时遇到了一些问题。

基本上,我有一个像这样的 docker-compose:

zookeeper:
  container_name: zookeeper_test
  image: wurstmeister/zookeeper #zookeeper:3.5.7
  ports:
  - 2181:2181

kafka:
  container_name: kafka_test
  image: wurstmeister/kafka #:2.13-2.6.0
  ports:
  - 9092:9092
  environment:
      KAFKA_ADVERTISED_HOST_NAME: kafka
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
  depends_on: 
      - zookeeper
      - kafkaui

kafkaui:
  container_name: kafka-ui_test
  image: provectuslabs/kafka-ui:latest
  environment: 
      - KAFKA_CLUSTERS_0_NAME=kafka
      - KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=kafka:9092
      - KAFKA_CLUSTERS_0_ZOOKEEPER= zookeeper:2181
  ports:
      - 6789:8080

test:
  container_name: test
  build:
      context: ./test
      dockerfile: Dockerfile
  depends_on:
      - kafka
  command: python test.py

测试是我的制作人:

from kafka import KafkaProducer
import json
from time import sleep

producer = KafkaProducer(bootstrap_servers='kafka:9092')
json_message = {"hello":"world"}


for i in range(1000):
   producer.send("INPUT", json.dumps(json_message).encode('utf-8')) 
   producer.flush()
   sleep(1)

通过 KafkaUI,我可以看到已发送的主题 INPUT。

在 Apache nifi 仪表板中,我使用以下参数设置了 ConsumerKafka_2.6: 卡夫卡经纪人:本地主机:9092 组号:1 主题:输入

然后我在“成功”时连接到这个漏斗,只是为了查看收到的消息。不幸的是,这样做,我没有看到任何收到的东西。我只是在 consumerkafka 框中看到了很多任务,但队列中没有任何元素连接到漏斗。我希望看到收到的 json,不是吗?我可以错过什么吗?

【问题讨论】:

  • 生产者是否在 NiFi 开始消费后运行?或者,NiFi 是从最早还是最晚开始消费?
  • 生产者在 NiFi 消费者之前开始。偏移重置设置为最新
  • 您是否尝试过 Kafka 控制台消费者来验证主题中有消息?将起始偏移设置为 EARLIEST。主题名称是什么? Python 生产者说 topic = INPUT 但你的 NiFi 配置说 Topic = SL.CPTI.INPUT。
  • 嗨,是的。我写了输入,因为我也尝试了不同的主题。不在乎……两个人的话题都是一样的。我还尝试了 Kafka ui 微服务,以查看该主题是否可用。所以这不是主题或发送消息的问题

标签: docker apache-kafka microservices apache-nifi


【解决方案1】:

根据@alexmark 的评论,NiFi 消费者已将偏移重置设置为最新。这意味着如果消费者组没有提交的偏移量,NiFi 将有效地忽略已经产生到主题的消息。

将偏移重置设置为最早将改变该行为以消费主题中的所有消息,之后将提交最后消费的偏移量,以便下次消费开始时它将是消费者停止的地方。

如果将偏移重置更改为最早没有任何效果,则 NiFi 很可能提交了之前运行时的偏移:更改消费者组(保留偏移重置为最早)是解决该问题的最简单方法(内部有工具Kafka 检查和修改提交的偏移量)

【讨论】:

  • 您好李维,感谢您的回复。我做了你写给我的事,但没有任何效果。卡夫卡仍在运行,我没有阻止它。我还添加了 apache nifi 的配置和工作流程的屏幕截图。
  • PS 你可以看到队列总是0
【解决方案2】:

我发现问题出在与 kafka 代理的连接上。当我 dockerize apache nifi(而不是将其用作 localhost)时,将其作为 kafka 浏览器 url kafka:9092,它可以完美运行。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-07-19
    • 2016-11-17
    • 1970-01-01
    • 2022-12-16
    • 2019-02-17
    • 1970-01-01
    • 2021-11-19
    • 2018-04-14
    相关资源
    最近更新 更多