【问题标题】:Unable to publish in Kafka无法在 Kafka 中发布
【发布时间】:2018-06-30 18:08:30
【问题描述】:

我想在 Kafka 主题中发布 我无法这样做,程序停止。 我收到此错误:

KafkaTimeoutError:60.0 秒后更新元数据失败。

def saveResults(response):
    entities_tweet = response["entities"]
    for entity in entities_tweet:
        try:
            for i in entity_dict:
                for j in entity_dict[i]:
                    if(entity["text"] in j):
                        entity["tweet"] = response["tweet"]
                        entity["tweetId"] = response["tweetId"]
                        entity["timeStamp"] = response["timeStamp"]
                        #entity["userProfile"] = response["userProfile"]
                        future = producer.send('argentina-iceland-june-16-watson', bytes(entity))
                        print("Published.")
                    else:
                        print("All ignored.")
                        future = producer.send('argentina-iceland-june-16-watson', bytes(entity))
                        print("Published")
        except Exception as e:
            print (e)
        finally:
            producer.flush()

但是,这是有效的:

from kafka import KafkaProducer
from kafka.errors import KafkaError

producer = KafkaProducer(bootstrap_servers=['broker1:1234'])

# Asynchronous by default
future = producer.send('my-topic', b'raw_bytes')

【问题讨论】:

  • 如果你想把 Mongo 读到 Kafka,看看 Debezium 项目

标签: apache-kafka kafka-python


【解决方案1】:

您似乎使用了不正确的 boostrap 服务器,它应该是 broker1:9092 而不是 broker1:1234...

【讨论】:

  • 改了没有改善。
  • 检查您是否真的可以从您运行生产者的机器访问此主机名/端口
  • 我可以。此外,问题在于多线程。你能帮我吗,可能给我一个如何使用python的多线程在mongo中插入数据的例子?还是我应该为此再问一个问题?
猜你喜欢
  • 2021-11-19
  • 2016-02-06
  • 1970-01-01
  • 2018-05-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-01-05
  • 1970-01-01
相关资源
最近更新 更多