Class GSKafkaConsumerThread

java.lang.Object
com.gigaspaces.dih.consumer.GSKafkaConsumerThread
All Implemented Interfaces:
Runnable

public class GSKafkaConsumerThread extends Object implements Runnable
  • Field Details

    • zookeeperAttributeStore

      public org.openspaces.zookeeper.attribute_store.ZooKeeperAttributeStore zookeeperAttributeStore
  • Constructor Details

  • Method Details

    • getMessageExecutor

      public GSMessageExecutor getMessageExecutor()
    • run

      public void run()
      Specified by:
      run in interface Runnable
    • isActive

      public boolean isActive()
    • setActive

      public void setActive(boolean active)
    • getMessageExecutionRetries

      public int getMessageExecutionRetries()
    • setStatus

      public void setStatus(GSKafkaConsumerStatusEnum newStatus)
    • getStatus

      public GSKafkaConsumerStatusEnum getStatus()
    • increaseNumberOfOperations

      public void increaseNumberOfOperations(int operations)
    • getTotalOperations

      public int getTotalOperations() throws IOException
      Throws:
      IOException
    • getKafkaBootstrapServers

      public String getKafkaBootstrapServers()
    • getKafkaTopic

      public String getKafkaTopic()
    • getKafkaConsumerGroup

      public String getKafkaConsumerGroup()
    • getKafkaProperties

      public Properties getKafkaProperties()
    • setOffset

      public void setOffset(String offset)
    • getPartitionsCount

      public int getPartitionsCount()
    • setPopulateDeletedObjectsTable

      public void setPopulateDeletedObjectsTable(boolean populateDeletedObjectsTable)