dispatcher

package
v0.7.1 Latest Latest
Warning

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

Go to latest
Published: Jul 2, 2019 License: Apache-2.0 Imports: 13 Imported by: 0

Documentation

Overview

Copyright 2018 The Knative Authors

Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at

http://www.apache.org/licenses/LICENSE-2.0

Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type KafkaCluster

type KafkaCluster interface {
	NewConsumer(groupID string, topics []string) (KafkaConsumer, error)

	GetConsumerMode() cluster.ConsumerMode
}

type KafkaConsumer

type KafkaConsumer interface {
	Messages() <-chan *sarama.ConsumerMessage
	Partitions() <-chan cluster.PartitionConsumer
	MarkOffset(msg *sarama.ConsumerMessage, metadata string)
	Close() (err error)
}

type KafkaDispatcher

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

func NewDispatcher

func NewDispatcher(args *KafkaDispatcherArgs) (*KafkaDispatcher, error)

func (*KafkaDispatcher) Start

func (d *KafkaDispatcher) Start(stopCh <-chan struct{}) error

Start starts the kafka dispatcher's message processing.

func (*KafkaDispatcher) UpdateConfig

func (d *KafkaDispatcher) UpdateConfig(config *multichannelfanout.Config) error

UpdateConfig is used by older kafka channel dispatcher controller that is based on ClusterChannelProvisioners model Remove this function when the older channel code is deleted

func (*KafkaDispatcher) UpdateHostToChannelMap added in v0.7.0

func (d *KafkaDispatcher) UpdateHostToChannelMap(config *multichannelfanout.Config) error

UpdateHostToChannelMap will be called by new CRD based kafka channel dispatcher controller, instead of UpdateConfig.

func (*KafkaDispatcher) UpdateKafkaConsumers added in v0.7.0

func (d *KafkaDispatcher) UpdateKafkaConsumers(config *multichannelfanout.Config) (map[eventingduck.SubscriberSpec]error, error)

UpdateKafkaConsumers will be called by new CRD based kafka channel dispatcher controller, instead of UpdateConfig.

type KafkaDispatcherArgs added in v0.7.0

type KafkaDispatcherArgs struct {
	ClientID     string
	Brokers      []string
	ConsumerMode cluster.ConsumerMode
	TopicFunc    TopicFunc
	Logger       *zap.Logger
}

type TopicFunc added in v0.7.0

type TopicFunc func(separator, namespace, name string) string

Jump to

Keyboard shortcuts

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