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

Kafka到ClickHouse的IPFIX数据导入延迟过高,求优化方案

Kafka到ClickHouse写入性能优化求助

我有大量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消费者代码优化

  1. 改用批量插入替代单条插入
    ClickHouse对批量插入优化极强,单条stmt.Exec性能极差。修改代码,先收集一批数据(比如10000条),然后用INSERT ... VALUES (...), (...), ...的方式一次性插入,或者使用ClickHouse的PrepareBatch方法(clickhouse-go支持),大幅减少插入请求次数。

  2. 移除事务,关闭ClickHouse debug模式
    MergeTree引擎不需要事务支持,当前代码中的事务不仅不会提升可靠性,还会增加额外开销。同时连接字符串中的debug=true要关闭,减少日志和调试开销。

  3. 优化Kafka消费方式
    用kafka-go的ReadBatch方法批量读取Kafka消息,一次读取多条(比如1000条),减少网络IO和消息处理的循环开销。同时调整消费者的MaxBytes参数,允许更大的批量读取。

  4. 优化JSON反序列化与类型断言

    • 当前dataField.V用interface{}存储,每次类型断言有性能损耗,自定义JSON反序列化逻辑,直接将字段解析为对应类型(比如字符串、uint64等),避免interface{}。
    • 若生产者可控,改用Protobuf序列化替代JSON,Protobuf的序列化/反序列化速度远快于JSON,还能减少消息体积。
  5. 调整连接池参数
    当前每个ingestClickHouse goroutine创建一个独立连接,改为使用连接池:

    connect.SetMaxOpenConns(20)
    connect.SetMaxIdleConns(10)
    connect.SetConnMaxLifetime(30 * time.Minute)
    

    合理设置连接数,避免连接过多或不足导致的性能瓶颈。

  6. 移除不必要的日志
    代码中log.Printf(dd.V.(string))在循环内执行,会极大拖慢性能,生产环境必须关闭Debug日志,移除这类高频日志输出。

  7. 调整通道容量
    当前通道ch容量为10000,当Kafka消费速度快于ClickHouse写入时,会阻塞消费者,建议调大到100000或更高,缓冲更多待处理数据。

二、ClickHouse表结构与配置优化

  1. 优化ORDER BY与分区键

    • 将ORDER BY router_ip改为ORDER BY (router_ip, timestamp),时间序列数据按时间排序能减少Merge操作的开销,提升写入和查询性能。
    • 若每日数据量极大,将分区键从toYYYYMMDD(timestamp)改为toYYYYMMDDhh(timestamp),按小时分区,避免单个分区过大导致Merge缓慢。
  2. 优化数据类型
    将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
    
  3. 调整索引设置

    • 当前minmax索引的GRANULARITY 1过于精细,会大幅增加索引存储和写入开销,改为GRANULARITY 8192(默认值)或更高。
    • 若该minmax索引不是业务必需,直接移除,减少写入时的索引维护工作。
  4. 调整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 -- 增大写入缓冲区
    
  5. 使用分布式表并行写入
    若ClickHouse是集群部署,创建分布式表,将数据并行写入多个节点,利用集群的横向扩展能力提升写入吞吐量。


内容的提问来源于stack exchange,提问作者Muhammed Elbuvaydani

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 17:29:54