【发布时间】:2020-06-02 16:19:30
【问题描述】:
producer = KafkaProducer(bootstrap_servers='kf-p1l-node3:9092,xxxxx,xxxxx',
value_serializer=lambda x: dumps(x).encode('utf-8')) # utf-8
consumer = KafkaConsumer( bootstrap_servers='rdwh-node1:49092,xxxxx,xxxxx',
# bootstrap_servers='kf-p1l-node3:9092,xxxxx,xxxxx',
auto_offset_reset=param["AUTO_OFFSET_RESET"],
consumer_timeout_ms=param["CONSUMER_TIMEOUT_MS"],
enable_auto_commit=False,
auto_commit_interval_ms=60000,
group_id=param["GROUP_ID"],
client_id=param["CLIENT_ID"]
)
consumer.subscribe([param["TOPIC_IN"]])
如果 KafkaProducer 和 KafkaConsumer 的 bootstrap_server 相同,则此代码有效。但是如果将 KafkaConsumer 更改为另一台服务器,它就不起作用了
【问题讨论】:
标签: python apache-kafka kafka-consumer-api