kafka

package
v0.0.25 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jan 9, 2023 License: Apache-2.0 Imports: 11 Imported by: 0

Documentation

Index

Constants

View Source
const (
	TopicCommentPublish      = "comment_publish"
	TopicCommentCacheRebuild = "comment_cache_rebuild"
	TopicCommentOperator     = "comment_operator"

	TopicRelationFollow       = "relation_follow"
	TopicRelationCacheRebuild = "relation_cache_rebuild"
	TopicRelationOperator     = "relation_operator"

	TopicOpusOperator = "opus_operator"

	EventTypeCreate        = 1
	EventTypeReply         = 2
	EventTypeListMissed    = 3
	EventTypeSubListMissed = 4
	EventTypeLike          = 5
	EventTypeHate          = 6
	EventTypeDelete        = 7
)

Variables

This section is empty.

Functions

This section is empty.

Types

type Consumer

type Consumer struct {
	// contains filtered or unexported fields
}

func InitConsumer

func InitConsumer(broker, topics []string, group string, fn func(context.Context, *KafkaMessage), apiLogger, excLogger *log.Logger) (*Consumer, error)

func (*Consumer) Start

func (consumer *Consumer) Start()

func (*Consumer) Stop

func (consumer *Consumer) Stop() error

type KafkaMessage

type KafkaMessage struct {
	TraceId              string   `protobuf:"bytes,1,opt,name=trace_id,json=traceId" json:"trace_id,omitempty"`
	EventType            int32    `protobuf:"varint,2,opt,name=event_type,json=eventType" json:"event_type,omitempty"`
	Message              []byte   `protobuf:"bytes,3,opt,name=message,proto3" json:"message,omitempty"`
	XXX_NoUnkeyedLiteral struct{} `json:"-"`
	XXX_unrecognized     []byte   `json:"-"`
	XXX_sizecache        int32    `json:"-"`
}

func (*KafkaMessage) Descriptor

func (*KafkaMessage) Descriptor() ([]byte, []int)

func (*KafkaMessage) GetEventType

func (m *KafkaMessage) GetEventType() int32

func (*KafkaMessage) GetMessage

func (m *KafkaMessage) GetMessage() []byte

func (*KafkaMessage) GetTraceId

func (m *KafkaMessage) GetTraceId() string

func (*KafkaMessage) ProtoMessage

func (*KafkaMessage) ProtoMessage()

func (*KafkaMessage) Reset

func (m *KafkaMessage) Reset()

func (*KafkaMessage) String

func (m *KafkaMessage) String() string

func (*KafkaMessage) XXX_DiscardUnknown

func (m *KafkaMessage) XXX_DiscardUnknown()

func (*KafkaMessage) XXX_Marshal

func (m *KafkaMessage) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*KafkaMessage) XXX_Merge

func (dst *KafkaMessage) XXX_Merge(src proto.Message)

func (*KafkaMessage) XXX_Size

func (m *KafkaMessage) XXX_Size() int

func (*KafkaMessage) XXX_Unmarshal

func (m *KafkaMessage) XXX_Unmarshal(b []byte) error

type Producer

type Producer struct {
	// contains filtered or unexported fields
}

func InitKafkaProducer

func InitKafkaProducer(addr []string, apiLogger, excLogger *log.Logger) (*Producer, error)

func (*Producer) SendKafkaMsg added in v0.0.25

func (producer *Producer) SendKafkaMsg(ctx context.Context, topic, key string, req proto.Message, eventType int32) error

func (*Producer) Stop

func (producer *Producer) Stop() error

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL