浏览代码

Add version check for headers

Evan Huus 7 年之前
父节点
当前提交
6d9cc201b5
共有 1 个文件被更改,包括 3 次插入0 次删除
  1. 3 0
      async_producer.go

+ 3 - 0
async_producer.go

@@ -271,6 +271,9 @@ func (p *asyncProducer) dispatcher() {
 		version := 1
 		if p.conf.Version.IsAtLeast(V0_11_0_0) {
 			version = 2
+		} else if msg.Headers != nil {
+			p.returnError(msg, ConfigurationError("Producing headers requires Kafka at least v0.11"))
+			continue
 		}
 		if msg.byteSize(version) > p.conf.Producer.MaxMessageBytes {
 			p.returnError(msg, ErrMessageSizeTooLarge)