public class SimpleKafkaConsumer<K,V> extends Object
| Modifier and Type | Class and Description |
|---|---|
static class |
SimpleKafkaConsumer.Status |
| Constructor and Description |
|---|
SimpleKafkaConsumer() |
| Modifier and Type | Method and Description |
|---|---|
void |
addTopicPartition(String topicPartition)
Add a topic-partition
|
void |
consumeFromPartition(String server,
String partitionString,
KafkaRecordProcessor<K,V>... kafkaRecordProcessors) |
int |
getAutoCommitIntervalMillis() |
org.apache.kafka.clients.consumer.KafkaConsumer<K,V> |
getConsumer() |
Map<String,Object> |
getExtraProps() |
String |
getGroupId() |
org.apache.kafka.common.serialization.Deserializer<K> |
getKeyDeserializer() |
int |
getPollTimeout() |
List<KafkaRecordProcessor<K,V>> |
getRecordProcessors() |
AtomicBoolean |
getRunning() |
String |
getServer() |
int |
getSessionTimeoutMillis() |
EnumCounter<SimpleKafkaConsumer.Status> |
getStatusCounter() |
String |
getTopic() |
org.apache.kafka.common.serialization.Deserializer<V> |
getValueDeserializer() |
void |
init() |
boolean |
isCommitSyncEnabled() |
boolean |
isEnableAutoCommit() |
boolean |
isSpecificPartitions() |
Set<String> |
listTopicNames()
List all topic names
|
void |
preInit() |
void |
reassignTopic(String topic)
Re assign the topics
|
void |
removeTopicPartition(String topicPartition)
Add a topic-partition
|
void |
setAutoCommitIntervalMillis(int autoCommitIntervalMillis) |
void |
setCommitSyncEnabled(boolean commitSyncEnabled) |
void |
setEnableAutoCommit(boolean enableAutoCommit) |
void |
setExtraProps(Map<String,Object> extraProps) |
void |
setGroupId(String groupId) |
void |
setKeyDeserializer(org.apache.kafka.common.serialization.Deserializer<K> keyDeserializer) |
void |
setPollTimeout(int pollTimeout) |
void |
setRecordProcessors(List<KafkaRecordProcessor<K,V>> recordProcessors) |
void |
setServer(String server) |
void |
setSessionTimeoutMillis(int sessionTimeoutMillis) |
void |
setSpecificPartitions(boolean specificPartitions) |
void |
setTopic(String topic) |
void |
setValueDeserializer(org.apache.kafka.common.serialization.Deserializer<V> valueDeserializer) |
void |
shutdown() |
void |
start() |
void |
update() |
@PostConstruct public void init()
public void start()
public void preInit()
public void update()
public void consumeFromPartition(String server, String partitionString, KafkaRecordProcessor<K,V>... kafkaRecordProcessors)
@PreDestroy public void shutdown()
public void addTopicPartition(String topicPartition)
public void removeTopicPartition(String topicPartition)
public void reassignTopic(String topic)
public String getServer()
public void setServer(String server)
server - the server to setpublic String getGroupId()
public void setGroupId(String groupId)
groupId - the groupId to setpublic boolean isEnableAutoCommit()
public void setEnableAutoCommit(boolean enableAutoCommit)
enableAutoCommit - the enableAutoCommit to setpublic boolean isCommitSyncEnabled()
public void setCommitSyncEnabled(boolean commitSyncEnabled)
commitSyncEnabled - the commitSyncEnabled to setpublic int getAutoCommitIntervalMillis()
public void setAutoCommitIntervalMillis(int autoCommitIntervalMillis)
autoCommitIntervalMillis - the autoCommitIntervalMillis to setpublic int getSessionTimeoutMillis()
public void setSessionTimeoutMillis(int sessionTimeoutMillis)
sessionTimeoutMillis - the sessionTimeoutMillis to setpublic org.apache.kafka.common.serialization.Deserializer<K> getKeyDeserializer()
public void setKeyDeserializer(org.apache.kafka.common.serialization.Deserializer<K> keyDeserializer)
keyDeserializer - the keyDeserializer to setpublic org.apache.kafka.common.serialization.Deserializer<V> getValueDeserializer()
public void setValueDeserializer(org.apache.kafka.common.serialization.Deserializer<V> valueDeserializer)
valueDeserializer - the valueDeserializer to setpublic String getTopic()
public void setTopic(String topic)
topic - the topic to setpublic boolean isSpecificPartitions()
public void setSpecificPartitions(boolean specificPartitions)
specificPartitions - the specificPartitions to setpublic int getPollTimeout()
public void setPollTimeout(int pollTimeout)
pollTimeout - the pollTimeout to setpublic org.apache.kafka.clients.consumer.KafkaConsumer<K,V> getConsumer()
public AtomicBoolean getRunning()
public List<KafkaRecordProcessor<K,V>> getRecordProcessors()
public void setRecordProcessors(List<KafkaRecordProcessor<K,V>> recordProcessors)
recordProcessors - the recordProcessors to setpublic void setExtraProps(Map<String,Object> extraProps)
extraProps - the extraProps to setpublic EnumCounter<SimpleKafkaConsumer.Status> getStatusCounter()
Copyright © 2016. All rights reserved.