GO的Kafka的问题,Local: Queue full?
小伙伴们对Golang编程感兴趣吗?是否正在学习相关知识点?如果是,那么本文《GO的Kafka的问题,Local: Queue full?》,就很适合你,本篇文章讲解的知识点主要包括。在之后的文章中也会多多分享相关知识点,希望对大家的知识积累有所帮助!
问题内容
今天线上报错Local: Queue full,导致接口无法使用,请求各位大佬指教
目前使用go连接kafka的库使用的是
package kfk
import (
"fmt"
"github.com/confluentinc/confluent-kafka-go/v2/kafka"
"strings"
"time"
)
// 发送消息到指定主题
func SendMessage(broker, topic string, tableName string, message []byte) error {
config := &kafka.ConfigMap{
"bootstrap.servers": strings.Join([]string{"localhost:9092"}, ","),
"acks": "1",
"delivery.timeout.ms": 3000,
"security.protocol": "PLAINTEXT",
}
producer, err := kafka.NewProducer(config)
if err != nil {
}
err = producer.Produce(&kafka.Message{
TopicPartition: kafka.TopicPartition{
Topic: &topic,
},
Key: []byte(tableName),
Value: message,
Timestamp: time.Now(),
}, nil)
if err != nil {
fmt.Println(err)
return err
}
return nil
}经过多次的测试,不管是不是进行消费,都是在110万条左右开始报错,并且不在之后的数据不能写入,但是重启服务就好了,分析原因:github.com/confluentinc/confluent-kafka-go 这个包使用了队列的概念,累计到110万时,超出了队列的缓存区,所以不能再继续写入,而重启服务则清空了这个缓冲区,让服务可用,但是再累计到110万时,服务还会不可用。
各位大佬,这是这个库的Bug么?是作者特意做的,还是我的配置不正确,可以通过什么配置可以解决呢?
有没有更好一点的kafka库呢? 其他库会有这种问题吗?
正确答案
这个库其实是包装了一下 c 的实现,最终其实这个报错应该也是 c 的库里返回的,并不是这个库本身的意思 详情见:https://github.com/confluentinc/librdkafka/blob/aa50e52a12bece3e399f69b8477fd0c8aadbfff1/src/rdkafka.c#L430
我就不追 C 的源码了,大胆来猜测一下,估计这个报错就是因为本地队列满了导致的(kafka 客户端在发送消息的时候,并不是收到之后马上就发送出去的,而是攒起来,一批一批发)。这样的话应该会有两种思路,一种是库里面本身支持了某个配置可以修改本地队列最大数量,或者 buffer 最大数量类似的参数;还有就是本身不提供这样的参数调整,通过前面的限流完成。
顺着这个思路去找一下文档,发现有
https://github.com/confluentinc/librdkafka/blob/master/INTRODUCTION.md
Compression Producer message compression is enabled through the compression.codec configuration property. Compression is performed on the batch of messages in the local queue, the larger the batch the higher likelyhood of a higher compression ratio. The local batch queue size is controlled through the batch.num.messages, batch.size, and linger.ms configuration properties as described in the High throughput chapter above.
估计就是 batch.num.messages 这个配置了,这个库没用过,所以不知道包装之后有没有这个配置,你可以找一下。
理论要掌握,实操不能落!以上关于《GO的Kafka的问题,Local: Queue full?》的详细介绍,大家都掌握了吧!如果想要继续提升自己的能力,那么就来关注golang学习网公众号吧!
基于go-zero进行微服务架构模式的分析与应用
- 上一篇
- 基于go-zero进行微服务架构模式的分析与应用
- 下一篇
- Gin框架的分布式锁和分布式事务详解
-
- Golang · Go问答 | 27分钟前 | 数据隔离 循环变量 Go测试 t.Parallel 并行子测试
- 并行子测试为什么会拿到同一个循环变量,应该怎样隔离数据
- 205浏览 收藏
-
- Golang · Go问答 | 1小时前 | Context · 并发编程 · 接口设计 · Go问答 · 生命周期 结构体 context.Context 向后兼容 Go context 取消传播
- 为什么不建议把 Context 保存进结构体,例外场景是什么
- 227浏览 收藏
-
- Golang · Go问答 | 2小时前 | golang · Context · 并发编程 · 超时控制 WithTimeout WithCancel Go context 取消传播 WithoutCancel
- WithCancel、WithTimeout 与 WithoutCancel 的边界怎么选
- 202浏览 收藏
-
- Golang · Go问答 | 2小时前 | goroutine · Context · 并发编程 · 故障排查 · Go问答 · channel WaitGroup Go context ctx.Done 阻塞排查 context.Canceled
- Context 已取消但函数仍不退出,通常漏查了哪些阻塞点
- 144浏览 收藏
-
- Golang · Go问答 | 3小时前 |
- 无缓冲和有缓冲 Channel 的选择应看吞吐还是同步语义
- 421浏览 收藏
-
- Golang · Go问答 | 3小时前 | channel · panic · 并发编程 · Go问答 · 并发安全 go channel关闭 send on closed channel Golang panic Channel关闭权 多生产者
- 向已关闭 Channel 发送为什么会 panic,关闭权应归谁
- 156浏览 收藏
-
- Golang · Go问答 | 4小时前 | 并发 · channel · goroutine · go · Context · context 并发限制 工作池 Go channel worker pool Goroutine生命周期
- 任务数很多时应该每任务一个协程还是固定工作池
- 458浏览 收藏
-
- 前端进阶之JavaScript设计模式
- 设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
- 543次学习
-
- GO语言核心编程课程
- 本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
- 516次学习
-
- 简单聊聊mysql8与网络通信
- 如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
- 500次学习
-
- JavaScript正则表达式基础与实战
- 在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
- 487次学习
-
- 从零制作响应式网站—Grid布局
- 本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
- 485次学习
-
- PubMedQA
- 深入了解PubMedQA生物医学问答数据集,涵盖其核心功能、使用方法及在临床决策、药物研发等场景的应用,助力提升NLP模型性能。
- 363次使用
-
- H2O EvalGPT
- H2O EvalGPT是H2O.ai推出的开源LLM评估平台,提供详细的大模型性能排行榜、行业特定基准测试及A/B测试功能,助您快速选择最适合项目的高性能大语言模型。
- 417次使用
-
- LMArena
- LMArena是加州大学伯克利分校推出的AI模型匿名评测平台。通过盲测投票机制,用户可对比不同大模型回答并生成实时排行榜,助力开发者优化模型及用户选择最佳AI工具。
- 430次使用
-
- HELM
- 深入了解斯坦福推出的HELM(Holistic Evaluation of Language Models)大模型评测体系。本文解析其核心功能、安装配置步骤及应用场景,涵盖准确性、公平性、鲁棒性等多维度指标,助力开发者全面优化语言模型性能。
- 385次使用
-
- MMBench
- MMBench是由上海人工智能实验室等机构联合推出的多模态基准测试平台,提供细粒度能力评估、大规模数据集及VLMEvalKit工具。本文详细介绍其核心功能、安装使用方法及应用场景,助力开发者全面评估多模态模型性能。
- 210次使用
-
- 用Nginx反向代理部署go写的网站。
- 2023-01-17 502浏览
-
- GoLand调式动态执行代码
- 2023-01-13 502浏览
-
- Go crypto/rand.Text 的长度为什么不是固定字符数
- 2026-10-04 501浏览
-
- Go strings.ToValidUTF8 清洗日志内容的边界
- 2026-10-03 501浏览
-
- Go tls.GetCertificate 为什么收不到空 ServerName 请求
- 2026-09-27 501浏览

