在使用 Golang 操作 Kafka 时,你可以使用 Sarama 库来设置消息的失效时间。以下是一个示例代码,演示如何在生产者端设置数据失效时间:
package mainimport ("log""time""github.com/Shopify/sarama"
)func main() {// Kafka broker地址brokers := []string{"localhost:9092"}// 创建配置config := sarama.NewConfig()// 设置消息的失效时间expirationTime := time.Hour * 24 // 一天的时间config.Message.MaxAge = expirationTime// 创建生产者producer, err := sarama.NewSyncProducer(brokers, config)if err != nil {log.Fatal("Failed to create producer:", err)}defer producer.Close()// 定义消息message := &sarama.ProducerMessage{Topic: "your_topic",Value: sarama.StringEncoder("Hello, Kafka!"),}// 发送消息partition, offset, err := producer.SendMessage(message)if err != nil {log.Println("Failed to send message:", err)} else {log.Printf("Message sent successfully! Partition:%d Offset:%d\n", partition, offset)}
}
上述示例中,我们首先创建了一个 sarama.Config 实例,并通过 config.Message.MaxAge 属性设置了消息的失效时间,此处设定为一天 (time.Hour * 24)。然后,我们创建了一个生产者实例并发送一条消息。
除了设置消息的失效时间,还可以在消费者端进行相关处理。可以使用 sarama.Consumer 接口提供的方法,结合 Message.Timestamp 属性来判断消息是否过期,并根据需要进行处理。