Kafka到ClickHouse的IPFIX数据导入延迟过高,求优化方案
我有大量IPFIX(NetFlow)记录存入Kafka,已用Go开发消费者程序将数据写入ClickHouse。目前程序每5分钟仅能插入200万条记录,但Kafka每5分钟流入记录超2000万条,导致两者间延迟已超过20小时,寻求可行的提速方案。
Go消费者代码
import ( "context" "database/sql" "encoding/json" "flag" "fmt" "log" // "os" // "strconv" "sync" "time" "github.com/ClickHouse/clickhouse-go" "github.com/segmentio/kafka-go" cluster "github.com/bsm/sarama-cluster" ) type options struct { Broker string Topic string Debug bool Workers int } type dataField struct { I int V interface{} } type Header struct { Version int Length int ExportTime int64 SequenceNo int DomainID int } type ipfix struct { AgentID string Header Header DataSets [][]dataField } type dIPFIXSample struct { device string sourceIPv4Address string sourceTransportPort uint64 postNATSourceIPv4Address string postNATSourceTransportPort uint64 destinationIPv4Address string postNATDestinationIPv4Address string postNATDestinationTransportPort uint64 dstport uint64 timestamp string postNATSourceIPv6Address string postNATDestinationIPv6Address string sourceIPv6Address string destinationIPv6Address string proto uint8 login string sessionid uint64 } var opts options func init() { flag.StringVar(&opts.Broker, "broker", "172.18.0.4:9092", "broker ipaddress:port") flag.StringVar(&opts.Topic, "topic", "vflow.ipfix", "kafka topic") flag.BoolVar(&opts.Debug, "debug", true, "enabled/disabled debug") flag.IntVar(&opts.Workers, "workers", 16, "workers number / partition number") flag.Parse() } func main() { var ( wg sync.WaitGroup ch = make(chan ipfix, 10000) ) for i := 0; i < 5; i++ { go ingestClickHouse(ch) } wg.Add(opts.Workers) for i := 0; i < opts.Workers; i++ { go func(ti int) { // create a new kafka reader with the broker and topic r := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{opts.Broker}, Topic: opts.Topic, GroupID: "mygroup", // start consuming from the earliest message StartOffset: 0, }) pCount := 0 count := 0 tik := time.Tick(10 * time.Second) for { select { case <-tik: if opts.Debug { log.Printf("partition GroupId#%d, rate=%d\n", ti, (count-pCount)/10) } pCount = count default: // read the next message from kafka m, err := r.ReadMessage(context.Background()) if err != nil { if err == kafka.ErrGenerationEnded { log.Println("generation ended") return } log.Println(err) continue } // log.Printf("Received message from Kafka: %s\n", string(m.Value)) // unmarshal the message into an ipfix struct objmap:= ipfix{} if err := json.Unmarshal(m.Value, &objmap); err != nil { log.Println(err) continue } fmt.Sprintf("kkkkkkkkkkkkkkkk%v",objmap); // send the ipfix struct to the ingestClickHouse goroutine ch <- objmap // go ingestClickHouse(ch) // mark the message as processed if err := r.CommitMessages(context.Background(), m); err != nil { log.Println(err) continue } count++ } } }(i) } wg.Wait() // close(ch) } func ingestClickHouse(ch chan ipfix) { var sample ipfix connect, err := sql.Open("clickhouse", "tcp://127.0.0.1:9000?debug=true&username=default&password=wawa123") if err != nil { log.Fatal(err) } if err := connect.Ping(); err != nil { if exception, ok := err.(*clickhouse.Exception); ok { log.Printf("[%d] %s \n%s\n", exception.Code, exception.Message, exception.StackTrace) } else { log.Println(err) } return } defer connect.Close() for { tx, err := connect.Begin() if err != nil { log.Fatal(err) } stmt, err := tx.Prepare("INSERT INTO natdb.natlogs (timestamp,router_ip,sourceIPv4Address, sourceTransportPort,postNATSourceIPv4Address,postNATSourceTransportPort,destinationIPv4Address,dstport,postNATDestinationIPv4Address, postNATDestinationTransportPort,postNATSourceIPv6Address,postNATDestinationIPv6Address,sourceIPv6Address,destinationIPv6Address,proto,login) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?,?)") if err != nil { log.Fatal(err) } for i := 0; i < 10000; i++ { sample = <-ch for _, data := range sample.DataSets { s := dIPFIXSample{} for _, dd := range data { switch dd.I { case 8: s.sourceIPv4Address = dd.V.(string) case 7: s.sourceTransportPort =uint64( dd.V.(float64)) case 225: s.postNATSourceIPv4Address = dd.V.(string) case 227: s.postNATSourceTransportPort = uint64(dd.V.(float64)) case 12: s.destinationIPv4Address=dd.V.(string) case 11: s.dstport=uint64(dd.V.(float64)) case 226: s.postNATDestinationIPv4Address=dd.V.(string) case 27: s.sourceIPv6Address=dd.V.(string) case 28: s.destinationIPv6Address=dd.V.(string) case 281: s.postNATSourceIPv6Address=dd.V.(string) case 282: s.postNATDestinationIPv6Address=dd.V.(string) case 2003: s.login =dd.V.(string) log.Printf(dd.V.(string)) case 228: s.postNATDestinationTransportPort=uint64(dd.V.(float64)) case 4: s.proto = uint8(dd.V.(float64)) } } timestamp := time.Unix(sample.Header.ExportTime, 0).Format("2006-01-02 15:04:05") if _, err := stmt.Exec( timestamp, sample.AgentID, s.sourceIPv4Address, s.sourceTransportPort, s.postNATSourceIPv4Address, s.postNATSourceTransportPort, s.destinationIPv4Address, s.dstport, s.postNATDestinationIPv4Address, s.postNATDestinationTransportPort, s.postNATSourceIPv6Address, s.postNATDestinationIPv6Address, s.sourceIPv6Address, s.destinationIPv6Address, s.proto, s.login, ); err != nil { log.Fatal(err) } } } go func(tx *sql.Tx) { if err := tx.Commit(); err != nil { log.Fatal(err) } }(tx) } }
ClickHouse表结构
CREATE TABLE natdb.natlogs ( `timestamp` DateTime, `router_ip` String, `sourceIPv4Address` String, `sourceTransportPort` UInt64, `postNATSourceIPv4Address` String, `postNATSourceTransportPort` UInt64, `destinationIPv4Address` String, `dstport` UInt64, `postNATDestinationIPv4Address` String, `postNATDestinationTransportPort` UInt64, `proto` UInt8, `login` String, `sessionid` String, `sourceIPv6Address` String, `destinationIPv6Address` String, `postNATSourceIPv6Address` String, `postNATDestinationIPv6Address` String, INDEX idx_natlogs_router_source_time_postnat (router_ip, sourceIPv4Address, timestamp, postNATSourceIPv4Address) TYPE minmax GRANULARITY 1 ) ENGINE = MergeTree PARTITION BY toYYYYMMDD(timestamp) ORDER BY router_ip SETTINGS index_granularity = 8192
优化方案
一、Go消费者代码优化
改用批量插入替代单条插入
ClickHouse对批量插入优化极强,单条stmt.Exec性能极差。修改代码,先收集一批数据(比如10000条),然后用INSERT ... VALUES (...), (...), ...的方式一次性插入,或者使用ClickHouse的PrepareBatch方法(clickhouse-go支持),大幅减少插入请求次数。移除事务,关闭ClickHouse debug模式
MergeTree引擎不需要事务支持,当前代码中的事务不仅不会提升可靠性,还会增加额外开销。同时连接字符串中的debug=true要关闭,减少日志和调试开销。优化Kafka消费方式
用kafka-go的ReadBatch方法批量读取Kafka消息,一次读取多条(比如1000条),减少网络IO和消息处理的循环开销。同时调整消费者的MaxBytes参数,允许更大的批量读取。优化JSON反序列化与类型断言
- 当前
dataField.V用interface{}存储,每次类型断言有性能损耗,自定义JSON反序列化逻辑,直接将字段解析为对应类型(比如字符串、uint64等),避免interface{}。 - 若生产者可控,改用Protobuf序列化替代JSON,Protobuf的序列化/反序列化速度远快于JSON,还能减少消息体积。
- 当前
调整连接池参数
当前每个ingestClickHousegoroutine创建一个独立连接,改为使用连接池:connect.SetMaxOpenConns(20) connect.SetMaxIdleConns(10) connect.SetConnMaxLifetime(30 * time.Minute)合理设置连接数,避免连接过多或不足导致的性能瓶颈。
移除不必要的日志
代码中log.Printf(dd.V.(string))在循环内执行,会极大拖慢性能,生产环境必须关闭Debug日志,移除这类高频日志输出。调整通道容量
当前通道ch容量为10000,当Kafka消费速度快于ClickHouse写入时,会阻塞消费者,建议调大到100000或更高,缓冲更多待处理数据。
二、ClickHouse表结构与配置优化
优化ORDER BY与分区键
- 将
ORDER BY router_ip改为ORDER BY (router_ip, timestamp),时间序列数据按时间排序能减少Merge操作的开销,提升写入和查询性能。 - 若每日数据量极大,将分区键从
toYYYYMMDD(timestamp)改为toYYYYMMDDhh(timestamp),按小时分区,避免单个分区过大导致Merge缓慢。
- 将
优化数据类型
将IPv4/IPv6地址从String改为ClickHouse原生的IPv4/IPv6类型,减少存储空间(每个IPv4仅占4字节,远小于字符串),降低IO开销:`sourceIPv4Address` IPv4, `postNATSourceIPv4Address` IPv4, `destinationIPv4Address` IPv4, `postNATDestinationIPv4Address` IPv4, `sourceIPv6Address` IPv6, `postNATSourceIPv6Address` IPv6, `destinationIPv6Address` IPv6, `postNATDestinationIPv6Address` IPv6调整索引设置
- 当前
minmax索引的GRANULARITY 1过于精细,会大幅增加索引存储和写入开销,改为GRANULARITY 8192(默认值)或更高。 - 若该
minmax索引不是业务必需,直接移除,减少写入时的索引维护工作。
- 当前
调整MergeTree引擎参数
修改表的SETTINGS,优化写入性能:ALTER TABLE natdb.natlogs MODIFY SETTINGS max_insert_block_size = 1048576, -- 增大插入块大小 index_granularity = 65536, -- 增大索引粒度 merge_with_ttl_timeout = 3600, -- 降低Merge频率 write_buffer_size = 134217728 -- 增大写入缓冲区使用分布式表并行写入
若ClickHouse是集群部署,创建分布式表,将数据并行写入多个节点,利用集群的横向扩展能力提升写入吞吐量。
内容的提问来源于stack exchange,提问作者Muhammed Elbuvaydani

