You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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 mydata topic, so Kafka auto-created it using the default num.partitions setting (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

  1. Check your topic's actual partition count
    Use Kafka's built-in command-line tool to inspect the mydata topic 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:2181 if it's running on your first Kafka node). Look for the PartitionCount field in the output—this will tell you how many partitions the topic actually has.

  2. 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.

  3. Quick Validation Checks

    • You already print XXXX: 12 to confirm msg.Partition is 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.

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 06:43:13