public class ZKLeaderConsumerPartitionCoordinator
extends de.zalando.paradox.nakadi.consumer.core.partitioned.impl.AbstractPartitionCoordinator
| Constructor and Description |
|---|
ZKLeaderConsumerPartitionCoordinator(ZKHolder zkHolder,
java.lang.String consumerName,
java.util.List<de.zalando.paradox.nakadi.consumer.core.http.handlers.EventErrorHandler> eventErrorHandlers) |
| Modifier and Type | Method and Description |
|---|---|
void |
close() |
void |
commit(de.zalando.paradox.nakadi.consumer.core.domain.EventTypeCursor cursor) |
void |
error(int statusCode,
java.lang.String content,
de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition) |
void |
error(java.lang.Throwable t,
de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition,
java.lang.String offset,
java.lang.String rawEvent) |
void |
finished(de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition) |
void |
flush(de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition) |
java.util.Optional<de.zalando.paradox.nakadi.consumer.core.partitioned.PartitionAdminService> |
getAdminService() |
java.lang.String |
getConsumerName() |
void |
init() |
void |
rebalance(de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartitions consumerPartitions,
java.util.Collection<de.zalando.paradox.nakadi.consumer.core.domain.NakadiPartition> nakadiPartitions) |
void |
setDeleteUnavailableCursors(boolean deleteUnavailableCursors) |
void |
setStartNewestAvailableOffset(boolean startNewestAvailableOffset) |
assignPartition, assignPartitions, getPartitionCommitCallback, getPartitionRebalanceListener, getPartitions, getPartitionsToAssign, getPartitionsToRevoke, registerCommitCallback, registerRebalanceListener, revokePartition, revokePartitions, unregisterCommitCallback, unregisterRebalanceListenerpublic ZKLeaderConsumerPartitionCoordinator(ZKHolder zkHolder, java.lang.String consumerName, java.util.List<de.zalando.paradox.nakadi.consumer.core.http.handlers.EventErrorHandler> eventErrorHandlers)
public void close()
public void init()
public void rebalance(de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartitions consumerPartitions,
java.util.Collection<de.zalando.paradox.nakadi.consumer.core.domain.NakadiPartition> nakadiPartitions)
public void finished(de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition)
finished in interface de.zalando.paradox.nakadi.consumer.core.partitioned.PartitionCoordinatorfinished in class de.zalando.paradox.nakadi.consumer.core.partitioned.impl.AbstractPartitionCoordinatorpublic void commit(de.zalando.paradox.nakadi.consumer.core.domain.EventTypeCursor cursor)
public void flush(de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition)
public void error(java.lang.Throwable t,
de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition,
@Nullable
java.lang.String offset,
java.lang.String rawEvent)
public void error(int statusCode,
java.lang.String content,
de.zalando.paradox.nakadi.consumer.core.domain.EventTypePartition eventTypePartition)
public void setStartNewestAvailableOffset(boolean startNewestAvailableOffset)
public void setDeleteUnavailableCursors(boolean deleteUnavailableCursors)
public java.lang.String getConsumerName()
public java.util.Optional<de.zalando.paradox.nakadi.consumer.core.partitioned.PartitionAdminService> getAdminService()