Apache Kafka是分布式流处理平台,以高吞吐、低延迟和水平可扩展性著称,单集群可支撑每秒百万级消息吞吐。在后端微服务架构中,Kafka作为消息中间件承担服务间异步通信、事件驱动和日志聚合等角色,是高并发设计和业务解耦的核心基础设施。
Kafka核心架构与分区并行机制
Kafka的架构由Producer、Consumer、Broker、Topic和Partition组成。Topic是逻辑消息分类,Partition是物理存储单元——每个Partition是一个有序的、不可变的消息序列。分区数量直接决定并行度:Consumer Group中的消费者数不能超过分区数,超出部分的消费者处于空闲状态。
消息写入Partition时根据Key哈希路由到固定分区,保证同一Key的消息顺序。无Key消息采用轮询策略。每个Partition有多个副本(Replica),其中Leader副本处理读写请求,Follower副本同步数据。ISR(In-Sync Replicas)集合中的副本与Leader保持同步,当Leader故障时从ISR中选举新Leader。
关键配置参数说明:
# server.properties - Broker核心配置
broker.id=1
listeners=PLAINTEXT://:9092
log.dirs=/data/kafka/logs
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=1048576
socket.receive.buffer.bytes=1048576
log.retention.hours=168
log.retention.bytes=10737418240
log.segment.bytes=1073741824
log.cleanup.policy=delete
default.replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false
num.partitions=12
Go语言Kafka生产者实现与可靠性保证
使用segmentio/kafka-go库实现高性能生产者。该库纯Go实现,无CGO依赖,适合容器化部署。生产者可靠性通过acks参数和重试机制保证。
package main
import (
"context"
"encoding/json"
"log"
"time"
"github.com/segmentio/kafka-go"
)
type OrderEvent struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
Amount float64 `json:"amount"`
Status string `json:"status"`
Timestamp time.Time `json:"timestamp"`
}
type KafkaProducer struct {
writer *kafka.Writer
}
func NewKafkaProducer(brokers []string, topic string) *KafkaProducer {
return &KafkaProducer{
writer: &kafka.Writer{
Addr: kafka.TCP(brokers...),
Topic: topic,
Balancer: &kafka.Hash{},
RequiredAcks: kafka.RequireAll,
Async: false,
BatchSize: 100,
BatchTimeout: 10 * time.Millisecond,
Compression: kafka.Snappy,
MaxAttempts: 5,
ReadTimeout: 10 * time.Second,
WriteTimeout: 10 * time.Second,
},
}
}
func (p *KafkaProducer) SendMessage(ctx context.Context, key string, event OrderEvent) error {
value, err := json.Marshal(event)
if err != nil {
return err
}
err = p.writer.WriteMessages(ctx, kafka.Message{
Key: []byte(key),
Value: value,
Headers: []kafka.Header{
{Key: "event-type", Value: []byte("order-created")},
{Key: "source", Value: []byte("order-service")},
{Key: "timestamp", Value: []byte(time.Now().Format(time.RFC3339))},
},
})
if err != nil {
log.Printf("Kafka send failed: %v", err)
return err
}
return nil
}
func (p *KafkaProducer) SendBatch(ctx context.Context, events []OrderEvent) error {
messages := make([]kafka.Message, len(events))
for i, event := range events {
value, _ := json.Marshal(event)
messages[i] = kafka.Message{
Key: []byte(event.OrderID),
Value: value,
}
}
return p.writer.WriteMessages(ctx, messages...)
}
func (p *KafkaProducer) Close() error {
return p.writer.Close()
}
消费者组与精确一次消费实现
消费者组(Consumer Group)机制使Kafka支持点对点和发布订阅两种模式。同一组内的消费者均衡消费分区,不同组独立消费全量消息。手动提交Offset实现精确一次消费(Exactly-Once),避免自动提交导致的重复消费问题。
package main
import (
"context"
"encoding/json"
"log"
"time"
"github.com/segmentio/kafka-go"
)
type KafkaConsumer struct {
reader *kafka.Reader
}
func NewKafkaConsumer(brokers []string, topic, groupID string) *KafkaConsumer {
return &KafkaConsumer{
reader: kafka.NewReader(kafka.ReaderConfig{
Brokers: brokers,
Topic: topic,
GroupID: groupID,
MinBytes: 10,
MaxBytes: 10 * 1024 * 1024,
MaxWait: 2 * time.Second,
CommitInterval: 0,
ReadLagInterval: -1,
}),
}
}
func (c *KafkaConsumer) Consume(ctx context.Context, handler func(OrderEvent) error) error {
for {
select {
case <-ctx.Done():
return ctx.Err()
default:
msg, err := c.reader.ReadMessage(ctx)
if err != nil {
log.Printf("Read message error: %v", err)
continue
}
var event OrderEvent
if err := json.Unmarshal(msg.Value, &event); err != nil {
log.Printf("Unmarshal error: %v, partition=%d, offset=%d",
err, msg.Partition, msg.Offset)
continue
}
if err := handler(event); err != nil {
log.Printf("Handler error: %v, will retry", err)
continue
}
if err := c.reader.CommitMessages(ctx, msg); err != nil {
log.Printf("Commit offset failed: %v", err)
}
log.Printf("Processed: partition=%d, offset=%d, order=%s",
msg.Partition, msg.Offset, event.OrderID)
}
}
}
func (c *KafkaConsumer) Close() error {
return c.reader.Close()
}
消费幂等与死信队列设计
Kafka的At-Least-Once语义意味着消费者可能收到重复消息。实现幂等消费的标准做法是在业务层使用唯一ID去重:
type OrderHandler struct {
redis *redis.Client
db *sql.DB
}
func (h *OrderHandler) Handle(event OrderEvent) error {
key := fmt.Sprintf("kafka:processed:%s", event.OrderID)
set, err := h.redis.SetNX(context.Background(), key, "1", 24*time.Hour).Result()
if err != nil {
return err
}
if !set {
return nil
}
_, err = h.db.ExecContext(context.Background(),
`INSERT INTO orders (id, user_id, amount, status, created_at)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (id) DO NOTHING`,
event.OrderID, event.UserID, event.Amount, event.Status, event.Timestamp)
return err
}
type DeadLetterQueue struct {
producer *kafka.Writer
}
func (dlq *DeadLetterQueue) Send(ctx context.Context, originalMsg kafka.Message, reason string) error {
headers := append(originalMsg.Headers,
kafka.Header{Key: "dlq-reason", Value: []byte(reason)},
kafka.Header{Key: "dlq-timestamp", Value: []byte(time.Now().Format(time.RFC3339))},
)
return dlq.producer.WriteMessages(ctx, kafka.Message{
Key: originalMsg.Key,
Value: originalMsg.Value,
Headers: headers,
})
}
生产环境中Kafka的吞吐量调优需要关注几个关键指标:Producer的batch.size和linger.ms控制批量发送粒度,Consumer的fetch.min.bytes和fetch.max.wait.ms控制拉取效率。合理设置分区数(通常为消费者数的2-3倍)可以在保证并行度的同时避免过度分区带来的元数据管理开销。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/kafka-gao-tun-tu-xiao-xi-dui-lie-jia-gou-she-ji-yu-go-yu/