TensorFlow 1.7使用Kafka API报错:Failed to consume:Broker: No more messages求助
Hey there, let's break down this error you're hitting when using TensorFlow 1.7's Kafka API. That "No more messages" error basically means your consumer is asking the Kafka broker for messages, but there's nothing left to pull for the specified topic/partition. Here are the most common fixes tailored to your code:
1. First, Verify Your Topic Partition Has Messages
Your code targets partition 0 of bt1_meeting_appeventlog_oracle (from the :0 in your topic string). Let's confirm there's actually data to consume:
- Use Kafka's built-in tool to check the latest offset for your partition:
Compare the returned offset for partition 0 with the starting offset you set (kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list kafka1:9092,kafka2:9092,kafka3:9092 --topic bt1_meeting_appeventlog_oracle --time -10in your topic string). If your starting offset is already at or above the latest offset, there's no data left to pull.
2. Tweak Your KafkaDataset Configuration
Fix the Consumer Group Setting
You've set group="None"—this runs your consumer in standalone mode, which means Kafka won't track your offset for you. If you want to keep consuming new messages as they're produced, use a valid consumer group name instead. Kafka will then remember where you left off, so you don't have to hardcode offsets every time:
temp = kafka.KafkaDataset( topics='bt1_meeting_appeventlog_oracle:0:0:-1', group="tf-meeting-log-consumer", # Replace with your own group name servers="kafka1:9092,kafka2:9092,kafka3:9092" )
Add a Timeout for Message Polling
By default, TensorFlow 1.7's KafkaDataset throws an error immediately when there are no messages. You can add a timeout to make it wait for new messages instead. Use KafkaOptions to set this:
kafka_options = tf.contrib.kafka.KafkaOptions( timeout=10000 # Wait 10 seconds (in milliseconds) before giving up ) temp = kafka.KafkaDataset( topics='bt1_meeting_appeventlog_oracle:0:0:-1', group="tf-meeting-log-consumer", servers="kafka1:9092,kafka2:9092,kafka3:9092", kafka_options=kafka_options )
This way, if the broker has no messages right now, your consumer will wait a bit for new ones instead of crashing immediately.
3. Catch the OutOfRangeError Gracefully
When your iterator runs out of messages, TensorFlow throws a tf.errors.OutOfRangeError. You can catch this to handle the situation cleanly instead of letting your program crash:
with tf.Session() as sess: try: while True: message = sess.run(next_element) # Process your message here print(f"Received message: {message}") except tf.errors.OutOfRangeError: print("No more messages available right now. Exiting gracefully.") # You could also add retry logic here if needed
4. Double-Check Kafka Cluster Connectivity
Make sure the machine running your TensorFlow code can actually reach the Kafka brokers. Test connectivity with a simple telnet command:
telnet kafka1 9092
If this fails, you've got a network issue blocking access to the broker, which would also cause the "no more messages" error (since the consumer can't reach the broker to pull data).
内容的提问来源于stack exchange,提问作者zhangzheng

