Package org.apache.pinot.tools.utils
Class KafkaStarterUtils
- java.lang.Object
-
- org.apache.pinot.tools.utils.KafkaStarterUtils
-
public class KafkaStarterUtils extends Object
-
-
Field Summary
Fields Modifier and Type Field Description static StringBROKER_IDstatic intDEFAULT_BROKER_IDstatic StringDEFAULT_KAFKA_BROKERstatic intDEFAULT_KAFKA_PORTstatic StringKAFKA_JSON_MESSAGE_DECODER_CLASS_NAMEstatic StringKAFKA_PRODUCER_CLASS_NAMEstatic StringKAFKA_SERVER_STARTABLE_CLASS_NAMEstatic StringKAFKA_STREAM_CONSUMER_FACTORY_CLASS_NAMEstatic StringKAFKA_STREAM_LEVEL_CONSUMER_CLASS_NAMEstatic StringPORTstatic StringZOOKEEPER_CONNECT
-
Method Summary
All Methods Static Methods Concrete Methods Modifier and Type Method Description static voidconfigureHostName(Properties configuration, String hostName)static voidconfigureOffsetsTopicReplicationFactor(Properties configuration, short replicationFactor)static voidconfigureTopicDeletion(Properties configuration, boolean topicDeletionEnabled)static voidconfigureTransactionStateLogReplicationFactor(Properties configuration, short replicationFactor)static PropertiesgetDefaultKafkaConfiguration()static PropertiesgetDefaultKafkaConfiguration(ZkStarter.ZookeeperInstance zookeeperInstance)static StringgetDefaultKafkaZKAddress()static PropertiesgetTopicCreationProps(int numKafkaPartitions)static org.apache.pinot.spi.stream.StreamDataServerStartablestartServer(int port, int brokerId, String zkStr, Properties baseConf)static List<org.apache.pinot.spi.stream.StreamDataServerStartable>startServers(int brokerCount, int port, String zkStr, Properties configuration)
-
-
-
Field Detail
-
DEFAULT_BROKER_ID
public static final int DEFAULT_BROKER_ID
- See Also:
- Constant Field Values
-
DEFAULT_KAFKA_PORT
public static final int DEFAULT_KAFKA_PORT
- See Also:
- Constant Field Values
-
DEFAULT_KAFKA_BROKER
public static final String DEFAULT_KAFKA_BROKER
- See Also:
- Constant Field Values
-
PORT
public static final String PORT
- See Also:
- Constant Field Values
-
BROKER_ID
public static final String BROKER_ID
- See Also:
- Constant Field Values
-
ZOOKEEPER_CONNECT
public static final String ZOOKEEPER_CONNECT
- See Also:
- Constant Field Values
-
KAFKA_SERVER_STARTABLE_CLASS_NAME
public static final String KAFKA_SERVER_STARTABLE_CLASS_NAME
-
KAFKA_PRODUCER_CLASS_NAME
public static final String KAFKA_PRODUCER_CLASS_NAME
-
KAFKA_STREAM_CONSUMER_FACTORY_CLASS_NAME
public static final String KAFKA_STREAM_CONSUMER_FACTORY_CLASS_NAME
-
KAFKA_STREAM_LEVEL_CONSUMER_CLASS_NAME
public static final String KAFKA_STREAM_LEVEL_CONSUMER_CLASS_NAME
-
KAFKA_JSON_MESSAGE_DECODER_CLASS_NAME
public static final String KAFKA_JSON_MESSAGE_DECODER_CLASS_NAME
- See Also:
- Constant Field Values
-
-
Method Detail
-
getDefaultKafkaConfiguration
public static Properties getDefaultKafkaConfiguration()
-
configureOffsetsTopicReplicationFactor
public static void configureOffsetsTopicReplicationFactor(Properties configuration, short replicationFactor)
-
configureTransactionStateLogReplicationFactor
public static void configureTransactionStateLogReplicationFactor(Properties configuration, short replicationFactor)
-
configureTopicDeletion
public static void configureTopicDeletion(Properties configuration, boolean topicDeletionEnabled)
-
configureHostName
public static void configureHostName(Properties configuration, String hostName)
-
getDefaultKafkaZKAddress
public static String getDefaultKafkaZKAddress()
-
getTopicCreationProps
public static Properties getTopicCreationProps(int numKafkaPartitions)
-
startServers
public static List<org.apache.pinot.spi.stream.StreamDataServerStartable> startServers(int brokerCount, int port, String zkStr, Properties configuration)
-
startServer
public static org.apache.pinot.spi.stream.StreamDataServerStartable startServer(int port, int brokerId, String zkStr, Properties baseConf)
-
getDefaultKafkaConfiguration
public static Properties getDefaultKafkaConfiguration(ZkStarter.ZookeeperInstance zookeeperInstance)
-
-