使用Sarama调用SendMessage指定12分区却始终返回0分区的问题
Hey there, let's dig into why your messages are landing in partition 0 instead of the partition 12 you specified.
First, let's recap your setup for context:
- Kafka Version: kafka_2.12-1.1.0
- Go Version: 1.9.1
- Sarama producer code (as you provided):
package kafka import ( "flag" "fmt" "log" "strings" "github.com/Shopify/sarama" ) var partition = flag.Int("partition", 12, "The partition to produce to.") func Start_producer(payload []byte) { flag.Parse() s := "mydata" topic := &s // brokers := &[]string{"172.25.33.175:9092,172.25.33.176:9092,172.25.33.177:9092"} // brokers := []string{"172.25.33.175:9092,172.25.33.176:9092,172.25.33.177:9092"} config := sarama.NewConfig() config.Producer.RequiredAcks = sarama.WaitForAll config.Producer.Retry.Max = 5 config.Producer.Return.Successes = true producer, err := sarama.NewSyncProducer(strings.Split("172.25.33.175:9092,172.25.33.176:9092,172.25.33.177:9092", ","), config) //default port if err != nil { log.Println("ERRR") panic(err) } defer func() { if err := producer.Close(); err != nil { panic(err) } }() msg := &sarama.ProducerMessage{ Topic: *topic, Value: sarama.StringEncoder(payload), Partition: int32(*partition), } fmt.Println("XXXX: ", msg.Partition) partition, offset, err := producer.SendMessage(msg) if err != nil { panic(err) } fmt.Println() fmt.Printf("Message is stored in topic(%s)/partition(%d)/offset(%d)\n", *topic, partition, offset) fmt.Println("--------------------------------------------------") fmt.Println(partition) }
The Most Likely Culprit: Your Topic Doesn't Have 12 Partitions
Kafka partitions are zero-indexed—so if your mydata topic has fewer than 12 partitions, the partition number 12 doesn't actually exist. The reason you're not getting an error but seeing messages in partition 0 is probably one of these two scenarios:
- Auto-created topic: You didn't manually create the
mydatatopic, so Kafka auto-created it using the defaultnum.partitionssetting (which is 1 for most default Kafka configurations). - Mistake during topic creation: When you created the topic, you set fewer than 12 partitions (like 1 or 3, matching your 3-node cluster) by accident.
Steps to Fix This
Check your topic's actual partition count
Use Kafka's built-in command-line tool to inspect themydatatopic details:kafka-topics.sh --describe --topic mydata --zookeeper <your-zookeeper-address>Replace
<your-zookeeper-address>with your ZooKeeper endpoint (e.g.,172.25.33.175:2181if it's running on your first Kafka node). Look for thePartitionCountfield in the output—this will tell you how many partitions the topic actually has.Expand the topic's partition count (if needed)
If the partition count is less than 12, you can increase it (note: Kafka doesn't support decreasing partition counts later):kafka-topics.sh --alter --topic mydata --partitions 12 --zookeeper <your-zookeeper-address>Once the partition expansion is done, run your producer code again—it should now correctly send messages to partition 12.
Quick Validation Checks
- You already print
XXXX: 12to confirmmsg.Partitionis set correctly, so that's good. - Take a look at your Kafka broker logs for any warnings about invalid partitions—this can help confirm if the issue was indeed non-existent partitions.
- You already print
A Quick Note on Error Behavior
Normally, specifying a non-existent partition should trigger an InvalidPartition error in Sarama, which would cause your code to panic. The fact that you're not seeing this might be due to compatibility quirks between older Sarama versions and Kafka 1.1.0, but the partition count mismatch is by far the most common reason for this behavior.
内容的提问来源于stack exchange,提问作者Muhammad Sohaib

