public class KafkaSource<K,V> extends Object
Constructor and Description |
---|
KafkaSource(io.vertx.mutiny.core.Vertx vertx,
String consumerGroup,
KafkaConnectorIncomingConfiguration config,
javax.enterprise.inject.Instance<KafkaConsumerRebalanceListener> consumerRebalanceListeners,
KafkaCDIEvents kafkaCDIEvents,
javax.enterprise.inject.Instance<DeserializationFailureHandler<?>> deserializationFailureHandlers,
int index) |
Modifier and Type | Method and Description |
---|---|
void |
closeQuietly() |
String |
getChannel() |
KafkaCommitHandler |
getCommitHandler() |
ReactiveKafkaConsumer<K,V> |
getConsumer()
For testing purpose only
|
io.smallrye.mutiny.Multi<IncomingKafkaRecord<K,V>> |
getStream() |
boolean |
hasSubscribers() |
void |
incomingTrace(IncomingKafkaRecord<K,V> kafkaRecord) |
void |
isAlive(HealthReport.HealthReportBuilder builder) |
void |
isReady(HealthReport.HealthReportBuilder builder) |
void |
reportFailure(Throwable failure,
boolean fatal) |
public KafkaSource(io.vertx.mutiny.core.Vertx vertx, String consumerGroup, KafkaConnectorIncomingConfiguration config, javax.enterprise.inject.Instance<KafkaConsumerRebalanceListener> consumerRebalanceListeners, KafkaCDIEvents kafkaCDIEvents, javax.enterprise.inject.Instance<DeserializationFailureHandler<?>> deserializationFailureHandlers, int index)
public void reportFailure(Throwable failure, boolean fatal)
public void incomingTrace(IncomingKafkaRecord<K,V> kafkaRecord)
public io.smallrye.mutiny.Multi<IncomingKafkaRecord<K,V>> getStream()
public void closeQuietly()
public void isAlive(HealthReport.HealthReportBuilder builder)
public void isReady(HealthReport.HealthReportBuilder builder)
public ReactiveKafkaConsumer<K,V> getConsumer()
public KafkaCommitHandler getCommitHandler()
public String getChannel()
public boolean hasSubscribers()
Copyright © 2018–2021 SmallRye. All rights reserved.