【问题标题】:Getting messages sent from python KafkaProducer获取从 python KafkaProducer 发送的消息
【发布时间】:2017-07-11 11:53:43
【问题描述】:

我的目标是从非文件源(即在程序中生成或通过 API 发送)获取数据并将其发送到 spark 流。为此,我通过python-based KafkaProducer 发送数据:

$ bin/zookeeper-server-start.sh config/zookeeper.properties &
$ bin/kafka-server-start.sh config/server.properties &
$ bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic my-topic
$ python 
Python 3.6.1| Anaconda custom (64-bit)
> from kafka import KafkaProducer
> import time
> producer = KafkaProducer(bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8'))
> producer.send(topic = 'my-topic', value = 'MESSAGE ACKNOWLEDGED', timestamp_ms = time.time())
> producer.close()
> exit()

我的问题是从消费者 shell 脚本中检查主题时没有出现任何内容:

$ bin/kafka-console-consumer.sh --bootstrap-server localhost:2181 --topic my-topic
^C$

这里有什么遗漏或错误吗?我是 spark/kafka/messaging 系统的新手,所以任何事情都会有所帮助。 Kafka 版本是 0.11.0.0 (Scala 2.11),并且没有对配置文件进行任何更改。

【问题讨论】:

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


    【解决方案1】:

    如果您在向主题发送消息后启动消费者,消费者可能会跳过该消息,因为它将设置主题偏移量(可以被视为读取的“起点”)到主题的结尾。要改变这种行为,请尝试添加 --from-beginning 选项:

    $ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic my-topic --from-beginning
    

    你也可以试试kafkacat,比Kafka的控制台消费者和生产者(恕我直言)更方便。使用kafkacat从Kafka读取消息可以使用以下命令:

    kafkacat -C -b 'localhost:9092' -o beginning -e -D '\n' -t 'my-topic'
    

    希望它会有所帮助。

    【讨论】:

    • 我添加了from-beginning,但结果是一样的。我还安装了 kafkacat,重做了步骤并运行了您的命令,但仍然找不到消息。
    • @user2361174 刚刚检查了您的示例,由于timestamp_ms=time.time(),生产者似乎没有发送任何内容——如果打开调试日志记录,日志中将出现以下消息:DEBUG:kafka.producer.kafka:Exception occurred during message send: <class 'struct.error'>。可能time.time() 以生产者意想不到的格式返回时间戳...因此删除该选项应该可以解决问题,即producer.send(topic='my-topic', value='MESSAGE ACKNOWLEDGED')(默认情况下将使用当前时间戳)。
    • 我去掉了时间戳,尽管消费者或 kafkacat 命令中仍然没有显示任何内容。我在这里保存了命令行输出:raw.githubusercontent.com/dretta/spark/master/kafka.log
    【解决方案2】:

    我发现了问题,value_serializer 默默地崩溃了,因为我没有将 json 模块导入解释器。对此有两种解决方案,一种是简单地导入模块,您将得到"MESSAGE ACKNOWLEDGED"(带引号)。或者,您可以完全删除 value_serializer 并将在下一行发送的 value 字符串转换为字节字符串(即 Python 3 的 b'MESSAGE ACKNOWLEDGED'),这样您将得到不带引号的消息。

    我还将 Kafka 切换到版本 0.10.2.1 (Scala 2.11),因为 Kafka-python 文档中没有确认它与版本 0.11.0.0 兼容

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-03-30
      • 2016-09-21
      • 2017-12-06
      • 1970-01-01
      • 2018-11-23
      • 2015-05-14
      • 1970-01-01
      相关资源
      最近更新 更多