dispatcher

package
v0.0.0-...-157142f Latest Latest
Warning

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

Go to latest
Published: Apr 3, 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(brokers []string, consumerMode cluster.ConsumerMode, logger *zap.Logger) (*KafkaDispatcher, error)

func (*KafkaDispatcher) ConfigDiff

func (d *KafkaDispatcher) ConfigDiff(updated *multichannelfanout.Config) string

ConfigDiff diffs the new config with the existing config. If there are no differences, then the empty string is returned. If there are differences, then a non-empty string is returned describing the differences.

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

Jump to

Keyboard shortcuts

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