【发布时间】:2022-01-16 12:39:00
【问题描述】:
我正在尝试让 kafka 消费者订阅 2 个主题,但每当我尝试分配多个主题时,我都会收到以下消息:
INFO:__main__:Controller module is running and listening...
WARNING:kafka.coordinator.consumer:group_id is None: disabling auto-commit.
INFO:kafka.consumer.subscription_state:Updating subscribed topics to: ('folder-data', 'password')
INFO:kafka.conn:<BrokerConnection node_id=bootstrap-0 host=172.17.0.1:9092 <connecting> [IPv4 ('172.17.0.1', 9092)]>: connecting to 172.17.0.1:9092 [('172.17.0.1', 9092) IPv4]
INFO:kafka.conn:Probing node bootstrap-0 broker version
INFO:kafka.conn:<BrokerConnection node_id=bootstrap-0 host=172.17.0.1:9092 <connecting> [IPv4 ('172.17.0.1', 9092)]>: Connection complete.
INFO:kafka.conn:Broker version identified as 2.5.0
INFO:kafka.conn:Set configuration api_version=(2, 5, 0) to skip auto check_version requests on startup
INFO:__main__:Sending message: <Logger __main__ (INFO)> to topics: find-password, analyze-folder
INFO:kafka.conn:<BrokerConnection node_id=1001 host=172.17.0.1:9092 <connecting> [IPv4 ('172.17.0.1', 9092)]>: connecting to 172.17.0.1:9092 [('172.17.0.1', 9092) IPv4]
INFO:kafka.conn:<BrokerConnection node_id=1001 host=172.17.0.1:9092 <connecting> [IPv4 ('172.17.0.1', 9092)]>: Connection complete.
INFO:kafka.conn:<BrokerConnection node_id=bootstrap-0 host=172.17.0.1:9092 <connected> [IPv4 ('172.17.0.1', 9092)]>: Closing connection.
INFO:__main__:Sent.
INFO:kafka.conn:<BrokerConnection node_id=bootstrap-0 host=172.17.0.1:9092 <connecting> [IPv4 ('172.17.0.1', 9092)]>: connecting to 172.17.0.1:9092 [('172.17.0.1', 9092) IPv4]
INFO:kafka.conn:<BrokerConnection node_id=bootstrap-0 host=172.17.0.1:9092 <connecting> [IPv4 ('172.17.0.1', 9092)]>: Connection complete.
WARNING:kafka.cluster:Topic password is not available during auto-create initialization
WARNING:kafka.cluster:Topic folder-data is not available during auto-create initialization
INFO:kafka.consumer.subscription_state:Updated partition assignment: []
由于某种原因,如果我之后重新启动容器,它就可以工作。
这是我的代码:
analye_folder = 'analyze-folder'
folder_data = 'folder-data'
find_password = 'find-password'
get_password = 'password'
consumer = KafkaConsumer (auto_offset_reset='earliest',
bootstrap_servers=bootstrap_servers,
api_version=(0,10))
consumer.subscribe([get_password, folder_data])
我该如何解决这个问题?
【问题讨论】:
标签: python apache-kafka kafka-python