openmeter / sink /buffer.go
Leon4gr45's picture
Upload folder using huggingface_hub (part 6)
d6f631f verified
Raw
History Blame Contribute Delete
1.97 kB
package sink
import (
"fmt"
"sync"
"github.com/confluentinc/confluent-kafka-go/v2/kafka"
sinkmodels "github.com/openmeterio/openmeter/openmeter/sink/models"
)
type SinkBuffer struct {
mu sync.Mutex
data map[string]sinkmodels.SinkMessage
}
func NewSinkBuffer() *SinkBuffer {
return &SinkBuffer{
data: map[string]sinkmodels.SinkMessage{},
}
}
func (b *SinkBuffer) Size() int {
b.mu.Lock()
defer b.mu.Unlock()
return len(b.data)
}
func (b *SinkBuffer) Add(message sinkmodels.SinkMessage) {
b.mu.Lock()
defer b.mu.Unlock()
// Unique identifier for each message (topic + partition + offset)
key := message.KafkaMessage.String()
b.data[key] = message
}
type MessageTransformerFunc func(*sinkmodels.SinkMessage)
func (b *SinkBuffer) Dequeue(transformers ...MessageTransformerFunc) []sinkmodels.SinkMessage {
b.mu.Lock()
defer b.mu.Unlock()
messages := make([]sinkmodels.SinkMessage, 0, len(b.data))
for key, message := range b.data {
for _, transformer := range transformers {
transformer(&message)
}
messages = append(messages, message)
delete(b.data, key)
}
return messages
}
// RemoveByPartitions removes messages from the buffer by partitions
// Useful when partitions are revoked.
func (b *SinkBuffer) RemoveByPartitions(partitions []kafka.TopicPartition) {
b.mu.Lock()
defer b.mu.Unlock()
partitionMap := map[string]bool{}
for _, topicPartition := range partitions {
key := topicPartitionKey(topicPartition)
partitionMap[key] = true
}
for key, message := range b.data {
topicKey := topicPartitionKey(message.KafkaMessage.TopicPartition)
if partitionMap[topicKey] {
delete(b.data, key)
}
}
}
func topicPartitionKey(partition kafka.TopicPartition) string {
var topic string
if partition.Topic != nil {
topic = *partition.Topic
}
return partitionKey(topic, partition.Partition)
}
func partitionKey(topic string, partition int32) string {
return fmt.Sprintf("%s-%d", topic, partition)
}