kinesisacquisition

package
v1.4.4-rc1 Latest Latest
Warning

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

Go to latest
Published: Dec 12, 2022 License: MIT Imports: 19 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type CloudWatchSubscriptionRecord

type CloudWatchSubscriptionRecord struct {
	MessageType         string                           `json:"messageType"`
	Owner               string                           `json:"owner"`
	LogGroup            string                           `json:"logGroup"`
	LogStream           string                           `json:"logStream"`
	SubscriptionFilters []string                         `json:"subscriptionFilters"`
	LogEvents           []CloudwatchSubscriptionLogEvent `json:"logEvents"`
}

type CloudwatchSubscriptionLogEvent

type CloudwatchSubscriptionLogEvent struct {
	ID        string `json:"id"`
	Message   string `json:"message"`
	Timestamp int64  `json:"timestamp"`
}

type KinesisConfiguration

type KinesisConfiguration struct {
	configuration.DataSourceCommonCfg `yaml:",inline"`
	StreamName                        string  `yaml:"stream_name"`
	StreamARN                         string  `yaml:"stream_arn"`
	UseEnhancedFanOut                 bool    `yaml:"use_enhanced_fanout"` //Use RegisterStreamConsumer and SubscribeToShard instead of GetRecords
	AwsProfile                        *string `yaml:"aws_profile"`
	AwsRegion                         string  `yaml:"aws_region"`
	AwsEndpoint                       string  `yaml:"aws_endpoint"`
	ConsumerName                      string  `yaml:"consumer_name"`
	FromSubscription                  bool    `yaml:"from_subscription"`
	MaxRetries                        int     `yaml:"max_retries"`
}

type KinesisSource

type KinesisSource struct {
	Config KinesisConfiguration
	// contains filtered or unexported fields
}

func (*KinesisSource) CanRun

func (k *KinesisSource) CanRun() error

func (*KinesisSource) Configure

func (k *KinesisSource) Configure(yamlConfig []byte, logger *log.Entry) error

func (*KinesisSource) ConfigureByDSN

func (k *KinesisSource) ConfigureByDSN(string, map[string]string, *log.Entry) error

func (*KinesisSource) DeregisterConsumer

func (k *KinesisSource) DeregisterConsumer() error

func (*KinesisSource) Dump

func (k *KinesisSource) Dump() interface{}

func (*KinesisSource) EnhancedRead

func (k *KinesisSource) EnhancedRead(out chan types.Event, t *tomb.Tomb) error

func (*KinesisSource) GetAggregMetrics

func (k *KinesisSource) GetAggregMetrics() []prometheus.Collector

func (*KinesisSource) GetMetrics

func (k *KinesisSource) GetMetrics() []prometheus.Collector

func (*KinesisSource) GetMode

func (k *KinesisSource) GetMode() string

func (*KinesisSource) GetName

func (k *KinesisSource) GetName() string

func (*KinesisSource) OneShotAcquisition

func (k *KinesisSource) OneShotAcquisition(out chan types.Event, t *tomb.Tomb) error

func (*KinesisSource) ParseAndPushRecords

func (k *KinesisSource) ParseAndPushRecords(records []*kinesis.Record, out chan types.Event, logger *log.Entry, shardId string)

func (*KinesisSource) ReadFromShard

func (k *KinesisSource) ReadFromShard(out chan types.Event, shardId string) error

func (*KinesisSource) ReadFromStream

func (k *KinesisSource) ReadFromStream(out chan types.Event, t *tomb.Tomb) error

func (*KinesisSource) ReadFromSubscription

func (k *KinesisSource) ReadFromSubscription(reader kinesis.SubscribeToShardEventStreamReader, out chan types.Event, shardId string, streamName string) error

func (*KinesisSource) RegisterConsumer

func (k *KinesisSource) RegisterConsumer() (*kinesis.RegisterStreamConsumerOutput, error)

func (*KinesisSource) StreamingAcquisition

func (k *KinesisSource) StreamingAcquisition(out chan types.Event, t *tomb.Tomb) error

func (*KinesisSource) SubscribeToShards

func (k *KinesisSource) SubscribeToShards(arn arn.ARN, streamConsumer *kinesis.RegisterStreamConsumerOutput, out chan types.Event) error

func (*KinesisSource) WaitForConsumerDeregistration

func (k *KinesisSource) WaitForConsumerDeregistration(consumerName string, streamARN string) error

func (*KinesisSource) WaitForConsumerRegistration

func (k *KinesisSource) WaitForConsumerRegistration(consumerARN string) error

Jump to

Keyboard shortcuts

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