如何在Golang中用回调函数异步写入ClickHouse并感知Kafka入数失败?
Golang操作ClickHouse异步插入及失败感知方案
1. Golang中借助回调函数实现ClickHouse异步数据插入
基于官方clickhouse-go/v2客户端,结合Go的goroutine和回调函数可以实现灵活的异步插入逻辑,具体实现如下:
步骤1:初始化ClickHouse客户端
开启客户端的异步插入模式:
import ( "context" "fmt" "time" "github.com/ClickHouse/clickhouse-go/v2" ) func getCHClient() clickhouse.Conn { conn, err := clickhouse.Open(&clickhouse.Options{ Addr: []string{"127.0.0.1:9000"}, Auth: clickhouse.Auth{ Database: "default", Username: "default", Password: "", }, AsyncInsert: true, // 启用异步插入 }) if err != nil { panic(err) } return conn }
步骤2:实现带回调的异步插入函数
将插入逻辑放入goroutine执行,完成后调用回调返回结果:
// 定义回调函数类型 type InsertCallback func(err error) func AsyncInsert(conn clickhouse.Conn, ctx context.Context, query string, args []interface{}, callback InsertCallback) { go func() { // 执行异步插入 err := conn.Exec(ctx, query, args...) // 插入完成后触发回调 callback(err) }() }
步骤3:使用示例
func main() { conn := getCHClient() defer conn.Close() ctx := context.Background() insertQuery := "INSERT INTO user (id, name) VALUES (?, ?)" insertArgs := []interface{}{1001, "alice"} // 调用带回调的异步插入 AsyncInsert(conn, ctx, insertQuery, insertArgs, func(err error) { if err != nil { fmt.Printf("异步插入失败: %v\n", err) // 此处可添加重试、日志上报、死信队列投递等逻辑 } else { fmt.Println("异步插入成功") } }) // 主goroutine需等待异步任务完成,实际业务中可使用sync.WaitGroup替代sleep time.Sleep(2 * time.Second) }
2. Kafka数据写入ClickHouse时,ACK导致异步插入失败的感知方案
针对ACK引发的异步插入失败,可通过以下方式感知并处理:
- 客户端回调直接捕获ACK错误
在插入回调中判断错误类型,针对性处理ACK失败场景:
AsyncInsert(conn, ctx, insertQuery, insertArgs, func(err error) { if err != nil { // 匹配ClickHouse客户端定义的ACK未确认错误 if errors.Is(err, clickhouse.ErrAsyncInsertNotAcked) { fmt.Println("异步插入未收到ClickHouse ACK,操作失败") // 将当前Kafka消息转存死信队列,或加入本地重试队列 } else { fmt.Printf("插入失败: %v\n", err) } } })
- 强制等待ACK响应
修改客户端配置,开启WaitForAcknowledge,此时异步插入会阻塞等待ACK,错误直接通过Exec方法返回:
conn, err := clickhouse.Open(&clickhouse.Options{ // ...其他配置 AsyncInsert: true, WaitForAcknowledge: true, // 等待ClickHouse返回ACK })
控制Kafka消息偏移量提交时机
自定义Kafka消费者逻辑时,不要立即提交偏移量:- 插入成功(收到ACK)后,再提交当前消息偏移量
- 插入失败(ACK超时/失败)时,不提交偏移量,让消费者重新拉取该消息;或标记消息为失败后提交偏移量,转存到死信队列
监控ClickHouse日志
在ClickHouse配置中设置async_insert_log_level=info,通过监控日志中的Failed to acknowledge async insert等关键字,实现告警或批量失败排查。
内容的提问来源于stack exchange,提问作者CharmCcc
相关产品推荐
相关产品推荐

