我有一个高吞吐量 kafka 生产者的用例,我想每秒推送数千条 json 消息。
我有一个 3 节点 kafka 集群,我正在使用最新的 kafka-python 库,并有以下方法来生成消息
def publish_to_kafka(topic):
data = get_data(topic)
producer = KafkaProducer(bootstrap_servers=['b1', 'b2', 'b3'],
value_serializer=lambda x: dumps(x).encode('utf-8'), compression_type='gzip')
try:
for obj in data:
producer.send(topic, value=obj)
except Exception as e:
logger.error(e)
finally:
producer.close()
Run Code Online (Sandbox Code Playgroud)
我的主题有 3 个分区。
方法有时可以正常工作,但会失败并出现错误“KafkaTimeoutError:无法在 60.0 秒后更新元数据。”
我需要更改哪些设置才能使其顺利工作?